1use crate::{Error, Result, Timescale, Timestamp, coding};
17use crate::{broadcast, cache, group, stats};
18
19use super::{Datagram, Requests};
20
21use super::Cap;
22pub use super::subscription::{Position, Subscription};
23
24use std::{
25 collections::{BTreeMap, HashSet, VecDeque},
26 ops::{Bound, RangeBounds},
27 sync::Arc,
28 sync::OnceLock,
29 sync::atomic::{AtomicBool, Ordering},
30 task::{Poll, ready},
31 time::Duration,
32};
33
34pub const DEFAULT_MAX_AGE: Duration = Duration::from_secs(5);
36
37const MAX_DATAGRAMS: usize = 64;
43
44const EVICT_SLACK: usize = 64;
47
48const EVICT_SCAN: usize = 4;
52
53#[derive(Clone, Copy)]
55pub(super) struct ExpiryScan {
56 start: usize,
57 width: usize,
60 now: u64,
61 max_ticks: u64,
62 gc: bool,
64}
65
66#[derive(Clone, Debug)]
77#[non_exhaustive]
78pub struct Info {
79 pub timescale: Timescale,
86 pub max_age: Duration,
106 pub priority: u8,
109}
110
111impl Default for Info {
112 fn default() -> Self {
113 Self {
114 timescale: Timescale::default(),
115 max_age: DEFAULT_MAX_AGE,
116 priority: 0,
117 }
118 }
119}
120
121impl Info {
122 pub fn with_timescale(mut self, timescale: Timescale) -> Self {
127 self.timescale = timescale;
128 self
129 }
130
131 pub fn with_max_age(mut self, max_age: Duration) -> Self {
133 self.max_age = max_age;
134 self
135 }
136
137 pub fn with_priority(mut self, priority: u8) -> Self {
139 self.priority = priority;
140 self
141 }
142}
143
144#[derive(Default)]
145pub(crate) struct TrackState {
146 info: Option<Info>,
149 published: bool,
152 claimed: bool,
155
156 broadcast: Arc<broadcast::Info>,
159
160 cache: Arc<cache::Track>,
164
165 lookup: BTreeMap<u64, Slot>,
171
172 arrival: VecDeque<(u64, u32)>,
177
178 evict: VecDeque<(u64, u32)>,
186
187 debt: u64,
192
193 datagrams: VecDeque<Datagram>,
196
197 datagram_offset: usize,
200
201 offset: usize,
204
205 max_sequence: Option<u64>,
208
209 latest_group: Option<u64>,
214
215 next_stamp: u32,
217
218 final_sequence: Option<u64>,
220
221 start_sequence: Option<u64>,
225
226 start_pending: bool,
231
232 resume: Option<Position>,
236
237 abort: Option<Error>,
239
240 subscriptions: kio::Shared<Subscriptions>,
244
245 fetch: kio::Shared<FetchState>,
248}
249
250struct Slot {
256 group: group::Producer,
257
258 stamp: u32,
263
264 visible: bool,
267}
268
269pub(crate) const CACHE_OVERHEAD: u64 = 2 * (size_of::<u64>() + size_of::<Slot>() + 2 * size_of::<(u64, u32)>()) as u64;
278
279type Subscriptions = Vec<kio::Consumer<Subscription>>;
281
282pub(crate) type FetchState = Requests<u64, PendingFetch>;
287
288pub(crate) struct PendingFetch {
290 priority: u8,
292
293 frame_start: u64,
297
298 result: kio::Producer<FetchOutcome>,
303}
304
305#[derive(Default)]
308pub(crate) struct FetchOutcome {
309 pub(crate) rejected: Option<Error>,
310}
311
312impl TrackState {
313 fn normalize_info(broadcast: &broadcast::Info, mut info: Info) -> Info {
314 info.max_age = info.max_age.min(broadcast.cache_duration);
315 info
316 }
317
318 fn accept(&mut self, info: Info) {
319 self.published = true;
320 self.install(info);
321 }
322
323 fn poll_info(&self) -> Poll<Result<Info>> {
324 if let Some(info) = &self.info {
325 Poll::Ready(Ok(info.clone()))
326 } else if let Some(err) = &self.abort {
327 Poll::Ready(Err(err.clone()))
330 } else {
331 Poll::Pending
332 }
333 }
334
335 fn poll_recv_group(&self, index: usize, min_sequence: u64) -> Poll<Result<Option<(group::Producer, usize)>>> {
341 let start = index.saturating_sub(self.offset);
342 for (i, (sequence, stamp)) in self.arrival.iter().enumerate().skip(start) {
343 if *sequence >= min_sequence
344 && let Some(slot) = self.lookup.get(sequence)
345 && slot.stamp == *stamp
346 && !slot.group.is_aborted()
347 {
348 return Poll::Ready(Ok(Some((slot.group.clone(), self.offset + i))));
349 }
350 }
351
352 if self.is_complete() {
354 Poll::Ready(Ok(None))
355 } else if let Some(err) = &self.abort {
356 Poll::Ready(Err(err.clone()))
357 } else {
358 Poll::Pending
359 }
360 }
361
362 fn poll_recv_datagram(&self, index: usize) -> Poll<Result<Option<(Datagram, usize)>>> {
368 let start = index.saturating_sub(self.datagram_offset);
369 if let Some(datagram) = self.datagrams.get(start) {
370 return Poll::Ready(Ok(Some((datagram.clone(), self.datagram_offset + start))));
371 }
372
373 if self.is_complete() {
375 Poll::Ready(Ok(None))
376 } else if let Some(err) = &self.abort {
377 Poll::Ready(Err(err.clone()))
378 } else {
379 Poll::Pending
380 }
381 }
382
383 fn push_datagram(&mut self, datagram: Datagram) {
385 if self.datagrams.len() == MAX_DATAGRAMS {
386 self.datagrams.pop_front();
387 self.datagram_offset += 1;
388 }
389 self.datagrams.push_back(datagram);
390 }
391
392 fn poll_next_in_range(
402 &self,
403 next_sequence: u64,
404 end_sequence: Option<u64>,
405 ) -> Poll<Result<Option<group::Producer>>> {
406 if let Some(end) = end_sequence
411 && end <= next_sequence
412 {
413 if let Some(err) = &self.abort {
414 return Poll::Ready(Err(err.clone()));
415 }
416 return Poll::Pending;
417 }
418
419 let best = self
420 .lookup
421 .range(next_sequence..)
422 .map(|(_, slot)| &slot.group)
423 .take_while(|group| super::subscription::before_end(group.sequence, end_sequence))
424 .find(|group| !group.is_aborted());
425
426 if let Some(group) = best {
427 return Poll::Ready(Ok(Some(group.clone())));
432 }
433
434 if let Some(err) = &self.abort {
436 return Poll::Ready(Err(err.clone()));
437 }
438 if let Some(fin) = self.final_sequence
441 && next_sequence >= fin
442 {
443 return Poll::Ready(Ok(None));
444 }
445 Poll::Pending
446 }
447
448 fn covering_group(&self, sequence: u64, frame_start: u64) -> Option<&group::Producer> {
455 let slot = self.lookup.get(&sequence)?;
456 let first = slot.group.live_first_frame()?;
457 (first as u64 <= frame_start).then_some(&slot.group)
458 }
459
460 fn max_age_bound(&self) -> Option<Duration> {
463 self.info.as_ref().map(|info| info.max_age)
464 }
465
466 fn live_edge(&self, cap: Option<u64>) -> Option<PresentationEdge> {
481 self.lookup
482 .range(..)
483 .rev()
484 .filter(|(seq, _)| super::subscription::before_end(**seq, cap))
485 .find_map(|(_, slot)| {
486 if !slot.visible || slot.group.is_aborted() {
487 return None;
488 }
489 let timestamp = slot.group.timestamp()?;
492 Some(PresentationEdge {
493 sequence: slot.group.sequence,
494 stamp: slot.stamp,
495 timestamp: slot.group.latest().unwrap_or(timestamp),
496 })
497 })
498 }
499
500 fn drift_edge(&self, cap: Option<u64>, outer: Option<(u64, Timestamp)>, successor: Option<Timestamp>) -> Edge {
504 Edge {
505 presentation: self.live_edge(cap),
506 outer,
507 cap,
508 successor,
509 }
510 }
511
512 fn holds(&self, sequence: u64, stamp: u32) -> bool {
514 self.lookup
515 .get(&sequence)
516 .is_some_and(|slot| slot.stamp == stamp && !slot.group.is_aborted())
517 }
518
519 fn reach(&self, sequence: u64, cap: Option<u64>, beyond: Option<Timestamp>) -> Option<Timestamp> {
534 match self.first_start(sequence.saturating_add(1), cap) {
535 Some(start) => start,
536 None => beyond,
537 }
538 }
539
540 fn first_start(&self, from: u64, cap: Option<u64>) -> Option<Option<Timestamp>> {
543 let slot = self
544 .lookup
545 .range(from..)
546 .map(|(_, slot)| slot)
547 .take_while(|slot| super::subscription::before_end(slot.group.sequence, cap))
548 .find(|slot| slot.visible && !slot.group.is_aborted())?;
549 Some(slot.group.timestamp())
550 }
551
552 fn served_start(&self, from: u64, cap: Option<u64>) -> Option<ServedStart> {
557 let slot = self
558 .lookup
559 .range(from..)
560 .map(|(_, slot)| slot)
561 .take_while(|slot| super::subscription::before_end(slot.group.sequence, cap))
562 .find(|slot| slot.visible && !slot.group.is_aborted())?;
563 Some(ServedStart {
564 sequence: slot.group.sequence,
565 stamp: slot.stamp,
566 timestamp: slot.group.timestamp()?,
567 })
568 }
569
570 fn is_stale(&self, sequence: u64, edge: &Edge, budget: Duration) -> bool {
597 if !self.lookup.contains_key(&sequence) {
598 return false;
599 }
600
601 let local = edge
605 .presentation
606 .filter(|live| self.holds(live.sequence, live.stamp))
607 .map(|live| (live.sequence, live.timestamp));
608 let Some((_, timestamp)) = local
612 .into_iter()
613 .chain(edge.outer)
614 .filter(|(live, _)| *live > sequence)
615 .max_by_key(|(live, _)| *live)
616 else {
617 return false;
618 };
619 self.reach(sequence, edge.cap, edge.successor)
620 .is_some_and(|reach| matches!(timestamp.checked_sub(reach), Ok(age) if Duration::from(age) >= budget))
621 }
622
623 fn poll_fetch_cached(&self, sequence: u64, frame_start: u64) -> Poll<Result<group::Consumer>> {
628 if let Some(group) = self.covering_group(sequence, frame_start) {
629 group.cache_refresh();
633 return Poll::Ready(Ok(group.consume()));
634 }
635
636 if let Some(err) = &self.abort {
637 return Poll::Ready(Err(err.clone()));
638 }
639
640 if self.final_sequence.is_some_and(|fin| sequence >= fin) {
642 return Poll::Ready(Err(Error::NotFound));
643 }
644
645 Poll::Pending
646 }
647
648 pub(super) fn evict_expired(&mut self) {
663 let scan = self.expiry_scan();
664 self.evict_expired_scan(scan);
665 }
666
667 pub(super) fn expiry_scan(&self) -> ExpiryScan {
669 ExpiryScan {
670 start: self.cache.next_expiry_scan(EVICT_SCAN),
671 width: EVICT_SCAN,
672 now: self.cache.pool().now(),
673 max_ticks: self.cache.pool().expiry_ticks(),
674 gc: false,
675 }
676 }
677
678 pub(super) fn expiry_scan_drain(&self) -> ExpiryScan {
681 ExpiryScan {
682 start: 0,
683 width: self.evict.len(),
684 now: self.cache.pool().now(),
685 max_ticks: self.cache.pool().expiry_ticks(),
686 gc: true,
687 }
688 }
689
690 #[cfg(test)]
691 pub(super) fn date_cache_accesses(&self, now: u64) {
692 for slot in self.lookup.values() {
693 slot.group.cache_accessed_tick(Some(now));
694 }
695 }
696
697 pub(super) fn expiry_mutation_due(&self, scan: ExpiryScan) -> bool {
702 let len = self.evict.len();
703 if len > 0 {
704 let start = scan.start % len;
705 let mut retained = 0;
706 for step in 0..len.min(scan.width) {
707 let (sequence, stamp) = self.evict[(start + step) % len];
708 let Some(slot) = self.lookup.get(&sequence) else {
709 continue;
710 };
711 if slot.stamp != stamp {
712 continue;
713 }
714 if slot.group.is_aborted()
715 || (Some(sequence) != self.latest_group
716 && slot
717 .group
718 .cache_accessed_tick(scan.gc.then_some(scan.now))
719 .is_some_and(|tick| scan.now.saturating_sub(tick) > scan.max_ticks))
720 {
721 return true;
722 }
723 retained += 1;
724 if !scan.gc && retained >= EVICT_SCAN {
725 break;
726 }
727 }
728 }
729
730 self.arrival
731 .front()
732 .is_some_and(|(sequence, stamp)| !self.is_current(*sequence, *stamp))
733 || self
734 .evict
735 .front()
736 .is_some_and(|(sequence, stamp)| !self.is_current(*sequence, *stamp))
737 || self.evict.len() > 2 * self.lookup.len() + EVICT_SLACK
738 }
739
740 pub(super) fn evict_expired_scan(&mut self, scan: ExpiryScan) {
743 let len = self.evict.len();
744 if len > 0 {
745 let start = scan.start % len;
746 let mut retained = 0;
747 for step in 0..len.min(scan.width) {
748 let (sequence, stamp) = self.evict[(start + step) % len];
749 let Some(slot) = self.lookup.get(&sequence) else {
750 continue;
751 };
752 if slot.stamp != stamp {
753 continue;
755 }
756 if slot.group.is_aborted() {
759 self.lookup.remove(&sequence);
760 continue;
761 }
762 if Some(sequence) == self.latest_group
763 || slot
764 .group
765 .cache_accessed_tick(scan.gc.then_some(scan.now))
766 .is_none_or(|tick| scan.now.saturating_sub(tick) <= scan.max_ticks)
767 {
768 retained += 1;
771 if !scan.gc && retained >= EVICT_SCAN {
772 break;
773 }
774 continue;
775 }
776 let slot = self.lookup.remove(&sequence).unwrap();
780 let _ = slot.group.abort(Error::Old);
781 }
782 }
783
784 while let Some((sequence, stamp)) = self.arrival.front() {
787 if self.lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp) {
788 break;
789 }
790 self.arrival.pop_front();
791 self.offset += 1;
792 }
793
794 while let Some((sequence, stamp)) = self.evict.front() {
796 if self.lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp) {
797 break;
798 }
799 self.evict.pop_front();
800 }
801
802 if self.evict.len() > 2 * self.lookup.len() + EVICT_SLACK {
805 let lookup = &self.lookup;
806 self.evict
807 .retain(|(sequence, stamp)| lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp));
808 }
809 }
810
811 fn is_current(&self, sequence: u64, stamp: u32) -> bool {
813 self.lookup.get(&sequence).is_some_and(|slot| slot.stamp == stamp)
814 }
815
816 fn clear_cache(&mut self) {
819 self.lookup.clear();
820 self.arrival.clear();
821 self.evict.clear();
822 self.latest_group = None;
823 self.debt = 0;
824 }
825
826 fn install(&mut self, info: Info) {
832 let info = Self::normalize_info(&self.broadcast, info);
833 self.info = Some(info);
834 }
835
836 fn spawn(broadcast: Arc<broadcast::Info>) -> kio::Producer<Self> {
843 let state = kio::Producer::new(Self {
844 broadcast: broadcast.clone(),
845 ..Default::default()
846 });
847 let cache = cache::Track::new(broadcast.pool.clone(), state.downgrade());
848 state.write().ok().expect("a new track is open").cache = cache;
849 state
850 }
851
852 fn claim_sequence(&mut self, sequence: u64, frame_start: u64) -> Result<()> {
861 if let Some(slot) = self.lookup.get(&sequence) {
862 if slot
866 .group
867 .live_first_frame()
868 .is_some_and(|first| first as u64 <= frame_start)
869 {
870 return Err(Error::Duplicate);
871 }
872 self.lookup.remove(&sequence);
873 }
874 Ok(())
875 }
876
877 fn insert_group(&mut self, group: &group::Producer, visible: bool) {
884 let sequence = group.sequence;
885 self.next_stamp = self.next_stamp.wrapping_add(1);
886 let stamp = self.next_stamp;
887
888 if self.latest_group.is_none_or(|latest| sequence >= latest) {
892 if let Some(latest) = self.latest_group
895 && sequence > latest
896 && let Some(prev) = self.lookup.get(&latest)
897 {
898 prev.group.cache_demote();
899 self.evict.push_back((latest, prev.stamp));
900 }
901 self.latest_group = Some(sequence);
902 } else {
903 group.cache_demote();
904 self.evict.push_back((sequence, stamp));
905 }
906
907 self.max_sequence = Some(self.max_sequence.map_or(sequence, |max| max.max(sequence)));
908 self.lookup.insert(
909 sequence,
910 Slot {
911 group: group.clone(),
912 stamp,
913 visible,
914 },
915 );
916 if visible {
917 self.arrival.push_back((sequence, stamp));
918 }
919 }
920
921 fn commit_group(&mut self, group: &group::Producer, visible: bool) {
925 self.charge_debt();
926 self.insert_group(group, visible);
927 self.evict_expired();
928 }
929
930 pub(super) fn charge_debt(&mut self) {
943 let written = self.cache.take_written();
944 let pool = self.cache.pool().clone();
945 match pool.accrue(written) {
946 Some(mut accrued) => {
947 if self.oldest_is_stale(&pool) {
948 accrued = accrued.saturating_mul(2);
949 }
950 self.debt = self.debt.saturating_add(accrued).min(pool.used());
953 self.pay_debt(&pool, written.saturating_mul(2));
956 }
957 None => self.debt = 0,
960 }
961 }
962
963 fn oldest_is_stale(&self, pool: &cache::Pool) -> bool {
967 let Some(average) = pool.average() else {
968 return false;
969 };
970 let Some((sequence, stamp)) = self.evict.front() else {
971 return false;
972 };
973 let Some(slot) = self.lookup.get(sequence) else {
974 return false;
975 };
976 slot.stamp == *stamp && !slot.group.is_aborted() && slot.group.cache_accessed() <= average
977 }
978
979 fn pay_debt(&mut self, pool: &cache::Pool, cap: u64) {
992 let average = pool.average().unwrap_or(0);
993 let mut paid = 0u64;
994 let mut scanned = 0usize;
995 for _ in 0..self.evict.len() {
996 if self.debt == 0 || paid >= cap || scanned >= EVICT_SCAN {
997 return;
998 }
999 let Some((sequence, stamp)) = self.evict.pop_front() else {
1000 return;
1001 };
1002 let Some(slot) = self.lookup.get(&sequence) else {
1003 continue;
1005 };
1006 if slot.stamp != stamp {
1007 continue;
1009 }
1010 if slot.group.is_aborted() {
1011 self.lookup.remove(&sequence);
1013 continue;
1014 }
1015 if Some(sequence) == self.latest_group {
1016 self.evict.push_back((sequence, stamp));
1018 continue;
1019 }
1020
1021 scanned += 1;
1022 if slot.group.cache_accessed() > average {
1026 self.evict.push_back((sequence, stamp));
1027 continue;
1028 }
1029 let size = slot.group.cache_size();
1032 if size > self.debt {
1033 self.evict.push_front((sequence, stamp));
1034 return;
1035 }
1036
1037 self.debt -= size;
1038 paid = paid.saturating_add(size);
1039 let slot = self.lookup.remove(&sequence).unwrap();
1040 let _ = slot.group.abort(Error::Evicted);
1041 }
1042 }
1043
1044 fn set_start(&mut self, start_sequence: Option<u64>, pending: bool) {
1049 self.start_sequence = start_sequence;
1050 self.start_pending = pending;
1051 }
1052
1053 fn set_final(&mut self, final_sequence: u64) -> Result<()> {
1056 if self.final_sequence.is_some() {
1057 return Err(Error::Closed);
1058 }
1059 if let Some(max) = self.max_sequence
1060 && final_sequence <= max
1061 {
1062 return Err(Error::ProtocolViolation);
1063 }
1064 self.final_sequence = Some(final_sequence);
1065 Ok(())
1066 }
1067
1068 fn is_complete(&self) -> bool {
1074 self.final_sequence
1075 .is_some_and(|fin| self.max_sequence.map_or(0, |max| max.saturating_add(1)) >= fin)
1076 }
1077
1078 fn resume_position(&self) -> Option<Position> {
1084 if self.resume.is_some() {
1087 return self.resume;
1088 }
1089
1090 let max = self.latest_group?;
1091 match self.lookup.get(&max).and_then(|slot| slot.group.resume_frame()) {
1092 Some(frame) => Some(Position {
1095 group: max,
1096 frame: frame as u64,
1097 }),
1098 None => Some(Position::group(max.saturating_add(1))),
1103 }
1104 }
1105
1106 fn poll_finished(&self) -> Poll<Result<u64>> {
1107 if let Some(fin) = self.final_sequence {
1108 Poll::Ready(Ok(fin))
1109 } else if let Some(err) = &self.abort {
1110 Poll::Ready(Err(err.clone()))
1111 } else {
1112 Poll::Pending
1113 }
1114 }
1115
1116 fn modify(producer: &kio::Producer<Self>) -> Result<kio::Mut<'_, Self>> {
1117 producer.write().map_err(|r| r.abort.clone().unwrap_or(Error::Dropped))
1118 }
1119
1120 pub(crate) fn insert_group_request(
1126 &mut self,
1127 sequence: u64,
1128 frame_start: u64,
1129 info: Option<Info>,
1130 ) -> Result<group::Producer> {
1131 if let Some(err) = &self.abort {
1132 return Err(err.clone());
1133 }
1134 if let Some(fin) = self.final_sequence
1135 && sequence >= fin
1136 {
1137 return Err(Error::Closed);
1138 }
1139
1140 if self.info.is_none() {
1144 self.install(info.unwrap_or_default());
1145 }
1146 let info = self.info.clone().unwrap();
1147
1148 self.claim_sequence(sequence, frame_start)?;
1150
1151 let group = group::Producer::new(group::Info { sequence }, info, self.cache.clone());
1152 group.cache_refresh();
1157 self.commit_group(&group, false);
1158 Ok(group)
1159 }
1160}
1161
1162fn commit_abort(mut state: kio::Mut<'_, TrackState>, err: Error) {
1165 state.resume = state.resume_position();
1168 state.abort = Some(err);
1169 state.clear_cache();
1170 state.datagrams.clear();
1171 state.close();
1172}
1173
1174#[derive(Clone)]
1176pub struct Producer {
1177 name: Arc<str>,
1178 info: Info,
1179 broadcast: Arc<broadcast::Info>,
1182 state: kio::Producer<TrackState>,
1183 prev_subscription: Option<Subscription>,
1184 alive: Arc<Alive>,
1186 stats: stats::Scope,
1190}
1191
1192impl Producer {
1193 pub(crate) fn publisher_priority(&self) -> u8 {
1195 self.info.priority
1196 }
1197
1198 pub(crate) fn new(
1206 broadcast: Arc<broadcast::Info>,
1207 name: impl Into<Arc<str>>,
1208 info: impl Into<Option<Info>>,
1209 ) -> Self {
1210 let name = name.into();
1211 let info = TrackState::normalize_info(&broadcast, info.into().unwrap_or_default());
1212 let state = TrackState::spawn(broadcast.clone());
1213 state.write().ok().expect("a new track is open").accept(info.clone());
1214 let alive = Alive::new(name.clone(), state.clone());
1215 alive.publish(None);
1216 Self {
1217 name,
1218 info,
1219 state,
1220 broadcast,
1221 prev_subscription: None,
1222 alive,
1223 stats: stats::Scope::default(),
1224 }
1225 }
1226
1227 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
1231 self.alive.publish(Some(&scope));
1232 self.stats = scope;
1233 self
1234 }
1235
1236 pub fn name(&self) -> &str {
1238 &self.name
1239 }
1240
1241 pub fn broadcast(&self) -> &broadcast::Info {
1243 &self.broadcast
1244 }
1245
1246 pub(crate) fn adopt_group(&mut self, group: group::Producer, visible: bool) -> Result<()> {
1251 let mut state = self.modify()?;
1252 if let Some(fin) = state.final_sequence
1253 && group.sequence >= fin
1254 {
1255 return Err(Error::Closed);
1256 }
1257 if state.lookup.contains_key(&group.sequence) {
1258 return Err(Error::Duplicate);
1259 }
1260 state.insert_group(&group, visible);
1261 Ok(())
1262 }
1263
1264 pub fn create_group(&self, group: group::Info) -> Result<group::Producer> {
1266 let mut state = self.modify()?;
1267 if let Some(fin) = state.final_sequence
1268 && group.sequence >= fin
1269 {
1270 return Err(Error::Closed);
1271 }
1272 let track = state.info.clone().unwrap();
1273
1274 state.claim_sequence(group.sequence, 0)?;
1276
1277 let group = group::Producer::new(group, track, state.cache.clone()).with_meter(self.stats.meter());
1278 state.commit_group(&group, true);
1279
1280 Ok(group)
1281 }
1282
1283 pub fn append_group(&self) -> Result<group::Producer> {
1285 let mut state = self.modify()?;
1286 let sequence = match state.max_sequence {
1287 Some(s) => s.checked_add(1).ok_or(coding::BoundsExceeded)?,
1288 None => 0,
1289 };
1290 if let Some(fin) = state.final_sequence
1291 && sequence >= fin
1292 {
1293 return Err(Error::Closed);
1294 }
1295
1296 let track = state.info.clone().unwrap();
1297
1298 let group =
1299 group::Producer::new(group::Info { sequence }, track, state.cache.clone()).with_meter(self.stats.meter());
1300 state.commit_group(&group, true);
1301
1302 Ok(group)
1303 }
1304
1305 pub fn append_datagram<B: crate::IntoBytes>(&mut self, timestamp: Timestamp, payload: B) -> Result<u64> {
1317 let payload = payload.into_bytes();
1318 if payload.len() > super::datagram::MAX_DATAGRAM_PAYLOAD {
1319 return Err(Error::FrameTooLarge);
1320 }
1321 let meter = self.stats.meter();
1323 let mut state = self.modify()?;
1324 let timescale = state.info.as_ref().unwrap().timescale;
1326 let timestamp = timestamp.convert(timescale).map_err(|_| Error::TimestampMismatch)?;
1327 let sequence = match state.max_sequence {
1328 Some(s) => s.checked_add(1).ok_or(coding::BoundsExceeded)?,
1329 None => 0,
1330 };
1331 if let Some(fin) = state.final_sequence
1332 && sequence >= fin
1333 {
1334 return Err(Error::Closed);
1335 }
1336 state.max_sequence = Some(sequence);
1337 meter.datagram(payload.len() as u64);
1338 state.push_datagram(Datagram {
1339 sequence,
1340 timestamp,
1341 payload,
1342 });
1343 let cache = state.cache.clone();
1344 drop(state);
1345 cache.settle(None);
1346 Ok(sequence)
1347 }
1348
1349 pub fn insert_datagram<B: crate::IntoBytes>(
1355 &mut self,
1356 sequence: u64,
1357 timestamp: Timestamp,
1358 payload: B,
1359 ) -> Result<()> {
1360 let payload = payload.into_bytes();
1361 if payload.len() > super::datagram::MAX_DATAGRAM_PAYLOAD {
1362 return Err(Error::FrameTooLarge);
1363 }
1364 let meter = self.stats.meter();
1366 let mut state = self.modify()?;
1367 let timescale = state.info.as_ref().unwrap().timescale;
1369 let timestamp = timestamp.convert(timescale).map_err(|_| Error::TimestampMismatch)?;
1370 if let Some(fin) = state.final_sequence
1371 && sequence >= fin
1372 {
1373 return Err(Error::Closed);
1374 }
1375 state.max_sequence = Some(state.max_sequence.unwrap_or(0).max(sequence));
1376 meter.datagram(payload.len() as u64);
1377 state.push_datagram(Datagram {
1378 sequence,
1379 timestamp,
1380 payload,
1381 });
1382 let cache = state.cache.clone();
1383 drop(state);
1384 cache.settle(None);
1385 Ok(())
1386 }
1387
1388 pub fn write_frame<B: crate::IntoBytes>(&mut self, timestamp: Timestamp, frame: B) -> Result<()> {
1393 let frame = crate::IntoBytes::into_bytes(frame);
1394 if frame.len() as u64 > group::MAX_CACHE_BYTES {
1395 return Err(Error::FrameTooLarge);
1396 }
1397 let mut group = self.append_group()?;
1398 group.write_frame(timestamp, frame)?;
1399 group.finish()?;
1400 Ok(())
1401 }
1402
1403 pub fn finish(&self) -> Result<()> {
1409 let mut state = self.modify()?;
1410 let final_sequence = match state.max_sequence {
1411 Some(max) => max.checked_add(1).ok_or(coding::BoundsExceeded)?,
1412 None => 0,
1413 };
1414 state.set_final(final_sequence)
1415 }
1416
1417 pub fn finish_at(&mut self, final_sequence: u64) -> Result<()> {
1430 self.modify()?.set_final(final_sequence)
1431 }
1432
1433 pub fn start_at(&mut self, sequence: impl Into<Option<u64>>) -> Result<()> {
1444 self.modify()?.set_start(sequence.into(), false);
1445 Ok(())
1446 }
1447
1448 pub(crate) fn request_start(&mut self, sequence: Option<u64>) -> Result<()> {
1453 self.modify()?.set_start(sequence, true);
1454 Ok(())
1455 }
1456
1457 #[cfg(test)]
1460 pub(crate) fn start_sequence(&self) -> Option<u64> {
1461 self.state.read().start_sequence
1462 }
1463
1464 pub fn final_sequence(&self) -> Option<u64> {
1469 self.state.read().final_sequence
1470 }
1471
1472 pub fn abort(self, err: Error) -> Result<()> {
1483 commit_abort(self.modify()?, err);
1484 Ok(())
1485 }
1486
1487 #[expect(
1495 clippy::result_large_err,
1496 reason = "return the owned producer without allocating on an idle check"
1497 )]
1498 pub fn abort_unused(self, err: Error) -> std::result::Result<(), Self> {
1499 match self.state.write_unused() {
1500 kio::Unused::Idle(guard) => {
1501 commit_abort(guard, err);
1502 return Ok(());
1503 }
1504 kio::Unused::Closed => return Ok(()),
1505 kio::Unused::Used => {}
1506 }
1507 Err(self)
1508 }
1509
1510 pub fn is_used(&self) -> bool {
1517 !self.is_closed() && self.state.is_used()
1518 }
1519
1520 pub async fn unused(&self) -> Result<()> {
1522 self.state.unused().await.map_err(|_| self.abort_reason())
1523 }
1524
1525 pub async fn used(&self) -> Result<()> {
1527 self.state.used().await.map_err(|_| self.abort_reason())
1528 }
1529
1530 pub async fn closed(&self) -> Error {
1532 kio::wait(|waiter| self.poll_closed(waiter)).await
1533 }
1534
1535 pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<Error> {
1537 self.state.poll_closed(waiter).map(|()| self.abort_reason())
1538 }
1539
1540 fn abort_reason(&self) -> Error {
1542 self.state.read().abort.clone().unwrap_or(Error::Dropped)
1543 }
1544
1545 pub fn is_closed(&self) -> bool {
1547 self.state.read().is_closed()
1548 }
1549
1550 pub fn latest(&self) -> Option<u64> {
1552 self.state.read().max_sequence
1553 }
1554
1555 pub fn is_clone(&self, other: &Self) -> bool {
1557 self.state.same_channel(&other.state)
1558 }
1559
1560 pub(crate) fn weak(&self) -> TrackWeak {
1562 TrackWeak {
1563 name: self.name.clone(),
1564 state: self.state.weak(),
1565 }
1566 }
1567
1568 pub fn demand(&self) -> Demand {
1576 Demand {
1577 name: self.name.clone(),
1578 state: self.state.weak(),
1579 }
1580 }
1581
1582 pub fn consume(&self) -> Consumer {
1587 Consumer::plain(self.name.clone(), self.state.consume())
1588 }
1589
1590 pub fn subscribe(&self, subscription: impl Into<Option<Subscription>>) -> Subscriber {
1600 let preferences = subscription.into().unwrap_or_default();
1601
1602 let info = self.info.clone();
1607 let min_sequence = floor_of(&preferences);
1608 let subscription = kio::Producer::new(preferences);
1609 register_subscription(self.state.read(), &subscription);
1610 let drift_anchor = kio::Producer::new(Anchor::default());
1611
1612 let broadcast = self.state.read().broadcast.clone();
1615 Subscriber {
1616 name: self.name.clone(),
1617 broadcast,
1618 info,
1619 inner: SubscriberKind::Plain(PlainSubscriber {
1620 state: self.state.consume(),
1621 subscription,
1622 min_sequence,
1623 index: 0,
1624 datagram_index: 0,
1625 next_sequence: 0,
1626 end_sequence: None,
1627 parked: BTreeMap::new(),
1628 stale_cap: None,
1629 drift_anchor,
1630 stale: stats::Content::default(),
1631 seek_pending: BTreeMap::new(),
1632 }),
1633 stats: stats::Scope::default(),
1635 _stats_sub: stats::Subscription::default(),
1636 }
1637 }
1638
1639 pub async fn subscription_changed(&mut self) -> Result<Option<Subscription>> {
1645 kio::wait(|waiter| self.poll_subscription_changed(waiter)).await
1646 }
1647
1648 pub fn subscription(&self) -> Option<Subscription> {
1656 let state = self.state.read();
1657 let (subs, bound) = (state.subscriptions.clone(), state.max_age_bound());
1658 drop(state);
1659 snapshot_subscription(&subs, bound)
1660 }
1661
1662 pub fn poll_subscription_changed(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Subscription>>> {
1666 if self.state.poll_closed(waiter).is_ready() {
1669 let abort = self.state.read().abort.clone();
1670 return Poll::Ready(Err(abort.unwrap_or(Error::Dropped)));
1671 }
1672
1673 let state = self.state.read();
1675 let (subs, bound) = (state.subscriptions.clone(), state.max_age_bound());
1676 drop(state);
1677
1678 let prev = &self.prev_subscription;
1679 let mut combined = None;
1680 let mut guard = ready!(subs.poll(waiter, |subs| {
1681 let next = combined_subscription(subs, bound, waiter);
1682 if &next == prev {
1683 Poll::Pending
1684 } else {
1685 combined = next;
1686 Poll::Ready(())
1687 }
1688 }));
1689 guard.retain(|sub| !sub.is_closed());
1691 drop(guard);
1692 self.prev_subscription = combined.clone();
1693 Poll::Ready(Ok(combined))
1694 }
1695
1696 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<Result<()>> {
1698 self.state.poll_unused(waiter).map(|used| match used {
1699 Some(()) => Ok(()),
1700 None => Err(self.abort_reason()),
1701 })
1702 }
1703
1704 pub fn dynamic(&self) -> Dynamic {
1708 Dynamic::new(self.name.clone(), self.state.clone(), self.alive.clone())
1709 }
1710
1711 fn modify(&self) -> Result<kio::Mut<'_, TrackState>> {
1712 TrackState::modify(&self.state)
1713 }
1714}
1715
1716fn poll_requested_group(
1720 state: &kio::Producer<TrackState>,
1721 fetch: &kio::Shared<FetchState>,
1722 waiter: &kio::Waiter,
1723) -> Poll<Result<group::Request>> {
1724 if let Poll::Ready(mut guard) = fetch.poll(waiter, |fetch| {
1726 if fetch.has_queued() {
1727 Poll::Ready(())
1728 } else {
1729 Poll::Pending
1730 }
1731 }) {
1732 let sequence = guard.pop().expect("predicate guaranteed a request");
1733 let pending = guard.get(&sequence).expect("popped key must be pending");
1737 let priority = pending.priority;
1738 let frame_start = pending.frame_start;
1739 let result = pending.result.clone();
1740 drop(guard);
1741 return Poll::Ready(Ok(group::Request {
1742 state: state.clone(),
1743 fetch: fetch.clone(),
1744 sequence,
1745 priority,
1746 frame_start,
1747 result,
1748 done: false,
1749 }));
1750 }
1751
1752 match state.poll_ref(waiter, |state| match &state.abort {
1754 Some(err) => Poll::Ready(err.clone()),
1755 None => Poll::Pending,
1756 }) {
1757 Poll::Ready(Ok(err)) => Poll::Ready(Err(err)),
1758 Poll::Ready(Err(closed)) => Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped))),
1759 Poll::Pending => Poll::Pending,
1760 }
1761}
1762
1763pub struct Dynamic {
1773 name: Arc<str>,
1774 state: kio::Producer<TrackState>,
1776 fetch: kio::Shared<FetchState>,
1778 alive: Arc<Alive>,
1781}
1782
1783impl Dynamic {
1784 fn new(name: Arc<str>, state: kio::Producer<TrackState>, alive: Arc<Alive>) -> Self {
1785 let fetch = state.read().fetch.clone();
1786 fetch.lock().add_handler();
1787 Self {
1788 name,
1789 state,
1790 fetch,
1791 alive,
1792 }
1793 }
1794
1795 pub fn name(&self) -> &str {
1797 &self.name
1798 }
1799
1800 pub async fn requested_group(&self) -> Result<group::Request> {
1806 kio::wait(|waiter| self.poll_requested_group(waiter)).await
1807 }
1808
1809 pub fn poll_requested_group(&self, waiter: &kio::Waiter) -> Poll<Result<group::Request>> {
1811 poll_requested_group(&self.state, &self.fetch, waiter)
1812 }
1813
1814 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
1816 self.state.poll_unused(waiter).map(|_| ())
1817 }
1818}
1819
1820impl Clone for Dynamic {
1821 fn clone(&self) -> Self {
1822 self.fetch.lock().add_handler();
1824 Self {
1825 name: self.name.clone(),
1826 state: self.state.clone(),
1827 fetch: self.fetch.clone(),
1828 alive: self.alive.clone(),
1829 }
1830 }
1831}
1832
1833impl Drop for Dynamic {
1834 fn drop(&mut self) {
1835 let mut fetch = self.fetch.lock();
1841 if fetch.remove_handler() {
1842 fetch.drain_queued();
1843 }
1844 }
1845}
1846
1847struct Alive {
1856 name: Arc<str>,
1857 state: kio::Producer<TrackState>,
1858
1859 published: AtomicBool,
1862
1863 stats: OnceLock<stats::Subscription>,
1866}
1867
1868impl Alive {
1869 fn new(name: Arc<str>, state: kio::Producer<TrackState>) -> Arc<Self> {
1870 Arc::new(Self {
1871 name,
1872 state,
1873 published: Default::default(),
1874 stats: Default::default(),
1875 })
1876 }
1877
1878 fn publish(&self, stats: Option<&stats::Scope>) {
1882 self.published.store(true, Ordering::Relaxed);
1883 if let Some(scope) = stats {
1884 let _ = self.stats.set(scope.subscribe());
1887 }
1888 }
1889}
1890
1891impl Drop for Alive {
1892 fn drop(&mut self) {
1893 if !self.published.load(Ordering::Relaxed) {
1895 return;
1896 }
1897 match self.state.write() {
1904 Ok(mut state) => {
1905 if state.final_sequence.is_some() || state.abort.is_some() {
1906 return;
1907 }
1908 tracing::warn!(
1909 track = %self.name,
1910 "track::Producer dropped without finish() or abort()"
1911 );
1912 state.resume = state.resume_position();
1914 state.clear_cache();
1915 state.datagrams.clear();
1916 }
1917 Err(state) => {
1918 if state.final_sequence.is_some() || state.abort.is_some() {
1919 return;
1920 }
1921 tracing::warn!(
1922 track = %self.name,
1923 "track::Producer dropped without finish() or abort()"
1924 );
1925 }
1926 }
1927 }
1928}
1929
1930fn combined_subscription(subs: &Subscriptions, bound: Option<Duration>, waiter: &kio::Waiter) -> Option<Subscription> {
1936 let mut combined = None;
1937 for sub in subs.iter() {
1938 if sub.is_closed() {
1943 continue;
1944 }
1945 let _ = sub.poll_closed(waiter);
1952 let _ = sub.poll(waiter, |_| Poll::<()>::Pending);
1953 if let Poll::Ready(merged) = sub.read().poll_combined(&combined) {
1954 combined = Some(merged);
1955 }
1956 }
1957 clamp_combined(combined, bound)
1958}
1959
1960fn snapshot_subscription(subs: &kio::Shared<Subscriptions>, bound: Option<Duration>) -> Option<Subscription> {
1962 let mut combined: Option<Subscription> = None;
1963 for sub in subs.read().iter() {
1964 if sub.is_closed() {
1966 continue;
1967 }
1968 if let Poll::Ready(merged) = sub.read().poll_combined(&combined) {
1969 combined = Some(merged);
1970 }
1971 }
1972 clamp_combined(combined, bound)
1973}
1974
1975fn servable_cap(cursor: Option<u64>, outer: Option<u64>) -> Option<u64> {
1990 super::subscription::min_some(cursor, outer)
1991}
1992
1993fn floor_of(subscription: &Subscription) -> u64 {
2002 subscription.start.map(|start| start.group).unwrap_or(0)
2003}
2004
2005fn clamp_max_age(mut max_age: Duration, bound: Option<Duration>) -> Duration {
2015 if let Some(bound) = bound {
2016 max_age = max_age.min(bound);
2017 }
2018 max_age
2019}
2020
2021fn clamp_combined(combined: Option<Subscription>, bound: Option<Duration>) -> Option<Subscription> {
2023 let mut combined = combined?;
2024 combined.max_age = clamp_max_age(combined.max_age, bound);
2025 Some(combined)
2026}
2027
2028fn register_subscription(state: kio::Ref<'_, TrackState>, subscription: &kio::Producer<Subscription>) {
2032 if state.is_closed() {
2033 return;
2034 }
2035 let subs = state.subscriptions.clone();
2036 drop(state);
2037 subs.lock().push(subscription.consume());
2038}
2039
2040#[derive(Clone)]
2042pub(crate) struct TrackWeak {
2043 name: Arc<str>,
2044 state: kio::ProducerWeak<TrackState>,
2045}
2046
2047impl TrackWeak {
2048 pub fn try_consume(&self) -> Option<Consumer> {
2055 Some(Consumer::plain(self.name.clone(), self.state.try_consume()?))
2056 }
2057
2058 pub(crate) fn name(&self) -> &Arc<str> {
2061 &self.name
2062 }
2063
2064 pub(crate) fn reject(&self, err: Error) -> bool {
2073 let Some(producer) = self.state.produce() else {
2074 return false;
2075 };
2076 let Ok(mut state) = producer.write() else {
2077 return false;
2078 };
2079 if state.published || state.claimed || state.abort.is_some() {
2080 return false;
2081 }
2082 state.abort = Some(err);
2083 state.close();
2084 true
2085 }
2086
2087 pub(crate) fn is_used(&self) -> bool {
2090 !self.state.is_closed() && self.state.is_used()
2091 }
2092
2093 pub(crate) fn poll_used(&self, waiter: &kio::Waiter) {
2096 let _ = self.state.poll_used(waiter);
2097 }
2098
2099 pub(crate) fn poll_unused(&self, waiter: &kio::Waiter) {
2102 let _ = self.state.poll_unused(waiter);
2103 }
2104}
2105
2106impl super::WeakEntry for TrackWeak {
2107 fn is_closed(&self) -> bool {
2108 self.state.is_closed()
2109 }
2110
2111 fn same_channel(&self, other: &Self) -> bool {
2112 self.state.same_channel(&other.state)
2113 }
2114}
2115
2116#[derive(Clone)]
2125pub struct Demand {
2126 name: Arc<str>,
2127 state: kio::ProducerWeak<TrackState>,
2128}
2129
2130impl Demand {
2131 pub fn name(&self) -> &str {
2133 &self.name
2134 }
2135
2136 pub async fn used(&self) -> Result<()> {
2138 self.state.used().await.map_err(|_| self.abort_reason())
2139 }
2140
2141 pub async fn unused(&self) -> Result<()> {
2143 self.state.unused().await.map_err(|_| self.abort_reason())
2144 }
2145
2146 pub async fn closed(&self) -> Error {
2148 self.state.closed().await;
2149 self.abort_reason()
2150 }
2151
2152 pub(crate) fn priority(&self) -> u8 {
2154 self.state.read().info.as_ref().map_or(0, |info| info.priority)
2156 }
2157
2158 pub fn is_used(&self) -> bool {
2160 self.state.is_used()
2161 }
2162
2163 pub fn poll_used(&self, waiter: &kio::Waiter) -> Poll<Result<()>> {
2165 self.state.poll_used(waiter).map(|used| match used {
2166 Some(()) => Ok(()),
2167 None => Err(self.abort_reason()),
2168 })
2169 }
2170
2171 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<Result<()>> {
2173 self.state.poll_unused(waiter).map(|used| match used {
2174 Some(()) => Ok(()),
2175 None => Err(self.abort_reason()),
2176 })
2177 }
2178
2179 pub(crate) fn is_closed(&self) -> bool {
2181 self.state.is_closed()
2182 }
2183
2184 pub(crate) fn poll_state(&self, waiter: &kio::Waiter) -> DemandState {
2191 loop {
2192 match self.state.poll_used(waiter) {
2193 Poll::Ready(None) => return DemandState::Closed,
2194 Poll::Pending => return DemandState::Idle,
2196 Poll::Ready(Some(())) => match self.state.poll_unused(waiter) {
2198 Poll::Ready(None) => return DemandState::Closed,
2199 Poll::Pending => return DemandState::Active,
2200 Poll::Ready(Some(())) => continue,
2203 },
2204 }
2205 }
2206 }
2207
2208 fn abort_reason(&self) -> Error {
2210 self.state.read().abort.clone().unwrap_or(Error::Dropped)
2211 }
2212}
2213
2214#[derive(Copy, Clone, Debug, Eq, PartialEq)]
2216pub(crate) enum DemandState {
2217 Active,
2219 Idle,
2221 Closed,
2223}
2224
2225#[derive(Clone)]
2236pub struct Consumer {
2237 name: Arc<str>,
2238 broadcast: Arc<broadcast::Info>,
2242 inner: ConsumerKind,
2243 stats: stats::Scope,
2246}
2247
2248#[derive(Clone)]
2249enum ConsumerKind {
2250 Plain(kio::Consumer<TrackState>),
2251 Spliced(super::resume::Consumer),
2252}
2253
2254impl Consumer {
2255 fn plain(name: Arc<str>, state: kio::Consumer<TrackState>) -> Self {
2256 let broadcast = state.read().broadcast.clone();
2257 Self {
2258 name,
2259 broadcast,
2260 inner: ConsumerKind::Plain(state),
2261 stats: stats::Scope::default(),
2262 }
2263 }
2264
2265 pub(crate) fn spliced(name: Arc<str>, broadcast: Arc<broadcast::Info>, resume: super::resume::Consumer) -> Self {
2267 Self {
2268 name,
2269 broadcast,
2270 inner: ConsumerKind::Spliced(resume),
2271 stats: stats::Scope::default(),
2272 }
2273 }
2274
2275 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
2278 self.stats = scope;
2279 self
2280 }
2281
2282 pub(crate) fn with_broadcast(mut self, broadcast: Arc<broadcast::Info>) -> Self {
2286 self.broadcast = broadcast;
2287 self
2288 }
2289
2290 pub(crate) fn cached_groups(&self) -> Vec<(group::Producer, bool)> {
2296 match &self.inner {
2297 ConsumerKind::Plain(state) => {
2298 let state = state.read();
2299 let mut out = Vec::with_capacity(state.lookup.len());
2300 for (sequence, stamp) in state.arrival.iter() {
2301 if let Some(slot) = state.lookup.get(sequence)
2302 && slot.stamp == *stamp
2303 && !slot.group.is_aborted()
2304 {
2305 out.push((slot.group.clone(), slot.visible));
2306 }
2307 }
2308 let mut copied: HashSet<u64> = out.iter().map(|(group, _)| group.sequence).collect();
2310 for (sequence, slot) in state.lookup.iter() {
2311 if !slot.group.is_aborted() && copied.insert(*sequence) {
2312 out.push((slot.group.clone(), slot.visible));
2313 }
2314 }
2315 out
2316 }
2317 ConsumerKind::Spliced(resume) => resume.cached_groups(),
2318 }
2319 }
2320
2321 pub(crate) fn cached_group(&self, sequence: u64, frame_start: u64) -> Option<group::Consumer> {
2326 match &self.inner {
2327 ConsumerKind::Plain(state) => {
2328 let state = state.read();
2329 let group = state.covering_group(sequence, frame_start)?;
2330 group.cache_refresh();
2331 Some(group.consume())
2332 }
2333 ConsumerKind::Spliced(resume) => resume.cached_group(sequence, frame_start),
2334 }
2335 }
2336
2337 pub(crate) fn cached_info(&self) -> Option<Info> {
2339 match &self.inner {
2340 ConsumerKind::Plain(state) => state.read().info.clone(),
2341 ConsumerKind::Spliced(resume) => resume.cached_info(),
2342 }
2343 }
2344
2345 pub fn name(&self) -> &str {
2347 &self.name
2348 }
2349
2350 pub fn broadcast(&self) -> &broadcast::Info {
2354 &self.broadcast
2355 }
2356
2357 pub fn subscribe(&self, subscription: impl Into<Option<Subscription>>) -> kio::Pending<Subscribing> {
2368 let subscription = kio::Producer::new(subscription.into().unwrap_or_default());
2369
2370 let inner = match &self.inner {
2371 ConsumerKind::Plain(state) => {
2372 register_subscription(state.read(), &subscription);
2375 SubscribingKind::Plain(state.clone())
2376 }
2377 ConsumerKind::Spliced(resume) => SubscribingKind::Spliced(resume.clone()),
2379 };
2380
2381 kio::Pending::new(Subscribing {
2382 name: self.name.clone(),
2383 broadcast: self.broadcast.clone(),
2384 inner,
2385 subscription,
2386 stats: self.stats.clone(),
2387 })
2388 }
2389
2390 pub(crate) fn poll_start(&self, waiter: &kio::Waiter) -> Poll<Option<u64>> {
2396 match &self.inner {
2397 ConsumerKind::Plain(state) => {
2398 let res = state.poll(waiter, |state| match state.start_pending && state.abort.is_none() {
2399 true => Poll::Pending,
2400 false => Poll::Ready(state.start_sequence),
2401 });
2402 match res {
2403 Poll::Ready(Ok(start)) => Poll::Ready(start),
2404 Poll::Ready(Err(state)) => Poll::Ready(state.start_sequence),
2405 Poll::Pending => Poll::Pending,
2406 }
2407 }
2408 ConsumerKind::Spliced(_) => Poll::Ready(None),
2409 }
2410 }
2411
2412 pub(crate) fn peek_latest(&self) -> Option<group::Consumer> {
2416 match &self.inner {
2417 ConsumerKind::Plain(state) => {
2418 let sequence = state.read().max_sequence?;
2419 self.peek_group(sequence)
2420 }
2421 ConsumerKind::Spliced(resume) => resume.peek_latest(),
2422 }
2423 }
2424
2425 pub(crate) fn live_edge(&self, cap: Option<u64>) -> Option<LiveEdge> {
2428 match &self.inner {
2429 ConsumerKind::Plain(state) => {
2430 let edge = state.read().live_edge(cap)?;
2431 Some(LiveEdge {
2432 sequence: edge.sequence,
2433 timestamp: edge.timestamp,
2434 stamp: edge.stamp,
2435 track: state.weak(),
2436 })
2437 }
2438 ConsumerKind::Spliced(resume) => resume.live_edge(cap),
2439 }
2440 }
2441
2442 pub(crate) fn served_start(&self, from: u64, cap: Option<u64>) -> Option<Successor> {
2445 match &self.inner {
2446 ConsumerKind::Plain(state) => {
2447 let served = state.read().served_start(from, cap)?;
2448 Some(Successor {
2449 sequence: served.sequence,
2450 timestamp: served.timestamp,
2451 stamp: served.stamp,
2452 track: state.weak(),
2453 })
2454 }
2455 ConsumerKind::Spliced(resume) => resume.served_start(from, cap),
2456 }
2457 }
2458
2459 pub(crate) fn peek_before(&self, sequence: u64) -> Option<group::Consumer> {
2463 match &self.inner {
2464 ConsumerKind::Plain(state) => {
2465 let state = state.read();
2466 state
2467 .lookup
2468 .range(..sequence)
2469 .rev()
2470 .map(|(_, slot)| &slot.group)
2471 .find(|group| !group.is_aborted())
2472 .map(|group| group.consume())
2473 }
2474 ConsumerKind::Spliced(resume) => resume.peek_before(sequence),
2475 }
2476 }
2477
2478 pub(crate) fn peek_group(&self, sequence: u64) -> Option<group::Consumer> {
2483 match &self.inner {
2484 ConsumerKind::Plain(state) => {
2485 let state = state.read();
2486 let slot = state.lookup.get(&sequence)?;
2487 if slot.group.is_aborted() {
2488 return None;
2489 }
2490 Some(slot.group.consume())
2491 }
2492 ConsumerKind::Spliced(resume) => resume.peek_group(sequence),
2493 }
2494 }
2495
2496 pub(crate) fn guard_group(
2498 &self,
2499 group: group::Consumer,
2500 subscription: kio::Consumer<Subscription>,
2501 anchor: kio::Consumer<Anchor>,
2502 bound: Option<u64>,
2503 ) -> group::Consumer {
2504 let ConsumerKind::Plain(state) = &self.inner else {
2505 return group;
2506 };
2507 let sequence = group.sequence;
2508 group.with_expiry(Arc::new(GroupExpiry {
2509 state: state.weak(),
2510 subscription,
2511 anchor,
2512 bound,
2513 sequence,
2514 }))
2515 }
2516
2517 pub(crate) fn poll_peek_group(&self, sequence: u64, waiter: &kio::Waiter) -> Poll<Option<group::Consumer>> {
2525 let ConsumerKind::Plain(state) = &self.inner else {
2526 return Poll::Pending;
2528 };
2529
2530 let res = state.poll(waiter, |state| {
2531 match state.lookup.get(&sequence) {
2532 Some(slot) if !slot.group.is_aborted() => Poll::Ready(Some(slot.group.consume())),
2533 Some(_) => Poll::Ready(None),
2535 None if state.final_sequence.is_some_and(|fin| sequence >= fin) => Poll::Ready(None),
2537 None if state.start_sequence.is_some_and(|start| sequence < start) => Poll::Ready(None),
2539 None => Poll::Pending,
2540 }
2541 });
2542
2543 match res {
2544 Poll::Ready(Ok(res)) => Poll::Ready(res),
2545 Poll::Ready(Err(_)) => Poll::Ready(None),
2547 Poll::Pending => Poll::Pending,
2548 }
2549 }
2550
2551 pub(crate) fn poll_serving_group(&self, sequence: u64, index: u64, waiter: &kio::Waiter) -> Poll<()> {
2566 let ConsumerKind::Plain(state) = &self.inner else {
2567 return Poll::Pending;
2569 };
2570 let res = state.poll(waiter, |state| match state.lookup.get(&sequence) {
2571 Some(slot) if !slot.group.is_aborted() => {
2572 let mut group = slot.group.consume();
2575 group.start_at(index);
2576 match group.index() == index {
2577 true => Poll::Ready(()),
2578 false => Poll::Pending,
2579 }
2580 }
2581 _ => Poll::Pending,
2582 });
2583 match res {
2584 Poll::Ready(Ok(())) => Poll::Ready(()),
2585 Poll::Ready(Err(_)) | Poll::Pending => Poll::Pending,
2588 }
2589 }
2590
2591 pub fn fetch_group(&self, sequence: u64, options: impl Into<Option<group::Fetch>>) -> kio::Pending<Fetching> {
2603 let options = options.into().unwrap_or_default();
2604
2605 self.stats.fetch();
2609
2610 let state = match &self.inner {
2611 ConsumerKind::Plain(state) => state,
2612 ConsumerKind::Spliced(resume) => {
2615 return kio::Pending::new(Fetching {
2616 inner: FetchingKind::Spliced(resume.fetch_group(sequence, options)),
2617 stats: self.stats.clone(),
2618 });
2619 }
2620 };
2621
2622 let mut result = None;
2623
2624 let (fetch, unresolved) = {
2628 let state = state.read();
2629 (
2630 state.fetch.clone(),
2631 state.poll_fetch_cached(sequence, options.frame_start).is_pending(),
2632 )
2633 };
2634
2635 if unresolved {
2636 let mut fetch = fetch.lock();
2637 if let Some(pending) = fetch.join(&sequence) {
2638 pending.priority = pending.priority.max(options.priority);
2648 pending.frame_start = pending.frame_start.min(options.frame_start);
2649 result = Some(pending.result.consume());
2650 } else {
2651 let producer = kio::Producer::<FetchOutcome>::default();
2655 let consumer = producer.consume();
2656 let attempt = PendingFetch {
2657 priority: options.priority,
2658 frame_start: options.frame_start,
2659 result: producer,
2660 };
2661 if fetch.insert(sequence, attempt).is_ok() {
2662 result = Some(consumer);
2663 }
2664 }
2665 }
2666
2667 kio::Pending::new(Fetching {
2668 inner: FetchingKind::Plain {
2669 state: state.clone(),
2670 fetch,
2671 sequence,
2672 frame_start: options.frame_start,
2673 result,
2674 },
2675 stats: self.stats.clone(),
2676 })
2677 }
2678
2679 pub fn query(&self) -> kio::Pending<Querying> {
2686 kio::Pending::new(Querying {
2687 inner: match &self.inner {
2688 ConsumerKind::Plain(state) => QueryingKind::Plain(state.clone()),
2689 ConsumerKind::Spliced(resume) => QueryingKind::Spliced(resume.clone()),
2690 },
2691 })
2692 }
2693
2694 pub fn latest(&self) -> Option<u64> {
2696 match &self.inner {
2697 ConsumerKind::Plain(state) => state.read().max_sequence,
2698 ConsumerKind::Spliced(resume) => resume.latest(),
2699 }
2700 }
2701
2702 pub(crate) fn resume_position(&self) -> Option<Position> {
2707 match &self.inner {
2708 ConsumerKind::Plain(state) => state.read().resume_position(),
2709 ConsumerKind::Spliced(resume) => resume.resume_position(),
2710 }
2711 }
2712
2713 pub(crate) fn poll_complete(&self, waiter: &kio::Waiter) -> Poll<Result<()>> {
2718 let ConsumerKind::Plain(state) = &self.inner else {
2719 return Poll::Pending;
2721 };
2722 match ready!(state.poll(waiter, |state| {
2723 if state.is_complete() {
2724 Poll::Ready(())
2725 } else {
2726 Poll::Pending
2727 }
2728 })) {
2729 Ok(_) => Poll::Ready(Ok(())),
2730 Err(closed) => Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped))),
2733 }
2734 }
2735}
2736
2737pub struct Subscribing {
2740 name: Arc<str>,
2741 broadcast: Arc<broadcast::Info>,
2742 inner: SubscribingKind,
2743 subscription: kio::Producer<Subscription>,
2744 stats: stats::Scope,
2745}
2746
2747enum SubscribingKind {
2748 Plain(kio::Consumer<TrackState>),
2749 Spliced(super::resume::Consumer),
2750}
2751
2752impl Subscribing {
2753 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<Subscriber>> {
2756 match &self.inner {
2757 SubscribingKind::Plain(state) => {
2758 let info = ready!(state.poll(waiter, |state| state.poll_info()))
2760 .map_err(|e| e.abort.clone().unwrap_or(Error::Dropped))??;
2761
2762 let drift_anchor = kio::Producer::new(Anchor::default());
2763 let min_sequence = floor_of(&self.subscription.read());
2764 Poll::Ready(Ok(Subscriber {
2765 name: self.name.clone(),
2766 broadcast: self.broadcast.clone(),
2767 info,
2768 inner: SubscriberKind::Plain(PlainSubscriber {
2769 state: state.clone(),
2770 subscription: self.subscription.clone(),
2771 min_sequence,
2772 index: 0,
2773 datagram_index: 0,
2774 next_sequence: 0,
2775 end_sequence: None,
2776 parked: BTreeMap::new(),
2777 stale_cap: None,
2778 drift_anchor,
2779 stale: stats::Content::default(),
2780 seek_pending: BTreeMap::new(),
2781 }),
2782 stats: self.stats.clone(),
2783 _stats_sub: self.stats.subscribe(),
2784 }))
2785 }
2786 SubscribingKind::Spliced(resume) => {
2787 let info = ready!(resume.poll_info(waiter))?;
2790
2791 Poll::Ready(Ok(Subscriber {
2792 name: self.name.clone(),
2793 broadcast: self.broadcast.clone(),
2794 info,
2795 inner: SubscriberKind::Spliced(Box::new(resume.subscribe_shared(self.subscription.clone()))),
2796 stats: self.stats.clone(),
2797 _stats_sub: self.stats.subscribe(),
2798 }))
2799 }
2800 }
2801 }
2802
2803 pub fn update(&mut self, subscription: Subscription) -> Result<()> {
2808 let mut state = self.subscription.write().map_err(|_| Error::Closed)?;
2809 *state = subscription;
2810 Ok(())
2811 }
2812}
2813
2814impl kio::Pollable for Subscribing {
2815 type Output = Result<Subscriber>;
2816
2817 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
2818 self.poll_ok(waiter)
2819 }
2820}
2821
2822pub struct Querying {
2825 inner: QueryingKind,
2826}
2827
2828enum QueryingKind {
2829 Plain(kio::Consumer<TrackState>),
2830 Spliced(super::resume::Consumer),
2831}
2832
2833impl Querying {
2834 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<Info>> {
2836 match &self.inner {
2837 QueryingKind::Plain(state) => {
2838 let info = ready!(state.poll(waiter, |state| state.poll_info()))
2840 .map_err(|e| e.abort.clone().unwrap_or(Error::Dropped))??;
2841 Poll::Ready(Ok(info))
2842 }
2843 QueryingKind::Spliced(resume) => resume.poll_info(waiter),
2844 }
2845 }
2846}
2847
2848impl kio::Pollable for Querying {
2849 type Output = Result<Info>;
2850
2851 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
2852 self.poll_ok(waiter)
2853 }
2854}
2855
2856impl group::Request {
2857 pub fn sequence(&self) -> u64 {
2859 self.sequence
2860 }
2861
2862 pub fn priority(&self) -> u8 {
2864 self.priority
2865 }
2866
2867 pub fn frame_start(&self) -> u64 {
2875 self.frame_start
2876 }
2877
2878 pub fn accept(mut self, info: impl Into<Option<Info>>) -> Result<group::Producer> {
2886 self.done = true;
2887 let res = TrackState::modify(&self.state)
2891 .and_then(|mut state| state.insert_group_request(self.sequence, self.frame_start, info.into()));
2892 self.remove();
2893 res
2894 }
2895
2896 pub fn reject(mut self, err: Error) {
2898 self.done = true;
2899 self.remove();
2902 if let Ok(mut outcome) = self.result.write() {
2903 outcome.rejected = Some(err);
2904 }
2905 }
2906
2907 fn remove(&self) {
2910 self.fetch
2911 .lock()
2912 .remove_if(&self.sequence, |pending| pending.result.same_channel(&self.result));
2913 }
2914}
2915
2916impl Drop for group::Request {
2917 fn drop(&mut self) {
2918 if self.done {
2919 return;
2920 }
2921 self.remove();
2922 if let Ok(mut outcome) = self.result.write() {
2923 outcome.rejected = Some(Error::Dropped);
2924 }
2925 }
2926}
2927
2928pub struct Fetching {
2934 inner: FetchingKind,
2935 stats: stats::Scope,
2938}
2939
2940enum FetchingKind {
2941 Plain {
2942 state: kio::Consumer<TrackState>,
2943 fetch: kio::Shared<FetchState>,
2944 sequence: u64,
2945 frame_start: u64,
2948 result: Option<kio::Consumer<FetchOutcome>>,
2950 },
2951 Spliced(kio::Pending<super::resume::Fetching>),
2953}
2954
2955impl kio::Pollable for Fetching {
2956 type Output = Result<group::Consumer>;
2957
2958 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
2959 let (state, fetch, sequence, frame_start, result) = match &self.inner {
2960 FetchingKind::Plain {
2961 state,
2962 fetch,
2963 sequence,
2964 frame_start,
2965 result,
2966 } => (state, fetch, *sequence, *frame_start, result.as_ref()),
2967 FetchingKind::Spliced(spliced) => {
2968 return kio::Pollable::poll(&**spliced, waiter)
2971 .map(|res| res.map(|group| group.with_meter(self.stats.meter())));
2972 }
2973 };
2974
2975 match state.poll(waiter, |state| state.poll_fetch_cached(sequence, frame_start)) {
2978 Poll::Ready(Ok(res)) => {
2979 return Poll::Ready(res.map(|mut group| {
2980 group.start_at(frame_start);
2984 group.with_meter(self.stats.meter())
2985 }));
2986 }
2987 Poll::Ready(Err(closed)) => {
2988 return Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped)));
2989 }
2990 Poll::Pending => {}
2991 }
2992
2993 let Some(result) = result else {
2995 return match fetch.poll(waiter, |fetch| match fetch.has_handlers() {
2998 false => Poll::Ready(()),
2999 true => Poll::Pending,
3000 }) {
3001 Poll::Ready(_guard) => Poll::Ready(Err(Error::NotFound)),
3002 Poll::Pending => Poll::Pending,
3003 };
3004 };
3005
3006 match result.poll(waiter, |outcome| match &outcome.rejected {
3009 Some(err) => Poll::Ready(err.clone()),
3010 None => Poll::Pending,
3011 }) {
3012 Poll::Ready(Ok(err)) => Poll::Ready(Err(err)),
3013 Poll::Ready(Err(_closed)) => Poll::Ready(Err(Error::NotFound)),
3014 Poll::Pending => Poll::Pending,
3015 }
3016 }
3017}
3018
3019pub struct Subscriber {
3047 name: Arc<str>,
3048 broadcast: Arc<broadcast::Info>,
3050 info: Info,
3051 inner: SubscriberKind,
3052 stats: stats::Scope,
3055 _stats_sub: stats::Subscription,
3058}
3059
3060enum SubscriberKind {
3061 Plain(PlainSubscriber),
3062 Spliced(Box<super::resume::Subscriber>),
3064}
3065
3066#[derive(Clone)]
3071struct Drift {
3072 budget: Duration,
3073 edge: Edge,
3074 outer: Option<LiveEdge>,
3075 successor: Option<Successor>,
3076}
3077
3078struct GroupExpiry {
3080 state: kio::ConsumerWeak<TrackState>,
3086 subscription: kio::Consumer<Subscription>,
3087 anchor: kio::Consumer<Anchor>,
3088 bound: Option<u64>,
3089 sequence: u64,
3090}
3091
3092impl group::Expiry for GroupExpiry {
3093 fn is_expired(&self, waiter: &kio::Waiter) -> bool {
3094 let mut max_age = Duration::default();
3095 let _ = self.subscription.poll(waiter, |subscription| {
3096 max_age = subscription.max_age;
3097 Poll::<()>::Pending
3098 });
3099
3100 let mut anchor = Anchor::default();
3101 let _ = self.anchor.poll(waiter, |current| {
3102 anchor = (**current).clone();
3103 Poll::<()>::Pending
3104 });
3105 let anchor = anchor.capped(self.bound);
3106 let cap = anchor.cap;
3107 let outer = anchor
3109 .edge
3110 .filter(LiveEdge::is_live)
3111 .map(|live| (live.sequence, live.timestamp));
3112 let successor = anchor.successor.as_ref().and_then(Successor::start);
3113
3114 let mut expired = false;
3115 let _ = self.state.poll(waiter, |state| {
3116 let budget = clamp_max_age(max_age, state.max_age_bound());
3117 loop {
3118 let edge = state.drift_edge(cap, outer, successor);
3119 expired = state.is_stale(self.sequence, &edge, budget);
3120 if expired {
3121 break;
3122 }
3123
3124 let mut timestamp_raced = false;
3130 if let Some(slot) = state.lookup.get(&self.sequence) {
3131 let group = &slot.group;
3132 if group.timestamp().is_none()
3133 && group.poll_timestamp(waiter).is_ready()
3134 && group.timestamp().is_some()
3135 {
3136 timestamp_raced = true;
3137 }
3138 }
3139 for (_, slot) in state
3140 .lookup
3141 .range((std::ops::Bound::Excluded(self.sequence), std::ops::Bound::Unbounded))
3142 {
3143 let group = &slot.group;
3144 if !super::subscription::before_end(group.sequence, cap) {
3145 break;
3146 }
3147 if slot.visible
3148 && !group.is_aborted()
3149 && group.timestamp().is_none()
3150 && group.poll_timestamp(waiter).is_ready()
3151 && group.timestamp().is_some()
3152 {
3153 timestamp_raced = true;
3154 break;
3155 }
3156 }
3157 if !timestamp_raced {
3158 break;
3159 }
3160 }
3161
3162 Poll::<()>::Pending
3165 });
3166
3167 expired
3168 }
3169}
3170
3171#[derive(Clone, Copy)]
3174struct Edge {
3175 presentation: Option<PresentationEdge>,
3177 outer: Option<(u64, Timestamp)>,
3180 cap: Option<u64>,
3183 successor: Option<Timestamp>,
3185}
3186
3187#[derive(Clone, Default, PartialEq)]
3190pub(crate) struct Anchor {
3191 pub cap: Option<u64>,
3193 pub edge: Option<LiveEdge>,
3198 pub successor: Option<Successor>,
3204}
3205
3206impl Anchor {
3207 pub fn capped(mut self, cap: Option<u64>) -> Self {
3210 let capped = servable_cap(cap, self.cap);
3211 if capped != self.cap {
3212 self.cap = capped;
3213 self.successor = None;
3214 }
3215 self
3216 }
3217}
3218
3219#[derive(Clone)]
3222pub(crate) struct LiveEdge {
3223 pub sequence: u64,
3224 pub timestamp: Timestamp,
3225 stamp: u32,
3226 track: kio::ConsumerWeak<TrackState>,
3227}
3228
3229impl LiveEdge {
3230 fn is_live(&self) -> bool {
3235 self.track.read().holds(self.sequence, self.stamp)
3236 }
3237}
3238
3239impl PartialEq for LiveEdge {
3240 fn eq(&self, other: &Self) -> bool {
3241 self.sequence == other.sequence
3242 && self.timestamp == other.timestamp
3243 && self.stamp == other.stamp
3244 && self.track.same_channel(&other.track)
3245 }
3246}
3247
3248struct ServedStart {
3251 sequence: u64,
3252 stamp: u32,
3253 timestamp: Timestamp,
3254}
3255
3256#[derive(Clone)]
3260pub(crate) struct Successor {
3261 sequence: u64,
3262 timestamp: Timestamp,
3263 stamp: u32,
3264 track: kio::ConsumerWeak<TrackState>,
3265}
3266
3267impl Successor {
3268 fn start(&self) -> Option<Timestamp> {
3272 let state = self.track.read();
3273 let slot = state.lookup.get(&self.sequence)?;
3274 if slot.stamp != self.stamp || slot.group.is_aborted() {
3275 return None;
3276 }
3277 slot.group.timestamp()
3278 }
3279}
3280
3281impl PartialEq for Successor {
3282 fn eq(&self, other: &Self) -> bool {
3283 self.sequence == other.sequence
3284 && self.timestamp == other.timestamp
3285 && self.stamp == other.stamp
3286 && self.track.same_channel(&other.track)
3287 }
3288}
3289
3290#[derive(Clone, Copy)]
3292struct PresentationEdge {
3293 sequence: u64,
3294 stamp: u32,
3297 timestamp: Timestamp,
3302}
3303
3304struct PlainSubscriber {
3306 state: kio::Consumer<TrackState>,
3307
3308 subscription: kio::Producer<Subscription>,
3309 index: usize,
3311 datagram_index: usize,
3313 min_sequence: u64,
3315 next_sequence: u64,
3318 end_sequence: Option<u64>,
3324 parked: BTreeMap<u64, group::Consumer>,
3329 stale_cap: Option<u64>,
3333 drift_anchor: kio::Producer<Anchor>,
3336 stale: stats::Content,
3342 seek_pending: BTreeMap<u64, stats::Content>,
3345}
3346
3347impl PlainSubscriber {
3348 fn anchor(&self, end: Option<u64>) -> Anchor {
3351 self.drift_anchor.read().clone().capped(end)
3352 }
3353
3354 fn update_drift_anchor(&mut self, outer: Anchor) {
3356 let anchor = outer.capped(self.end_sequence);
3357 if *self.drift_anchor.read() != anchor
3359 && let Ok(mut current) = self.drift_anchor.write()
3360 {
3361 *current = anchor;
3362 }
3363 }
3364
3365 fn poll<F, R>(&self, waiter: &kio::Waiter, f: F) -> Poll<Result<R>>
3367 where
3368 F: Fn(&kio::Ref<'_, TrackState>) -> Poll<Result<R>>,
3369 {
3370 Poll::Ready(match ready!(self.state.poll(waiter, f)) {
3371 Ok(res) => res,
3372 Err(state) => Err(state.abort.clone().unwrap_or(Error::Dropped)),
3374 })
3375 }
3376
3377 fn take_stale(&mut self) -> stats::Content {
3379 std::mem::take(&mut self.stale)
3380 }
3381
3382 fn note_stale(&mut self, group: &group::Consumer) {
3385 self.stale.add(group.content());
3386 }
3387
3388 fn poll_drift(&self, anchor: Anchor, waiter: &kio::Waiter) -> Poll<Result<Drift>> {
3399 let mut max_age = Duration::default();
3400 let _ = self.subscription.poll(waiter, |subscription| {
3401 max_age = subscription.max_age;
3402 Poll::<()>::Pending
3403 });
3404 let cap = anchor.cap;
3405 let outer = anchor.edge;
3406 let successor = anchor.successor;
3407 self.poll(waiter, |state| {
3408 Poll::Ready(Ok(Drift {
3411 budget: clamp_max_age(max_age, state.max_age_bound()),
3412 edge: state.drift_edge(cap, None, None),
3413 outer: outer.clone(),
3414 successor: successor.clone(),
3415 }))
3416 })
3417 }
3418
3419 fn poll_stale(&self, group: &group::Consumer, drift: &Drift, waiter: &kio::Waiter) -> Poll<Result<bool>> {
3422 let outer = drift
3425 .outer
3426 .as_ref()
3427 .filter(|live| live.is_live())
3428 .map(|live| (live.sequence, live.timestamp));
3429 let successor = drift.successor.as_ref().and_then(Successor::start);
3430 let presentation = drift.edge.presentation;
3431 let cap = drift.edge.cap;
3432 let budget = drift.budget;
3433 self.poll(waiter, move |state| {
3434 let edge = Edge {
3435 presentation,
3436 outer,
3437 cap,
3438 successor,
3439 };
3440 Poll::Ready(Ok(state.is_stale(group.sequence, &edge, budget)))
3441 })
3442 }
3443
3444 fn with_expiry(&self, group: group::Consumer) -> group::Consumer {
3445 let sequence = group.sequence;
3446 group.with_expiry(Arc::new(GroupExpiry {
3447 state: self.state.weak(),
3448 subscription: self.subscription.consume(),
3449 anchor: self.drift_anchor.consume(),
3450 bound: None,
3451 sequence,
3452 }))
3453 }
3454
3455 fn poll_recv_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
3456 let watch = |group: &group::Consumer| match group.poll_closed(waiter) {
3464 Poll::Pending => true,
3465 Poll::Ready(()) => !group.is_aborted(),
3466 };
3467
3468 let min_sequence = self.min_sequence;
3473 self.parked
3474 .retain(|sequence, group| *sequence >= min_sequence && watch(group));
3475
3476 let drift = ready!(self.poll_drift(self.anchor(self.end_sequence), waiter))?;
3478
3479 loop {
3480 let consumer = match self.parked.keys().next().copied() {
3483 Some(sequence) if super::subscription::before_end(sequence, self.end_sequence) => {
3484 let group = self.parked.remove(&sequence).expect("just looked it up");
3485 group.cache_refresh();
3487 group
3488 }
3489 _ => {
3490 let Some((producer, found_index)) =
3491 ready!(self.poll(waiter, |state| state.poll_recv_group(self.index, self.min_sequence))?)
3492 else {
3493 if self.parked.is_empty() {
3496 return Poll::Ready(Ok(None));
3497 }
3498 return Poll::Pending;
3499 };
3500 let consumer = producer.consume();
3501 consumer.cache_refresh();
3504 self.index = found_index + 1;
3505
3506 if !super::subscription::before_end(consumer.sequence, self.end_sequence) {
3509 if watch(&consumer) {
3513 self.parked.insert(consumer.sequence, consumer);
3514 }
3515 continue;
3516 }
3517 consumer
3518 }
3519 };
3520
3521 if ready!(self.poll_stale(&consumer, &drift, waiter))? {
3524 self.stale.add(consumer.content());
3525 continue;
3526 }
3527 return Poll::Ready(Ok(Some(self.with_expiry(consumer))));
3528 }
3529 }
3530
3531 fn poll_recv_datagram(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Datagram>>> {
3532 let Some((datagram, found_index)) =
3533 ready!(self.poll(waiter, |state| state.poll_recv_datagram(self.datagram_index))?)
3534 else {
3535 return Poll::Ready(Ok(None));
3536 };
3537
3538 self.datagram_index = found_index + 1;
3539 Poll::Ready(Ok(Some(datagram)))
3540 }
3541
3542 fn poll_next_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
3543 let floor = self.next_sequence.max(self.min_sequence);
3544 let Some(group) = ready!(self.poll_seek_group(floor, self.end_sequence, waiter))? else {
3545 return Poll::Ready(Ok(None));
3546 };
3547 self.next_sequence = group.sequence.saturating_add(1);
3548 self.commit_seek_stale(self.next_sequence);
3550 group.cache_refresh();
3552 Poll::Ready(Ok(Some(group)))
3553 }
3554
3555 fn poll_seek_group(
3568 &mut self,
3569 floor: u64,
3570 end: Option<u64>,
3571 waiter: &kio::Waiter,
3572 ) -> Poll<Result<Option<group::Consumer>>> {
3573 let mut floor = floor.max(self.min_sequence);
3574 let end = super::subscription::min_some(end, self.end_sequence);
3575 let drift = ready!(self.poll_drift(self.anchor(end), waiter))?;
3577
3578 loop {
3579 let Some(producer) = ready!(self.poll(waiter, |state| state.poll_next_in_range(floor, end))?) else {
3580 return Poll::Ready(Ok(None));
3587 };
3588 let group = producer.consume();
3589
3590 if ready!(self.poll_stale(&group, &drift, waiter))? {
3593 self.seek_pending.insert(group.sequence, group.content());
3598 floor = group.sequence.saturating_add(1);
3599 continue;
3600 }
3601
3602 self.seek_pending.remove(&group.sequence);
3605 return Poll::Ready(Ok(Some(self.with_expiry(group))));
3606 }
3607 }
3608
3609 fn commit_seek_stale(&mut self, committed: u64) {
3620 let keep = self.seek_pending.split_off(&committed);
3621 for (_, content) in std::mem::replace(&mut self.seek_pending, keep) {
3622 self.stale.add(content);
3623 }
3624 }
3625
3626 fn discard_seek_conviction(&mut self, sequence: u64) {
3630 self.seek_pending.remove(&sequence);
3631 }
3632}
3633
3634#[derive(Clone)]
3640pub struct Control {
3641 subscription: kio::Producer<Subscription>,
3642}
3643
3644impl Control {
3645 pub fn subscription(&self) -> Subscription {
3647 self.subscription.read().clone()
3648 }
3649
3650 pub fn update(&self, subscription: Subscription) -> Result<()> {
3655 let mut state = self.subscription.write().map_err(|_| Error::Closed)?;
3656 *state = subscription;
3657 Ok(())
3658 }
3659}
3660
3661impl Subscriber {
3662 pub fn info(&self) -> &Info {
3667 &self.info
3668 }
3669
3670 pub fn name(&self) -> &str {
3672 &self.name
3673 }
3674
3675 pub fn broadcast(&self) -> &broadcast::Info {
3679 &self.broadcast
3680 }
3681
3682 fn count_stale(&mut self, meter: &stats::Meter) {
3688 if meter.is_tracked() {
3693 meter.stale(self.take_stale());
3694 }
3695 }
3696
3697 pub(crate) fn take_stale(&mut self) -> stats::Content {
3700 match &mut self.inner {
3701 SubscriberKind::Plain(plain) => plain.take_stale(),
3702 SubscriberKind::Spliced(spliced) => spliced.take_stale(),
3703 }
3704 }
3705
3706 pub(crate) fn set_anchor(&mut self, anchor: Anchor) {
3717 match &mut self.inner {
3718 SubscriberKind::Plain(plain) => {
3719 plain.stale_cap = anchor.cap;
3720 plain.update_drift_anchor(anchor);
3721 }
3722 SubscriberKind::Spliced(spliced) => spliced.set_anchor(anchor),
3723 }
3724 }
3725
3726 pub(crate) fn commit_seek_stale(&mut self, committed: u64) {
3733 match &mut self.inner {
3734 SubscriberKind::Plain(plain) => plain.commit_seek_stale(committed),
3735 SubscriberKind::Spliced(spliced) => spliced.commit_seek_stale(committed),
3736 }
3737 }
3738
3739 pub(crate) fn discard_seek_conviction(&mut self, sequence: u64) {
3744 match &mut self.inner {
3745 SubscriberKind::Plain(plain) => plain.discard_seek_conviction(sequence),
3746 SubscriberKind::Spliced(spliced) => spliced.discard_seek_conviction(sequence),
3747 }
3748 }
3749
3750 pub(crate) fn poll_stale(&mut self, group: &group::Consumer, waiter: &kio::Waiter) -> Poll<Result<bool>> {
3756 let plain = match &mut self.inner {
3757 SubscriberKind::Plain(plain) => plain,
3758 SubscriberKind::Spliced(spliced) => return spliced.poll_stale(group, waiter),
3759 };
3760 let drift = ready!(plain.poll_drift(plain.anchor(plain.end_sequence), waiter))?;
3761 let stale = ready!(plain.poll_stale(group, &drift, waiter))?;
3762 if stale {
3763 plain.note_stale(group);
3764 }
3765 Poll::Ready(Ok(stale))
3766 }
3767
3768 pub fn control(&self) -> Control {
3770 Control {
3771 subscription: match &self.inner {
3772 SubscriberKind::Plain(plain) => plain.subscription.clone(),
3773 SubscriberKind::Spliced(spliced) => spliced.prefs(),
3774 },
3775 }
3776 }
3777
3778 pub fn poll_recv_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
3807 let meter = self.stats.meter();
3808 let res = match &mut self.inner {
3809 SubscriberKind::Plain(plain) => plain.poll_recv_group(waiter),
3810 SubscriberKind::Spliced(spliced) => spliced.poll_recv_group(waiter),
3811 };
3812 self.count_stale(&meter);
3813 res.map(|res| res.map(|group| group.map(|group| group.with_meter(meter))))
3814 }
3815
3816 pub async fn recv_group(&mut self) -> Result<Option<group::Consumer>> {
3823 kio::wait(|waiter| self.poll_recv_group(waiter)).await
3824 }
3825
3826 pub fn poll_recv_datagram(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Datagram>>> {
3837 let meter = self.stats.meter();
3838 let res = match &mut self.inner {
3839 SubscriberKind::Plain(plain) => plain.poll_recv_datagram(waiter),
3840 SubscriberKind::Spliced(spliced) => spliced.poll_recv_datagram(waiter),
3841 };
3842 if let Poll::Ready(Ok(Some(datagram))) = &res {
3845 meter.datagram(datagram.payload.len() as u64);
3846 }
3847 res
3848 }
3849
3850 pub async fn recv_datagram(&mut self) -> Result<Option<Datagram>> {
3857 kio::wait(|waiter| self.poll_recv_datagram(waiter)).await
3858 }
3859
3860 fn poll_next_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
3862 let meter = self.stats.meter();
3863 let res = match &mut self.inner {
3864 SubscriberKind::Plain(plain) => plain.poll_next_group(waiter),
3865 SubscriberKind::Spliced(spliced) => spliced.poll_next_group(waiter),
3866 };
3867 self.count_stale(&meter);
3868 res.map(|res| res.map(|group| group.map(|group| group.with_meter(meter))))
3869 }
3870
3871 pub(crate) fn poll_seek_group(
3876 &mut self,
3877 floor: u64,
3878 end: Option<u64>,
3879 waiter: &kio::Waiter,
3880 ) -> Poll<Result<Option<group::Consumer>>> {
3881 let meter = self.stats.meter();
3882 let res = match &mut self.inner {
3883 SubscriberKind::Plain(plain) => plain.poll_seek_group(floor, end, waiter),
3884 SubscriberKind::Spliced(spliced) => spliced.poll_seek_group(floor, end, waiter),
3885 };
3886 self.count_stale(&meter);
3887 res.map(|res| res.map(|group| group.map(|group| group.with_meter(meter))))
3888 }
3889
3890 pub fn ordered(self) -> Ordered {
3900 Ordered { inner: self }
3901 }
3902
3903 pub fn is_clone(&self, other: &Self) -> bool {
3905 match (&self.inner, &other.inner) {
3906 (SubscriberKind::Plain(a), SubscriberKind::Plain(b)) => a.state.same_channel(&b.state),
3907 (SubscriberKind::Spliced(a), SubscriberKind::Spliced(b)) => a.is_clone(b),
3908 _ => false,
3909 }
3910 }
3911
3912 pub fn poll_finished(&mut self, waiter: &kio::Waiter) -> Poll<Result<u64>> {
3914 match &mut self.inner {
3915 SubscriberKind::Plain(plain) => plain.poll(waiter, |state| state.poll_finished()),
3916 SubscriberKind::Spliced(spliced) => spliced.poll_finished(waiter),
3917 }
3918 }
3919
3920 pub async fn finished(&mut self) -> Result<u64> {
3928 kio::wait(|waiter| self.poll_finished(waiter)).await
3929 }
3930
3931 pub fn set_groups(&mut self, groups: impl RangeBounds<u64>) {
3939 let (start, end) = super::subscription::sequence_bounds(groups);
3940 self.raise_start_to(start);
3941 self.end_at(end.map_or(Bound::Unbounded, Bound::Excluded));
3942 }
3943
3944 pub(crate) fn start_at(&mut self, sequence: u64) {
3951 match &mut self.inner {
3952 SubscriberKind::Plain(plain) => plain.min_sequence = sequence,
3953 SubscriberKind::Spliced(spliced) => spliced.start_at(sequence),
3954 }
3955 }
3956
3957 pub(crate) fn raise_start_to(&mut self, sequence: u64) {
3963 match &mut self.inner {
3964 SubscriberKind::Plain(plain) => plain.min_sequence = plain.min_sequence.max(sequence),
3965 SubscriberKind::Spliced(spliced) => spliced.raise_start_to(sequence),
3966 }
3967 }
3968
3969 pub(crate) fn end_at(&mut self, end: impl Into<Cap>) {
3983 let end = end.into();
3984 match &mut self.inner {
3985 SubscriberKind::Plain(plain) => {
3986 plain.end_sequence = end.exclusive();
3987 let outer = Anchor {
3990 cap: plain.stale_cap,
3991 ..plain.drift_anchor.read().clone()
3992 };
3993 plain.update_drift_anchor(outer);
3994 }
3995 SubscriberKind::Spliced(spliced) => spliced.end_at(end),
3996 }
3997 }
3998
3999 pub fn subscription(&self) -> Subscription {
4001 self.control().subscription()
4002 }
4003
4004 pub fn update(&mut self, subscription: Subscription) -> Result<()> {
4010 match &mut self.inner {
4011 SubscriberKind::Plain(plain) => {
4012 let mut state = plain.subscription.write().map_err(|_| Error::Closed)?;
4013 *state = subscription;
4014 }
4015 SubscriberKind::Spliced(spliced) => spliced.update(subscription),
4016 }
4017 Ok(())
4018 }
4019
4020 pub fn latest(&self) -> Option<u64> {
4022 match &self.inner {
4023 SubscriberKind::Plain(plain) => plain.state.read().max_sequence,
4024 SubscriberKind::Spliced(spliced) => spliced.latest(),
4025 }
4026 }
4027}
4028
4029pub struct Ordered {
4056 inner: Subscriber,
4057}
4058
4059impl Ordered {
4060 pub fn info(&self) -> &Info {
4062 self.inner.info()
4063 }
4064
4065 pub fn name(&self) -> &str {
4067 self.inner.name()
4068 }
4069
4070 pub fn broadcast(&self) -> &broadcast::Info {
4072 self.inner.broadcast()
4073 }
4074
4075 pub fn poll_next_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
4087 self.inner.poll_next_group(waiter)
4088 }
4089
4090 pub async fn next_group(&mut self) -> Result<Option<group::Consumer>> {
4092 kio::wait(|waiter| self.poll_next_group(waiter)).await
4093 }
4094
4095 pub fn poll_recv_datagram(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Datagram>>> {
4103 self.inner.poll_recv_datagram(waiter)
4104 }
4105
4106 pub async fn recv_datagram(&mut self) -> Result<Option<Datagram>> {
4113 kio::wait(|waiter| self.poll_recv_datagram(waiter)).await
4114 }
4115
4116 pub fn set_groups(&mut self, groups: impl RangeBounds<u64>) {
4120 self.inner.set_groups(groups);
4121 }
4122
4123 pub fn control(&self) -> Control {
4126 self.inner.control()
4127 }
4128
4129 pub fn subscription(&self) -> Subscription {
4131 self.inner.subscription()
4132 }
4133
4134 pub fn update(&mut self, subscription: Subscription) -> Result<()> {
4138 self.inner.update(subscription)
4139 }
4140
4141 pub fn poll_finished(&mut self, waiter: &kio::Waiter) -> Poll<Result<u64>> {
4143 self.inner.poll_finished(waiter)
4144 }
4145
4146 pub async fn finished(&mut self) -> Result<u64> {
4151 kio::wait(|waiter| self.poll_finished(waiter)).await
4152 }
4153
4154 pub fn latest(&self) -> Option<u64> {
4156 self.inner.latest()
4157 }
4158
4159 pub fn is_clone(&self, other: &Self) -> bool {
4161 self.inner.is_clone(&other.inner)
4162 }
4163}
4164
4165pub struct Request {
4177 name: Arc<str>,
4178 broadcast: Arc<broadcast::Info>,
4180 state: kio::Producer<TrackState>,
4181
4182 prev_subscription: Option<Subscription>,
4184
4185 alive: Arc<Alive>,
4188
4189 _dynamic: Dynamic,
4194
4195 stats: stats::Scope,
4198
4199 resolving_start: bool,
4202}
4203
4204impl Request {
4205 pub(crate) fn new(broadcast: Arc<broadcast::Info>, name: impl Into<Arc<str>>) -> Self {
4206 let name = name.into();
4207 let state = TrackState::spawn(broadcast.clone());
4208 let alive = Alive::new(name.clone(), state.clone());
4209 let dynamic = Dynamic::new(name.clone(), state.clone(), alive.clone());
4210 Self {
4211 name,
4212 broadcast,
4213 state,
4214 prev_subscription: None,
4215 alive,
4216 _dynamic: dynamic,
4217 stats: stats::Scope::default(),
4218 resolving_start: false,
4219 }
4220 }
4221
4222 pub(crate) fn resolving_start(mut self) -> Self {
4227 self.resolving_start = true;
4228 self
4229 }
4230
4231 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
4234 self.stats = scope;
4235 self
4236 }
4237
4238 pub fn name(&self) -> &str {
4240 &self.name
4241 }
4242
4243 pub fn consume(&self) -> Consumer {
4245 Consumer::plain(self.name.clone(), self.state.consume())
4246 }
4247
4248 pub fn dynamic(&self) -> Dynamic {
4252 Dynamic::new(self.name.clone(), self.state.clone(), self.alive.clone())
4253 }
4254
4255 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
4258 self.state.poll_unused(waiter).map(|_| ())
4259 }
4260
4261 pub(crate) fn claim(self) -> Self {
4263 if let Ok(mut state) = self.state.write() {
4264 state.claimed = true;
4265 }
4266 self
4267 }
4268
4269 pub(crate) fn reject_unused(&self, err: Error) -> bool {
4273 match self.state.write_unused() {
4274 kio::Unused::Idle(guard) => {
4275 commit_abort(guard, err);
4276 true
4277 }
4278 kio::Unused::Closed => true,
4279 kio::Unused::Used => false,
4280 }
4281 }
4282
4283 pub fn accept(self, info: impl Into<Option<Info>>) -> Producer {
4290 let info = TrackState::normalize_info(&self.broadcast, info.into().unwrap_or_default());
4291 if let Ok(mut state) = self.state.write() {
4294 state.accept(info.clone());
4295 state.start_pending = self.resolving_start;
4296 }
4297 self.alive.publish(Some(&self.stats));
4300 Producer {
4301 name: self.name,
4302 info,
4303 broadcast: self.broadcast,
4304 state: self.state,
4305 prev_subscription: None,
4306 alive: self.alive,
4307 stats: self.stats,
4308 }
4309 }
4310
4311 pub fn reject(self, err: Error) {
4313 if let Ok(mut state) = self.state.write() {
4314 state.abort = Some(err);
4315 }
4316 }
4317
4318 pub fn subscription(&self) -> Option<Subscription> {
4321 let state = self.state.read();
4322 let (subs, bound) = (state.subscriptions.clone(), state.max_age_bound());
4323 drop(state);
4324 snapshot_subscription(&subs, bound)
4325 }
4326
4327 pub async fn subscription_changed(&mut self) -> Option<Subscription> {
4330 kio::wait(|waiter| self.poll_subscription_changed(waiter)).await
4331 }
4332
4333 pub fn poll_subscription_changed(&mut self, waiter: &kio::Waiter) -> Poll<Option<Subscription>> {
4335 let state = self.state.read();
4336 let (subs, bound) = (state.subscriptions.clone(), state.max_age_bound());
4337 drop(state);
4338
4339 let prev = &self.prev_subscription;
4340 let mut combined = None;
4341 let mut guard = ready!(subs.poll(waiter, |subs| {
4342 let next = combined_subscription(subs, bound, waiter);
4343 if &next == prev {
4344 Poll::Pending
4345 } else {
4346 combined = next;
4347 Poll::Ready(())
4348 }
4349 }));
4350 guard.retain(|sub| !sub.is_closed());
4352 drop(guard);
4353 self.prev_subscription = combined.clone();
4354 Poll::Ready(combined)
4355 }
4356
4357 pub(super) fn weak(&self) -> TrackWeak {
4358 TrackWeak {
4359 name: self.name.clone(),
4360 state: self.state.weak(),
4361 }
4362 }
4363}
4364
4365#[cfg(test)]
4366use futures::FutureExt;
4367
4368#[cfg(test)]
4369#[allow(missing_docs)] impl Subscriber {
4371 pub fn assert_group(&mut self) -> group::Consumer {
4372 self.recv_group()
4373 .now_or_never()
4374 .expect("group would have blocked")
4375 .expect("would have errored")
4376 .expect("track was closed")
4377 }
4378
4379 pub fn assert_no_group(&mut self) {
4380 assert!(
4381 self.recv_group().now_or_never().is_none(),
4382 "recv_group would not have blocked"
4383 );
4384 }
4385
4386 pub fn assert_not_closed(&mut self) {
4387 assert!(self.finished().now_or_never().is_none(), "should not be closed");
4388 }
4389
4390 pub fn assert_closed(&mut self) {
4391 assert!(self.finished().now_or_never().is_some(), "should be closed");
4392 }
4393
4394 pub fn assert_error(&mut self) {
4396 assert!(
4397 self.finished().now_or_never().expect("should not block").is_err(),
4398 "should be error"
4399 );
4400 }
4401
4402 pub fn assert_is_clone(&self, other: &Self) {
4403 assert!(self.is_clone(other), "should be clone");
4404 }
4405
4406 pub fn assert_not_clone(&self, other: &Self) {
4407 assert!(!self.is_clone(other), "should not be clone");
4408 }
4409}
4410
4411#[cfg(test)]
4412mod test {
4413 use super::*;
4414 use crate::frame;
4415 use crate::model::test_tracing::count_drop_warnings;
4416 use std::time::Duration;
4417
4418 fn track_producer(name: impl Into<Arc<str>>, info: impl Into<Option<Info>>) -> Producer {
4421 Producer::new(Arc::new(broadcast::Info::default()), name, info)
4422 }
4423
4424 fn replay() -> Subscription {
4426 Subscription::default().with_max_age(Duration::from_secs(30))
4427 }
4428
4429 fn live_groups(state: &TrackState) -> usize {
4431 state.lookup.len()
4432 }
4433
4434 fn first_live_sequence(state: &TrackState) -> u64 {
4436 state
4437 .arrival
4438 .iter()
4439 .find(|(sequence, stamp)| state.lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp))
4440 .map(|(sequence, _)| *sequence)
4441 .unwrap()
4442 }
4443
4444 fn recv_datagram(dg: &mut Subscriber) -> Datagram {
4446 dg.recv_datagram()
4447 .now_or_never()
4448 .expect("datagram would have blocked")
4449 .expect("would have errored")
4450 .expect("track was closed")
4451 }
4452
4453 #[tokio::test]
4457 async fn peek_resolves_below_the_declared_start() {
4458 let mut producer = track_producer("test", None);
4459 let consumer = producer.consume();
4460
4461 let mut cached = producer.create_group(group::Info { sequence: 1 }).unwrap();
4462 cached.write_frame(Timestamp::ZERO, b"backfill".to_vec()).unwrap();
4463 cached.finish().unwrap();
4464
4465 let waiter = kio::Waiter::noop();
4467 assert!(consumer.poll_peek_group(0, &waiter).is_pending());
4468
4469 producer.start_at(3).unwrap();
4472 assert!(matches!(consumer.poll_peek_group(0, &waiter), Poll::Ready(None)));
4473 assert!(matches!(consumer.poll_peek_group(1, &waiter), Poll::Ready(Some(_))));
4474 assert!(consumer.poll_peek_group(3, &waiter).is_pending());
4475
4476 producer.start_at(4).unwrap();
4479 assert!(matches!(consumer.poll_peek_group(3, &waiter), Poll::Ready(None)));
4480 producer.start_at(0).unwrap();
4481 assert!(consumer.poll_peek_group(0, &waiter).is_pending());
4482 }
4483
4484 #[tokio::test]
4485 async fn append_datagram_shares_group_sequence() {
4486 let mut producer = track_producer("test", None);
4487 let ts = Timestamp::from_millis(10).unwrap();
4488
4489 assert_eq!(producer.append_group().unwrap().sequence, 0);
4491 assert_eq!(producer.append_datagram(ts, &b"a"[..]).unwrap(), 1);
4492 assert_eq!(producer.append_group().unwrap().sequence, 2);
4493 assert_eq!(producer.append_datagram(ts, &b"b"[..]).unwrap(), 3);
4494 assert_eq!(producer.latest(), Some(3));
4495 }
4496
4497 #[tokio::test]
4498 async fn append_datagram_roundtrip() {
4499 let mut producer = track_producer("test", None);
4500 let mut dg = producer.subscribe(None);
4501
4502 let ts = Timestamp::from_millis(42).unwrap();
4503 let seq = producer.append_datagram(ts, &b"hello"[..]).unwrap();
4504
4505 let got = recv_datagram(&mut dg);
4506 assert_eq!(got.sequence, seq);
4507 assert_eq!(got.timestamp, ts);
4508 assert_eq!(&got.payload[..], b"hello");
4509 }
4510
4511 #[tokio::test]
4512 async fn insert_datagram_preserves_sequence() {
4513 let mut producer = track_producer("test", None);
4514 let mut dg = producer.subscribe(None);
4515
4516 let ts = Timestamp::from_millis(5).unwrap();
4517 producer
4519 .insert_datagram(100, ts, bytes::Bytes::from_static(b"x"))
4520 .unwrap();
4521
4522 assert_eq!(recv_datagram(&mut dg).sequence, 100);
4523 assert_eq!(producer.append_group().unwrap().sequence, 101);
4525 }
4526
4527 #[tokio::test]
4528 async fn insert_datagram_leaves_a_gap() {
4529 let mut producer = track_producer("test", None);
4530 let mut dg = producer.subscribe(None);
4531 let ts = Timestamp::from_millis(0).unwrap();
4532
4533 producer
4534 .insert_datagram(10, ts, bytes::Bytes::from_static(b"gap"))
4535 .unwrap();
4536 assert_eq!(recv_datagram(&mut dg).sequence, 10);
4537 assert_eq!(producer.append_datagram(ts, &b"next"[..]).unwrap(), 11);
4538 assert_eq!(producer.append_group().unwrap().sequence, 12);
4539 }
4540
4541 #[tokio::test]
4542 async fn insert_datagram_out_of_order_does_not_rewind() {
4543 let mut producer = track_producer("test", None);
4544 let mut dg = producer.subscribe(None);
4545 let ts = Timestamp::from_millis(0).unwrap();
4546
4547 producer
4548 .insert_datagram(10, ts, bytes::Bytes::from_static(b"high"))
4549 .unwrap();
4550 producer
4551 .insert_datagram(5, ts, bytes::Bytes::from_static(b"low"))
4552 .unwrap();
4553
4554 assert_eq!(recv_datagram(&mut dg).sequence, 10);
4555 assert_eq!(recv_datagram(&mut dg).sequence, 5);
4556 assert_eq!(producer.append_datagram(ts, &b"next"[..]).unwrap(), 11);
4557 }
4558
4559 #[tokio::test]
4560 async fn insert_datagram_duplicate_is_best_effort() {
4561 let mut producer = track_producer("test", None);
4562 let mut dg = producer.subscribe(None);
4563 let ts = Timestamp::from_millis(0).unwrap();
4564
4565 producer
4566 .insert_datagram(3, ts, bytes::Bytes::from_static(b"first"))
4567 .unwrap();
4568 producer
4569 .insert_datagram(3, ts, bytes::Bytes::from_static(b"again"))
4570 .unwrap();
4571
4572 assert_eq!(&recv_datagram(&mut dg).payload[..], b"first");
4573 assert_eq!(&recv_datagram(&mut dg).payload[..], b"again");
4574 assert_eq!(producer.append_datagram(ts, &b"next"[..]).unwrap(), 4);
4575 }
4576
4577 #[tokio::test]
4578 async fn insert_datagram_stale_does_not_rewind_after_append() {
4579 let mut producer = track_producer("test", None);
4580 let mut dg = producer.subscribe(None);
4581 let ts = Timestamp::from_millis(0).unwrap();
4582
4583 assert_eq!(producer.append_datagram(ts, &b"0"[..]).unwrap(), 0);
4584 assert_eq!(producer.append_datagram(ts, &b"1"[..]).unwrap(), 1);
4585 producer
4586 .insert_datagram(0, ts, bytes::Bytes::from_static(b"stale"))
4587 .unwrap();
4588
4589 assert_eq!(recv_datagram(&mut dg).sequence, 0);
4590 assert_eq!(recv_datagram(&mut dg).sequence, 1);
4591 assert_eq!(recv_datagram(&mut dg).sequence, 0);
4592 assert_eq!(producer.append_datagram(ts, &b"2"[..]).unwrap(), 2);
4593 assert_eq!(producer.append_group().unwrap().sequence, 3);
4594 }
4595
4596 #[tokio::test]
4597 async fn insert_datagram_cloned_producers_share_counter() {
4598 let mut producer = track_producer("test", None);
4599 let mut other = producer.clone();
4600 let mut dg = producer.subscribe(None);
4601 let ts = Timestamp::from_millis(0).unwrap();
4602
4603 producer
4604 .insert_datagram(4, ts, bytes::Bytes::from_static(b"a"))
4605 .unwrap();
4606 assert_eq!(other.append_datagram(ts, &b"b"[..]).unwrap(), 5);
4607 other.insert_datagram(8, ts, bytes::Bytes::from_static(b"c")).unwrap();
4608 assert_eq!(producer.append_group().unwrap().sequence, 9);
4609
4610 assert_eq!(recv_datagram(&mut dg).sequence, 4);
4611 assert_eq!(recv_datagram(&mut dg).sequence, 5);
4612 assert_eq!(recv_datagram(&mut dg).sequence, 8);
4613 }
4614
4615 #[test]
4616 fn insert_datagram_after_finish_is_closed() {
4617 let mut producer = track_producer("test", None);
4618 let ts = Timestamp::from_millis(0).unwrap();
4619 producer.finish().unwrap();
4620 assert!(matches!(
4621 producer.insert_datagram(0, ts, bytes::Bytes::from_static(b"x")),
4622 Err(Error::Closed)
4623 ));
4624 assert!(matches!(producer.append_datagram(ts, &b"x"[..]), Err(Error::Closed)));
4625 }
4626
4627 #[test]
4628 fn insert_datagram_after_abort_fails() {
4629 let producer = track_producer("test", None);
4630 let mut other = producer.clone();
4631 let ts = Timestamp::from_millis(0).unwrap();
4632 producer.abort(Error::Cancel).unwrap();
4633 assert!(other.insert_datagram(0, ts, bytes::Bytes::from_static(b"x")).is_err());
4634 }
4635
4636 #[test]
4637 fn insert_datagram_respects_finish_at() {
4638 let mut producer = track_producer("test", None);
4639 let ts = Timestamp::from_millis(0).unwrap();
4640 producer.finish_at(10).unwrap();
4641 producer
4642 .insert_datagram(5, ts, bytes::Bytes::from_static(b"ok"))
4643 .unwrap();
4644 assert!(matches!(
4645 producer.insert_datagram(10, ts, bytes::Bytes::from_static(b"late")),
4646 Err(Error::Closed)
4647 ));
4648 assert_eq!(producer.append_group().unwrap().sequence, 6);
4649 }
4650
4651 #[test]
4653 fn resume_position_uses_the_latest_group() {
4654 let mut datagram_only = track_producer("datagram-only", None);
4655 let datagram_only_consumer = datagram_only.consume();
4656 datagram_only
4657 .insert_datagram(8, Timestamp::ZERO, bytes::Bytes::from_static(b"x"))
4658 .unwrap();
4659 assert_eq!(
4660 datagram_only_consumer.resume_position(),
4661 None,
4662 "a datagram creates no group position to resume"
4663 );
4664
4665 let mut producer = track_producer("mixed", None);
4666 let consumer = producer.consume();
4667 let mut group = producer.create_group(group::Info { sequence: 3 }).unwrap();
4668 group
4669 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"head"))
4670 .unwrap();
4671 producer
4672 .insert_datagram(8, Timestamp::ZERO, bytes::Bytes::from_static(b"x"))
4673 .unwrap();
4674
4675 assert_eq!(
4676 consumer.resume_position(),
4677 Some(Position { group: 3, frame: 1 }),
4678 "the replacement must continue the open group"
4679 );
4680 group.abort(Error::Cancel).unwrap();
4681 }
4682
4683 #[tokio::test]
4686 async fn recv_datagram_leaves_the_ordered_cursor_alone() {
4687 let mut producer = track_producer("test", None);
4688 let mut datagrams = producer.subscribe(None);
4689 let mut subscriber = producer.subscribe(None).ordered();
4690 let ts = Timestamp::from_millis(5).unwrap();
4691
4692 producer
4693 .insert_datagram(5, ts, bytes::Bytes::from_static(b"x"))
4694 .unwrap();
4695 assert_eq!(recv_datagram(&mut datagrams).sequence, 5);
4696
4697 producer.create_group(group::Info { sequence: 3 }).unwrap();
4698 producer.create_group(group::Info { sequence: 6 }).unwrap();
4699
4700 let mut next = || {
4701 subscriber
4702 .next_group()
4703 .now_or_never()
4704 .expect("group would have blocked")
4705 .expect("would have errored")
4706 .expect("track was closed")
4707 .sequence
4708 };
4709 assert_eq!(next(), 3, "the datagram at sequence 5 did not consume group 3");
4710 assert_eq!(next(), 6);
4711 }
4712
4713 #[tokio::test]
4714 async fn datagram_normalized_to_track_timescale() {
4715 let info = Info::default().with_timescale(Timescale::MICRO);
4716 let mut producer = track_producer("test", info);
4717 let mut dg = producer.subscribe(None);
4718
4719 producer
4721 .append_datagram(Timestamp::from_millis(2).unwrap(), &b"z"[..])
4722 .unwrap();
4723 let got = recv_datagram(&mut dg);
4724 assert_eq!(got.timestamp.scale(), Timescale::MICRO);
4725 assert_eq!(got.timestamp.value(), 2_000);
4726 }
4727
4728 #[tokio::test]
4729 async fn datagram_rejects_oversized() {
4730 let mut producer = track_producer("test", None);
4731 let big = bytes::Bytes::from(vec![0u8; crate::model::datagram::MAX_DATAGRAM_PAYLOAD + 1]);
4732 let ts = Timestamp::from_millis(0).unwrap();
4733 assert!(matches!(
4734 producer.append_datagram(ts, big.clone()),
4735 Err(Error::FrameTooLarge)
4736 ));
4737 assert!(matches!(
4738 producer.insert_datagram(0, ts, big),
4739 Err(Error::FrameTooLarge)
4740 ));
4741 }
4742
4743 #[tokio::test]
4744 async fn datagram_fanout_to_subscribers() {
4745 let mut producer = track_producer("test", None);
4746 let mut a = producer.subscribe(None);
4748 let mut b = producer.subscribe(None);
4749 let ts = Timestamp::from_millis(1).unwrap();
4750
4751 producer.append_datagram(ts, &b"first"[..]).unwrap();
4752 producer.append_datagram(ts, &b"second"[..]).unwrap();
4753
4754 assert_eq!(&recv_datagram(&mut a).payload[..], b"first");
4756 assert_eq!(&recv_datagram(&mut a).payload[..], b"second");
4757 assert_eq!(&recv_datagram(&mut b).payload[..], b"first");
4758 assert_eq!(&recv_datagram(&mut b).payload[..], b"second");
4759 }
4760
4761 #[test]
4762 fn datagram_buffer_drops_oldest_at_capacity() {
4763 let mut producer = track_producer("test", None);
4764 let mut slow = producer.subscribe(None);
4765 let mut fast = producer.subscribe(None);
4766 let count = MAX_DATAGRAMS * 3;
4767 for sequence in 0..count {
4768 producer.append_datagram(Timestamp::ZERO, b"x".as_slice()).unwrap();
4769 assert_eq!(recv_datagram(&mut fast).sequence, sequence as u64);
4770 }
4771 assert_eq!(producer.state.read().datagrams.len(), MAX_DATAGRAMS);
4772 for sequence in count - MAX_DATAGRAMS..count {
4773 assert_eq!(recv_datagram(&mut slow).sequence, sequence as u64);
4774 }
4775 assert!(slow.poll_recv_datagram(&kio::Waiter::noop()).is_pending());
4776 }
4777
4778 #[tokio::test]
4779 async fn datagram_recv_pends_until_written() {
4780 let mut producer = track_producer("test", None);
4781 let mut dg = producer.subscribe(None);
4782
4783 assert!(
4784 dg.recv_datagram().now_or_never().is_none(),
4785 "should block with no datagrams"
4786 );
4787
4788 producer
4789 .append_datagram(Timestamp::from_millis(0).unwrap(), &b"go"[..])
4790 .unwrap();
4791 assert_eq!(&recv_datagram(&mut dg).payload[..], b"go");
4792 }
4793
4794 #[tokio::test]
4798 async fn datagram_wire_roundtrip_between_tracks() {
4799 use crate::coding::{Decode, Encode};
4800 use crate::lite;
4801
4802 let version = lite::Version::Lite05;
4803
4804 let mut origin = track_producer("test", None);
4806 let mut origin_dg = origin.subscribe(None);
4807 let ts = Timestamp::from_millis(7).unwrap();
4808 let seq = origin.append_datagram(ts, &b"payload"[..]).unwrap();
4809
4810 let d = recv_datagram(&mut origin_dg);
4811 let body = lite::Datagram {
4812 subscribe: 5,
4813 sequence: d.sequence,
4814 timestamp: d.timestamp.value(),
4815 payload: d.payload.clone(),
4816 }
4817 .encode_bytes(version)
4818 .unwrap();
4819
4820 let mut slice = &body[..];
4822 let wire = lite::Datagram::decode(&mut slice, version).unwrap();
4823 let mut downstream = track_producer("test", None);
4824 let mut downstream_dg = downstream.subscribe(None);
4825 downstream
4826 .insert_datagram(
4827 wire.sequence,
4828 Timestamp::new(wire.timestamp, Timescale::MILLI).unwrap(),
4829 wire.payload,
4830 )
4831 .unwrap();
4832
4833 let got = recv_datagram(&mut downstream_dg);
4834 assert_eq!(got.sequence, seq);
4835 assert_eq!(got.timestamp, ts);
4836 assert_eq!(&got.payload[..], b"payload");
4837 }
4838
4839 #[tokio::test]
4840 async fn evict_expired_groups() {
4841 let producer = track_producer("test", None);
4842
4843 producer.append_group().unwrap(); producer.append_group().unwrap(); producer.append_group().unwrap(); {
4849 let state = producer.state.read();
4850 assert_eq!(live_groups(&state), 3);
4851 assert_eq!(state.offset, 0);
4852 }
4853
4854 crate::model::clock::advance(cache::DEFAULT_EXPIRY + Duration::from_secs(1));
4856
4857 producer.append_group().unwrap(); {
4864 let state = producer.state.read();
4865 assert_eq!(live_groups(&state), 1);
4866 assert_eq!(first_live_sequence(&state), 3);
4867 assert_eq!(state.offset, 3);
4868 assert!(!state.lookup.contains_key(&0));
4869 assert!(!state.lookup.contains_key(&1));
4870 assert!(!state.lookup.contains_key(&2));
4871 assert!(state.lookup.contains_key(&3));
4872 }
4873 }
4874
4875 #[tokio::test]
4879 async fn aging_out_a_finished_group_keeps_the_clean_end() {
4880 let producer = track_producer("test", None);
4881 let mut group = producer.create_group(group::Info { sequence: 0 }).unwrap();
4882 let mut consumer = group.consume();
4883
4884 group
4885 .write_frame(Timestamp::from_millis(0).unwrap(), b"hello".as_slice())
4886 .unwrap();
4887 assert_eq!(consumer.next_frame().await.unwrap().unwrap().size, 5);
4888
4889 crate::model::clock::advance(cache::DEFAULT_EXPIRY * 2);
4891 group.finish().unwrap();
4892 let _next = producer.create_group(group::Info { sequence: 1 }).unwrap();
4893
4894 assert!(consumer.next_frame().await.unwrap().is_none());
4895 }
4896
4897 #[tokio::test]
4901 async fn active_reader_survives_expiry() {
4902 let producer = track_producer("test", None);
4903 let mut subscriber = producer.subscribe(None);
4904
4905 let mut group = producer.create_group(0u64.into()).unwrap();
4907 for _ in 0..10 {
4908 group.write_frame(Timestamp::ZERO, b"x".as_slice()).unwrap();
4909 }
4910 group.finish().unwrap();
4911 let mut reading = subscriber.assert_group();
4912
4913 producer.create_group(1u64.into()).unwrap().finish().unwrap();
4915
4916 for seq in 2..12u64 {
4919 crate::model::clock::advance(cache::DEFAULT_EXPIRY / 2);
4920 let frame = reading.next_frame().await;
4921 assert!(
4922 matches!(frame, Ok(Some(_))),
4923 "an actively-read group must not expire mid-read (step {seq})"
4924 );
4925 producer.create_group(seq.into()).unwrap().finish().unwrap();
4926 }
4927
4928 let state = producer.state.read();
4929 assert!(state.lookup.contains_key(&0), "the read group survived");
4930 assert!(!state.lookup.contains_key(&1), "the unread group still expired");
4931 }
4932
4933 #[tokio::test]
4938 async fn slow_prefetch_reader_survives_expiry() {
4939 let producer = track_producer("test", None);
4940 let mut subscriber = producer.subscribe(None);
4941
4942 let mut group = producer.create_group(0u64.into()).unwrap();
4943 for _ in 0..20 {
4944 group.write_frame(Timestamp::ZERO, b"x".as_slice()).unwrap();
4945 }
4946 group.finish().unwrap();
4947 let mut reading = subscriber.assert_group();
4948
4949 for seq in 1..20u64 {
4952 crate::model::clock::advance(cache::DEFAULT_EXPIRY / 2);
4953 let frame = reading.read_frame().await;
4954 assert!(
4955 matches!(frame, Ok(Some(_))),
4956 "a slow prefetch reader must not expire mid-read (step {seq})"
4957 );
4958 producer.create_group(seq.into()).unwrap().finish().unwrap();
4959 }
4960 }
4961
4962 #[tokio::test]
4966 async fn delivery_restarts_the_expiry_clock() {
4967 let producer = track_producer("test", None);
4968 let mut subscriber = producer.subscribe(replay());
4969
4970 let mut group = producer.create_group(0u64.into()).unwrap();
4971 group.write_frame(Timestamp::ZERO, b"x".as_slice()).unwrap();
4972 group.finish().unwrap();
4973 producer.create_group(1u64.into()).unwrap().finish().unwrap();
4975
4976 crate::model::clock::advance(cache::DEFAULT_EXPIRY - Duration::from_secs(1));
4978 let mut reading = subscriber.assert_group();
4979
4980 crate::model::clock::advance(cache::DEFAULT_EXPIRY - Duration::from_secs(1));
4983 producer.create_group(2u64.into()).unwrap().finish().unwrap();
4984
4985 let frame = reading.read_frame().await.unwrap();
4986 assert!(frame.is_some(), "a just-delivered group must not expire unread");
4987 }
4988
4989 #[tokio::test]
4993 async fn streaming_frame_writes_keep_the_group_alive() {
4994 let producer = track_producer("test", None);
4995 let mut straggler = producer.create_group(0u64.into()).unwrap();
4996 producer.create_group(1u64.into()).unwrap().finish().unwrap();
4998
4999 let mut frame = straggler
5000 .create_frame(frame::Info {
5001 size: 10,
5002 timestamp: Timestamp::ZERO,
5003 })
5004 .unwrap();
5005 for seq in 2..12u64 {
5008 crate::model::clock::advance(cache::DEFAULT_EXPIRY / 2);
5009 frame.write(bytes::Bytes::from_static(b"x")).unwrap();
5010 producer.create_group(seq.into()).unwrap().finish().unwrap();
5011 }
5012 frame.finish().unwrap();
5013 straggler.finish().unwrap();
5014
5015 let state = producer.state.read();
5016 assert!(
5017 state.lookup.contains_key(&0),
5018 "a group streaming a frame survives expiry"
5019 );
5020 }
5021
5022 #[tokio::test]
5027 async fn coalesced_frame_completion_keeps_the_group_alive() {
5028 let producer = track_producer("test", None);
5029 let mut straggler = producer.create_group(0u64.into()).unwrap();
5030 producer.create_group(1u64.into()).unwrap().finish().unwrap();
5032
5033 let mut frame = straggler
5034 .create_frame_owned(frame::Info {
5035 size: 3,
5036 timestamp: Timestamp::ZERO,
5037 })
5038 .unwrap();
5039
5040 crate::model::clock::advance(cache::DEFAULT_EXPIRY + Duration::from_secs(1));
5044 frame.write(bytes::Bytes::from_static(b"abc")).unwrap();
5045 frame.finish().unwrap();
5046 straggler.finish().unwrap();
5047
5048 producer.create_group(2u64.into()).unwrap().finish().unwrap();
5050
5051 let state = producer.state.read();
5052 assert!(
5053 state.lookup.contains_key(&0),
5054 "a group whose frame just completed must not expire"
5055 );
5056 }
5057
5058 #[tokio::test]
5061 async fn parked_reoffer_restarts_the_expiry_clock() {
5062 let producer = track_producer("test", None);
5063 let mut subscriber = producer.subscribe(None);
5064 subscriber.set_groups(..1);
5065
5066 for seq in 0..2u64 {
5067 let mut group = producer.create_group(seq.into()).unwrap();
5068 group.write_frame(Timestamp::ZERO, b"x".as_slice()).unwrap();
5069 group.finish().unwrap();
5070 }
5071
5072 assert_eq!(subscriber.assert_group().sequence, 0);
5074 subscriber.assert_no_group();
5075
5076 crate::model::clock::advance(cache::DEFAULT_EXPIRY - Duration::from_secs(1));
5078 subscriber.set_groups(..2);
5079 let mut reading = subscriber.assert_group();
5080 assert_eq!(reading.sequence, 1);
5081
5082 crate::model::clock::advance(cache::DEFAULT_EXPIRY - Duration::from_secs(1));
5085 producer.create_group(2u64.into()).unwrap().finish().unwrap();
5086
5087 let frame = reading.read_frame().await.unwrap();
5088 assert!(frame.is_some(), "a just-re-offered group must not expire unread");
5089 }
5090
5091 #[tokio::test]
5092 async fn evict_keeps_max_sequence() {
5093 let producer = track_producer("test", None);
5094 producer.append_group().unwrap(); crate::model::clock::advance(cache::DEFAULT_EXPIRY + Duration::from_secs(1));
5098
5099 producer.append_group().unwrap(); {
5103 let state = producer.state.read();
5104 assert_eq!(live_groups(&state), 1);
5105 assert_eq!(first_live_sequence(&state), 1);
5106 assert_eq!(state.offset, 1);
5107 }
5108 }
5109
5110 #[tokio::test]
5111 async fn no_eviction_when_fresh() {
5112 let producer = track_producer("test", None);
5113 producer.append_group().unwrap(); producer.append_group().unwrap(); producer.append_group().unwrap(); {
5118 let state = producer.state.read();
5119 assert_eq!(live_groups(&state), 3);
5120 assert_eq!(state.offset, 0);
5121 }
5122 }
5123
5124 #[tokio::test]
5125 async fn consumer_skips_evicted_groups() {
5126 let producer = track_producer("test", None);
5127 producer.append_group().unwrap(); let mut consumer = producer.subscribe(None);
5130
5131 crate::model::clock::advance(cache::DEFAULT_EXPIRY + Duration::from_secs(1));
5132 producer.append_group().unwrap(); let group = consumer.assert_group();
5136 assert_eq!(group.sequence, 1);
5137 }
5138
5139 fn track_producer_expiring(name: impl Into<Arc<str>>, expiry: impl Into<Option<Duration>>) -> Producer {
5141 track_producer_pooled(name, cache::Pool::new(cache::Config::default().with_expiry(expiry)))
5142 }
5143
5144 fn track_producer_pooled(name: impl Into<Arc<str>>, pool: cache::Pool) -> Producer {
5146 Producer::new(
5147 Arc::new(broadcast::Info {
5148 pool,
5149 ..Default::default()
5150 }),
5151 name,
5152 None,
5153 )
5154 }
5155
5156 #[tokio::test]
5160 async fn pool_sweep_expires_without_a_write() {
5161 let pool = cache::Pool::new(cache::Config::default().with_expiry(Duration::from_secs(1)));
5162 let producer = track_producer_pooled("test", pool.clone());
5163 let mut stalled = producer.append_group().unwrap(); stalled.write_frame(Timestamp::ZERO, b"x".as_slice()).unwrap();
5165 producer.append_group().unwrap(); crate::model::clock::advance(Duration::from_secs(2));
5168 pool.sweep();
5169
5170 assert!(
5171 !producer.state.read().lookup.contains_key(&0),
5172 "the sweep reclaimed an idle open group with no write behind it"
5173 );
5174 }
5175
5176 #[test]
5177 fn cache_gc_dates_activity_before_expiring_it() {
5178 let pool = cache::Pool::new(cache::Config::default().with_expiry(Duration::from_secs(1)));
5179 let now = crate::model::clock::now();
5180 assert_eq!(pool.gc(now), Some(now + Duration::from_millis(500)));
5181 let producer = track_producer_pooled("test", pool.clone());
5182 let mut group = producer.append_group().unwrap();
5183 group.write_frame(Timestamp::ZERO, b"first".as_slice()).unwrap();
5184 producer.append_group().unwrap();
5185 let later = now + Duration::from_secs(60);
5187 pool.gc(later);
5188 assert!(producer.state.read().lookup.contains_key(&0));
5189 pool.gc(later + Duration::from_secs(2));
5190 assert!(!producer.state.read().lookup.contains_key(&0));
5191 }
5192
5193 #[test]
5194 fn cache_gc_reaches_old_entries_behind_a_fresh_front() {
5195 let expiry = Duration::from_secs(1);
5196 let pool = cache::Pool::new(cache::Config::default().with_expiry(expiry));
5197 let producer = track_producer_pooled("test", pool.clone());
5198 let now = crate::model::clock::now();
5199 let groups: Vec<_> = (0..EVICT_SCAN * 3).map(|_| producer.append_group().unwrap()).collect();
5200 producer.append_group().unwrap();
5201 pool.gc(now);
5202 for group in &groups[..EVICT_SCAN * 2] {
5204 group.cache_refresh();
5205 }
5206 pool.gc(now + expiry * 2);
5207 let state = producer.state.read();
5208 for sequence in 0..EVICT_SCAN * 2 {
5209 assert!(state.lookup.contains_key(&(sequence as u64)), "fresh front survives");
5210 }
5211 for sequence in EVICT_SCAN * 2..EVICT_SCAN * 3 {
5212 assert!(!state.lookup.contains_key(&(sequence as u64)), "old tail is reclaimed");
5213 }
5214 }
5215
5216 #[tokio::test]
5220 async fn pool_sweep_drains_a_deep_backlog() {
5221 let pool = cache::Pool::new(cache::Config::default().with_expiry(Duration::from_secs(1)));
5222 let producer = track_producer_pooled("test", pool.clone());
5223
5224 let backlog = 4 * EVICT_SCAN;
5226 for _ in 0..backlog {
5227 let mut group = producer.append_group().unwrap();
5228 group.write_frame(Timestamp::ZERO, b"x".as_slice()).unwrap();
5229 }
5230 producer.append_group().unwrap(); crate::model::clock::advance(Duration::from_secs(2));
5233 pool.sweep();
5234
5235 let state = producer.state.read();
5236 let stale = (0..backlog as u64).filter(|seq| state.lookup.contains_key(seq)).count();
5237 assert_eq!(stale, 0, "one sweep reclaimed the whole idle backlog");
5238 }
5239
5240 #[tokio::test]
5241 async fn pool_expiry_controls_eviction() {
5242 let producer = track_producer_expiring("test", Duration::from_secs(1));
5244 producer.append_group().unwrap(); crate::model::clock::advance(Duration::from_secs(2));
5248 producer.append_group().unwrap(); let state = producer.state.read();
5252 assert_eq!(live_groups(&state), 1);
5253 assert_eq!(first_live_sequence(&state), 1);
5254 }
5255
5256 #[tokio::test]
5257 async fn small_frame_write_expires_idle_siblings() {
5258 let producer = track_producer_expiring("test", Duration::from_secs(1));
5259 producer.append_group().unwrap().finish().unwrap(); let mut live = producer.append_group().unwrap(); crate::model::clock::advance(Duration::from_secs(2));
5263 live.write_frame(Timestamp::ZERO, b"x".as_slice()).unwrap();
5264
5265 let expired = !producer.state.read().lookup.contains_key(&0);
5266 assert!(expired, "a small frame write runs expiry");
5267 }
5268
5269 #[tokio::test]
5270 async fn fresh_expiry_scan_does_not_wake_track_consumers() {
5271 use std::sync::atomic::{AtomicBool, Ordering};
5272
5273 let producer = track_producer_expiring("test", cache::DEFAULT_EXPIRY);
5274 producer.append_group().unwrap().finish().unwrap();
5275 let mut live = producer.append_group().unwrap();
5276 let mut consumer = producer.subscribe(None);
5277 assert_eq!(consumer.assert_group().sequence, 0);
5278 assert_eq!(consumer.assert_group().sequence, 1);
5279
5280 let woken = Arc::new(AtomicBool::new(false));
5281 let waiter = kio::Waiter::new(futures::task::waker(Arc::new(FlagWake(woken.clone()))));
5282 assert!(consumer.poll_recv_group(&waiter).is_pending());
5283
5284 live.write_frame(Timestamp::ZERO, b"x".as_slice()).unwrap();
5285 assert!(
5286 !woken.load(Ordering::SeqCst),
5287 "a no-op expiry scan must not wake track consumers"
5288 );
5289 }
5290
5291 #[tokio::test]
5292 async fn streaming_frame_write_expires_idle_siblings() {
5293 let producer = track_producer_expiring("test", Duration::from_secs(1));
5294 producer.append_group().unwrap().finish().unwrap(); let mut live = producer.append_group().unwrap(); let mut frame = live
5297 .create_frame(frame::Info {
5298 size: 1,
5299 timestamp: Timestamp::ZERO,
5300 })
5301 .unwrap();
5302
5303 crate::model::clock::advance(Duration::from_secs(2));
5304 frame.write(b"x".as_slice()).unwrap();
5305
5306 let expired = !producer.state.read().lookup.contains_key(&0);
5307 assert!(expired, "a streamed chunk runs expiry");
5308 }
5309
5310 #[tokio::test]
5311 async fn appended_datagram_expires_idle_groups() {
5312 let mut producer = track_producer_expiring("test", Duration::from_secs(1));
5313 producer.append_group().unwrap().finish().unwrap(); producer.append_group().unwrap().finish().unwrap(); crate::model::clock::advance(Duration::from_secs(2));
5317 producer.append_datagram(Timestamp::ZERO, b"x".as_slice()).unwrap();
5318
5319 let expired = !producer.state.read().lookup.contains_key(&0);
5320 assert!(expired, "an appended datagram runs expiry");
5321 }
5322
5323 #[tokio::test]
5324 async fn forwarded_datagram_expires_idle_groups() {
5325 let mut producer = track_producer_expiring("test", Duration::from_secs(1));
5326 producer.append_group().unwrap().finish().unwrap(); producer.append_group().unwrap().finish().unwrap(); crate::model::clock::advance(Duration::from_secs(2));
5330 producer
5331 .insert_datagram(2, Timestamp::ZERO, bytes::Bytes::from_static(b"x"))
5332 .unwrap();
5333
5334 let expired = !producer.state.read().lookup.contains_key(&0);
5335 assert!(expired, "a forwarded datagram runs expiry");
5336 }
5337
5338 #[tokio::test]
5342 async fn max_age_does_not_drive_wall_eviction() {
5343 let producer = track_producer("test", Info::default().with_max_age(Duration::from_secs(1)));
5344 producer.append_group().unwrap(); crate::model::clock::advance(Duration::from_secs(10));
5348 producer.append_group().unwrap(); let state = producer.state.read();
5351 assert_eq!(live_groups(&state), 2, "max_age is media time, not a wall clock");
5352 }
5353
5354 #[tokio::test]
5356 async fn disabled_pool_expiry_never_reclaims() {
5357 let producer = track_producer_expiring("test", None);
5358 producer.append_group().unwrap(); crate::model::clock::advance(Duration::from_secs(3600));
5361 producer.append_group().unwrap(); let state = producer.state.read();
5364 assert_eq!(live_groups(&state), 2);
5365 }
5366
5367 #[test]
5368 fn max_age_clamped_to_cache() {
5369 let producer = track_producer("test", Info::default().with_max_age(Duration::from_secs(2)));
5370
5371 let mut subscriber = producer.subscribe(Subscription::default().with_max_age(Duration::from_secs(10)));
5375 assert_eq!(subscriber.subscription().max_age, Duration::from_secs(10));
5376 assert_eq!(producer.subscription().unwrap().max_age, Duration::from_secs(2));
5377
5378 subscriber
5380 .update(Subscription::default().with_max_age(Duration::from_millis(500)))
5381 .unwrap();
5382 assert_eq!(producer.subscription().unwrap().max_age, Duration::from_millis(500));
5383
5384 subscriber
5385 .update(Subscription::default().with_max_age(Duration::ZERO))
5386 .unwrap();
5387 assert_eq!(producer.subscription().unwrap().max_age, Duration::ZERO);
5388 }
5389
5390 fn track_producer_capped(name: impl Into<Arc<str>>, info: Info, cap: Duration) -> Producer {
5393 Producer::new(
5394 Arc::new(broadcast::Info {
5395 cache_duration: cap,
5396 ..Default::default()
5397 }),
5398 name,
5399 info,
5400 )
5401 }
5402
5403 #[test]
5404 fn origin_cache_duration_clamps_max_age() {
5405 let capped = track_producer_capped(
5408 "test",
5409 Info::default().with_max_age(Duration::from_secs(60)),
5410 Duration::from_secs(1),
5411 );
5412 assert_eq!(capped.state.read().max_age_bound(), Some(Duration::from_secs(1)));
5413 assert_eq!(capped.subscribe(None).info().max_age, Duration::from_secs(1));
5414
5415 let under = track_producer_capped(
5416 "test",
5417 Info::default().with_max_age(Duration::from_millis(500)),
5418 Duration::from_secs(1),
5419 );
5420 assert_eq!(under.state.read().max_age_bound(), Some(Duration::from_millis(500)));
5421 }
5422
5423 #[tokio::test]
5426 async fn origin_cache_duration_does_not_wall_evict() {
5427 let producer = track_producer_capped(
5428 "test",
5429 Info::default().with_max_age(Duration::from_secs(60)),
5430 Duration::from_secs(1),
5431 );
5432 producer.append_group().unwrap(); crate::model::clock::advance(Duration::from_secs(2));
5436 producer.append_group().unwrap(); let state = producer.state.read();
5439 assert_eq!(live_groups(&state), 2);
5440 }
5441
5442 #[test]
5443 fn max_age_clamped_via_every_update_path() {
5444 let producer = track_producer("test", Info::default().with_max_age(Duration::from_secs(2)));
5445 let over = Subscription::default().with_max_age(Duration::from_secs(10));
5446
5447 let mut subscriber = producer.subscribe(over.clone());
5450 assert_eq!(producer.subscription().unwrap().max_age, Duration::from_secs(2));
5451
5452 subscriber.control().update(over.clone()).unwrap();
5453 assert_eq!(producer.subscription().unwrap().max_age, Duration::from_secs(2));
5454
5455 subscriber.update(over).unwrap();
5456 assert_eq!(producer.subscription().unwrap().max_age, Duration::from_secs(2));
5457 }
5458
5459 #[test]
5460 fn max_age_aggregate_clamps_across_subscribers() {
5461 let producer = track_producer("test", Info::default().with_max_age(Duration::from_secs(2)));
5462
5463 let _a = producer.subscribe(Subscription::default().with_max_age(Duration::from_millis(500)));
5466 let _b = producer.subscribe(Subscription::default().with_max_age(Duration::from_secs(10)));
5467
5468 assert_eq!(producer.subscription().unwrap().max_age, Duration::from_secs(2));
5469 }
5470
5471 fn append_at(producer: &mut Producer, millis: u64) -> u64 {
5474 let mut group = producer.append_group().unwrap();
5475 group
5476 .write_frame(Timestamp::from_millis(millis).unwrap(), bytes::Bytes::from_static(b"x"))
5477 .unwrap();
5478 group.finish().unwrap();
5479 group.sequence
5480 }
5481
5482 fn drain(subscriber: &mut Subscriber) -> Vec<u64> {
5484 let mut sequences = Vec::new();
5485 while let Some(Ok(Some(group))) = subscriber.recv_group().now_or_never() {
5486 sequences.push(group.sequence);
5487 }
5488 sequences
5489 }
5490
5491 #[test]
5492 fn real_time_skips_a_backlog_to_the_live_edge() {
5493 let mut producer = track_producer("test", None);
5494 for second in 0..5 {
5495 append_at(&mut producer, second * 1000);
5496 }
5497
5498 let mut subscriber = producer.subscribe(None);
5502 assert_eq!(drain(&mut subscriber), vec![4]);
5503
5504 append_at(&mut producer, 5000);
5506 assert_eq!(drain(&mut subscriber), vec![5]);
5507 }
5508
5509 #[test]
5510 fn real_time_skips_a_backlog_after_catching_up() {
5511 let mut producer = track_producer("test", None);
5512 append_at(&mut producer, 0);
5513 let mut subscriber = producer.subscribe(None);
5514 assert_eq!(drain(&mut subscriber), vec![0]);
5515
5516 for second in 1..6 {
5519 append_at(&mut producer, second * 1000);
5520 }
5521 assert_eq!(drain(&mut subscriber), vec![5]);
5522 }
5523
5524 #[test]
5525 fn a_newer_edge_changes_an_active_catch_up() {
5526 let mut producer = track_producer("test", None);
5527 for second in 0..5 {
5528 append_at(&mut producer, second * 1000);
5529 }
5530 let mut subscriber = producer.subscribe(Subscription::default().with_max_age(Duration::from_secs(1)));
5531 assert_eq!(
5532 subscriber
5533 .recv_group()
5534 .now_or_never()
5535 .unwrap()
5536 .unwrap()
5537 .unwrap()
5538 .sequence,
5539 3
5540 );
5541
5542 append_at(&mut producer, 5000);
5545 append_at(&mut producer, 10000);
5546 assert_eq!(drain(&mut subscriber), vec![5, 6]);
5547 }
5548
5549 #[test]
5550 fn a_growing_edge_changes_an_active_catch_up() {
5551 let mut producer = track_producer("test", None);
5552 append_at(&mut producer, 0);
5553 append_at(&mut producer, 1000);
5554 let mut edge = producer.append_group().unwrap();
5555 edge.write_frame(Timestamp::from_millis(2000).unwrap(), bytes::Bytes::from_static(b"a"))
5556 .unwrap();
5557
5558 let mut subscriber = producer.subscribe(Subscription::default().with_max_age(Duration::from_secs(2)));
5559 assert_eq!(
5560 subscriber
5561 .recv_group()
5562 .now_or_never()
5563 .unwrap()
5564 .unwrap()
5565 .unwrap()
5566 .sequence,
5567 0
5568 );
5569
5570 edge.write_frame(Timestamp::from_millis(5000).unwrap(), bytes::Bytes::from_static(b"b"))
5573 .unwrap();
5574 assert_eq!(drain(&mut subscriber), vec![2]);
5575 }
5576
5577 #[test]
5578 fn a_budget_admits_groups_within_it() {
5579 let mut producer = track_producer("test", None);
5580 for second in 0..5 {
5581 append_at(&mut producer, second * 1000);
5582 }
5583
5584 let mut subscriber = producer.subscribe(Subscription::default().with_max_age(Duration::from_secs(2)));
5589 assert_eq!(drain(&mut subscriber), vec![2, 3, 4]);
5590 }
5591
5592 #[test]
5593 fn a_budget_reaches_back_over_the_cache() {
5594 let mut producer = track_producer("test", None);
5595 for second in 0..5 {
5596 append_at(&mut producer, second * 1000);
5597 }
5598
5599 let mut live = producer.subscribe(None);
5602 assert_eq!(drain(&mut live), vec![4]);
5603
5604 let budget = Subscription::default().with_max_age(Duration::from_secs(2));
5609 let mut subscriber = producer.subscribe(budget);
5610 assert_eq!(drain(&mut subscriber), vec![2, 3, 4]);
5611 }
5612
5613 #[test]
5614 fn a_named_start_is_a_floor_not_a_request() {
5615 let mut producer = track_producer("test", None);
5616 for second in 0..5 {
5617 append_at(&mut producer, second * 1000);
5618 }
5619
5620 let named = Subscription::default().with_start(Position::group(1));
5624 let mut subscriber = producer.subscribe(named);
5625 assert_eq!(drain(&mut subscriber), vec![4]);
5626
5627 let floored = Subscription::default()
5629 .with_start(Position::group(3))
5630 .with_max_age(Duration::from_secs(10));
5631 let mut subscriber = producer.subscribe(floored);
5632 assert_eq!(drain(&mut subscriber), vec![3, 4]);
5633
5634 let slack = Subscription::default()
5636 .with_start(Position::group(1))
5637 .with_max_age(Duration::from_secs(2));
5638 let mut subscriber = producer.subscribe(slack);
5639 assert_eq!(drain(&mut subscriber), vec![2, 3, 4]);
5640 }
5641
5642 #[test]
5643 fn a_floor_above_the_live_edge_waits_there() {
5644 let mut producer = track_producer("test", None);
5645 for second in 0..3 {
5646 append_at(&mut producer, second * 1000);
5647 }
5648
5649 let resumed = Subscription::default()
5652 .with_start(Position::group(7))
5653 .with_max_age(Duration::from_secs(10));
5654 let mut subscriber = producer.subscribe(resumed);
5655 assert_eq!(drain(&mut subscriber), Vec::<u64>::new());
5656 append_at(&mut producer, 3000); assert_eq!(drain(&mut subscriber), Vec::<u64>::new());
5658 for second in 4..8 {
5659 append_at(&mut producer, second * 1000);
5660 }
5661 assert_eq!(drain(&mut subscriber), vec![7]);
5662 }
5663
5664 #[test]
5665 fn a_late_lower_group_within_the_budget_is_delivered() {
5666 let producer = track_producer("test", None);
5667 for (sequence, millis) in [(5, 0), (6, 1000), (7, 2000)] {
5668 let mut group = producer.create_group(group::Info { sequence }).unwrap();
5669 group
5670 .write_frame(Timestamp::from_millis(millis).unwrap(), bytes::Bytes::from_static(b"x"))
5671 .unwrap();
5672 group.finish().unwrap();
5673 }
5674
5675 let mut subscriber = producer.subscribe(Subscription::default().with_max_age(Duration::from_secs(5)));
5676 assert_eq!(drain(&mut subscriber), vec![5, 6, 7]);
5677
5678 let mut late = producer.create_group(group::Info { sequence: 4 }).unwrap();
5682 late.write_frame(Timestamp::from_millis(500).unwrap(), bytes::Bytes::from_static(b"late"))
5683 .unwrap();
5684 late.finish().unwrap();
5685 assert_eq!(drain(&mut subscriber), vec![4]);
5686 }
5687
5688 #[test]
5689 fn drift_is_measured_in_presentation_time_not_arrival_time() {
5690 let mut producer = track_producer("test", None);
5691 for second in 0..4 {
5695 append_at(&mut producer, second * 1000);
5696 }
5697
5698 let mut subscriber = producer.subscribe(Subscription::default().with_max_age(Duration::from_millis(1500)));
5699 assert_eq!(drain(&mut subscriber), vec![1, 2, 3]);
5700 }
5701
5702 #[tokio::test]
5703 async fn a_stamped_successor_expires_an_unstamped_group() {
5704 let mut producer = track_producer("test", None);
5705 let mut subscriber = producer.subscribe(None);
5706 producer.append_group().unwrap(); append_at(&mut producer, 1000); assert_eq!(drain(&mut subscriber), vec![1]);
5713 }
5714
5715 #[tokio::test]
5716 async fn a_handed_out_group_expires_while_its_first_frame_is_stalled() {
5717 let mut producer = track_producer("test", None);
5718 let mut subscriber = producer.subscribe(None);
5719 producer.append_group().unwrap();
5720
5721 let mut stalled = subscriber.recv_group().await.unwrap().expect("stalled group");
5722 let pending = tokio::spawn(async move { stalled.read_frame().await });
5723 tokio::task::yield_now().await;
5724 assert!(
5725 !pending.is_finished(),
5726 "the empty live edge still waits for its first frame"
5727 );
5728
5729 crate::model::clock::advance(Duration::from_secs(1));
5730 append_at(&mut producer, 1000);
5731
5732 let result = pending.await.unwrap();
5737 assert!(matches!(result, Ok(None)), "the held group ends: {result:?}");
5738 }
5739
5740 #[tokio::test]
5745 async fn a_handed_out_group_wakes_when_a_newer_group_gets_its_first_timestamp() {
5746 let mut producer = track_producer("test", None);
5747 let mut subscriber = producer.subscribe(Subscription::default().with_max_age(Duration::from_millis(500)));
5748 let mut old = producer.append_group().unwrap();
5749 old.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"old"))
5750 .unwrap();
5751
5752 append_at(&mut producer, 1000);
5754
5755 let mut held = subscriber.recv_group().await.unwrap().expect("old group");
5756 assert!(held.read_frame().await.unwrap().is_some());
5757
5758 let mut live = producer.append_group().unwrap();
5759 let pending = tokio::spawn(async move { held.read_frame().await });
5760 tokio::task::yield_now().await;
5761 assert!(!pending.is_finished(), "the newer group has no timestamp yet");
5762
5763 live.write_frame(
5765 Timestamp::from_millis(2000).unwrap(),
5766 bytes::Bytes::from_static(b"live"),
5767 )
5768 .unwrap();
5769 tokio::task::yield_now().await;
5770
5771 assert!(pending.is_finished(), "the new presentation edge wakes the held reader");
5772 let result = pending.await.unwrap();
5774 assert!(matches!(result, Ok(None)), "the held group ends: {result:?}");
5775 }
5776
5777 #[tokio::test]
5788 async fn real_time_reads_a_live_stream_without_truncating_it() {
5789 let producer = track_producer("test", None);
5790 let mut subscriber = producer.subscribe(None);
5791
5792 let gop = |n: u64| {
5793 [
5794 Timestamp::from_millis(n * 2000).unwrap(),
5795 Timestamp::from_millis(n * 2000 + 1900).unwrap(),
5796 ]
5797 };
5798 let write = |group: &mut group::Producer, timestamp| {
5799 group.write_frame(timestamp, bytes::Bytes::from_static(b"x")).unwrap();
5800 };
5801
5802 let mut open = producer.append_group().unwrap();
5803 write(&mut open, gop(0)[0]);
5804 write(&mut open, gop(0)[1]);
5805 let mut reading = subscriber.recv_group().await.unwrap().expect("the live group");
5806
5807 let mut read = Vec::new();
5808 for n in 1..5u64 {
5809 let sequence = reading.sequence;
5810 let mut frames = 0;
5811 while let Some(res) = reading.read_frame().now_or_never() {
5812 match res.expect("no truncation while draining") {
5813 Some(_) => frames += 1,
5814 None => panic!("group {sequence} ended early"),
5815 }
5816 }
5817 read.push((sequence, frames));
5818
5819 let next = {
5820 let mut end = std::pin::pin!(reading.read_frame());
5823 assert!(futures::poll!(end.as_mut()).is_pending(), "parked on the FIN");
5824
5825 let mut opened = producer.append_group().unwrap();
5828 write(&mut opened, gop(n)[0]);
5829 let verdict = futures::poll!(end.as_mut());
5830
5831 open.finish().unwrap();
5832 let res = match verdict {
5833 Poll::Ready(res) => res,
5834 Poll::Pending => end.await,
5835 };
5836 assert!(
5837 matches!(res, Ok(None)),
5838 "group {sequence} ends at the boundary rather than failing: {res:?}"
5839 );
5840 opened
5841 };
5842 let mut next = next;
5843 write(&mut next, gop(n)[1]);
5844
5845 reading = subscriber.recv_group().await.unwrap().expect("the next live group");
5846 open = next;
5847 }
5848
5849 assert_eq!(read, vec![(0, 2), (1, 2), (2, 2), (3, 2)], "every frame of every group");
5850 }
5851
5852 #[tokio::test]
5862 async fn a_budget_is_measured_from_the_readers_position() {
5863 let producer = track_producer("test", None);
5864 let mut subscriber = producer.subscribe(Subscription::default().with_max_age(Duration::from_secs(1)));
5865
5866 let mut open = producer.append_group().unwrap();
5867 open.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"key"))
5868 .unwrap();
5869 open.write_frame(
5870 Timestamp::from_millis(1900).unwrap(),
5871 bytes::Bytes::from_static(b"tail"),
5872 )
5873 .unwrap();
5874
5875 let mut reading = subscriber.recv_group().await.unwrap().expect("the live group");
5876 assert!(reading.read_frame().await.unwrap().is_some());
5877 assert!(reading.read_frame().await.unwrap().is_some());
5878
5879 let mut end = std::pin::pin!(reading.read_frame());
5880 assert!(futures::poll!(end.as_mut()).is_pending(), "parked at 1900ms");
5881
5882 let mut next = producer.append_group().unwrap();
5884 next.write_frame(Timestamp::from_millis(2000).unwrap(), bytes::Bytes::from_static(b"key"))
5885 .unwrap();
5886 assert!(
5887 futures::poll!(end.as_mut()).is_pending(),
5888 "a reader inside its budget is not expired by the next group opening"
5889 );
5890
5891 open.write_frame(
5893 Timestamp::from_millis(1950).unwrap(),
5894 bytes::Bytes::from_static(b"late"),
5895 )
5896 .unwrap();
5897 let late = end.await.expect("the straggler is not truncated");
5898 assert_eq!(
5899 late.map(|frame| frame.timestamp),
5900 Some(Timestamp::from_millis(1950).unwrap())
5901 );
5902 }
5903
5904 #[tokio::test]
5908 async fn an_ended_group_stays_ended_when_probed_again() {
5909 let mut producer = track_producer("test", None);
5910 let mut subscriber = producer.subscribe(None);
5911 let mut open = producer.append_group().unwrap();
5912 open.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"a"))
5913 .unwrap();
5914
5915 let mut group = subscriber.recv_group().await.unwrap().expect("group");
5916 assert!(group.read_frame().await.unwrap().is_some());
5917
5918 let probes = tokio::spawn(async move {
5919 let first = group.read_frame().await;
5920 let second = group.read_frame().await;
5921 let finished = group.finished().await;
5922 (first, second, finished)
5923 });
5924 tokio::task::yield_now().await;
5925
5926 append_at(&mut producer, 1000);
5927
5928 let (first, second, finished) = probes.await.unwrap();
5929 assert!(matches!(first, Ok(None)), "the group ends: {first:?}");
5930 assert!(matches!(second, Ok(None)), "and stays ended: {second:?}");
5931 assert!(matches!(finished, Ok(1)), "reporting what it delivered: {finished:?}");
5932 }
5933
5934 #[tokio::test]
5935 async fn a_drained_group_finishes_cleanly_after_the_live_edge_advances() {
5936 let mut producer = track_producer("test", None);
5937 let mut subscriber = producer.subscribe(None);
5938 append_at(&mut producer, 0);
5939
5940 let mut group = subscriber.recv_group().await.unwrap().expect("first group");
5941 assert!(group.read_frame().await.unwrap().is_some());
5942
5943 append_at(&mut producer, 1000);
5944
5945 assert!(group.read_frame().await.unwrap().is_none());
5946 assert!(!group.latency_expired());
5947 }
5948
5949 #[tokio::test]
5950 async fn a_handed_out_partial_frame_expires_while_its_payload_is_stalled() {
5951 let mut producer = track_producer("test", None);
5952 let mut subscriber = producer.subscribe(None);
5953 let mut source = producer.append_group().unwrap();
5954 let mut writing = source
5955 .create_frame(frame::Info {
5956 size: 6,
5957 timestamp: Timestamp::ZERO,
5958 })
5959 .unwrap();
5960 writing.write(bytes::Bytes::from_static(b"old")).unwrap();
5961
5962 let mut group = subscriber.recv_group().await.unwrap().expect("partial group");
5963 let mut frame = group.next_frame().await.unwrap().expect("partial frame");
5964 assert_eq!(
5965 frame.read_chunk().await.unwrap(),
5966 Some(bytes::Bytes::from_static(b"old"))
5967 );
5968 let pending = tokio::spawn(async move { frame.read_chunk().await });
5969 tokio::task::yield_now().await;
5970 assert!(!pending.is_finished(), "the partial payload is still stalled");
5971
5972 crate::model::clock::advance(Duration::from_secs(1));
5973 append_at(&mut producer, 1000);
5974
5975 let result = pending.await.unwrap();
5976 assert!(
5977 matches!(result, Err(Error::Old)),
5978 "the in-flight frame expires: {result:?}"
5979 );
5980 writing.abort(Error::Cancel).unwrap();
5981 }
5982
5983 #[test]
5984 fn max_age_bounds_the_budget() {
5985 let mut producer = track_producer("test", Info::default().with_max_age(Duration::from_millis(500)));
5989 append_at(&mut producer, 0);
5990 append_at(&mut producer, 1000);
5991 append_at(&mut producer, 2000);
5992
5993 let mut subscriber = producer.subscribe(Subscription::default().with_max_age(Duration::from_secs(10)));
5996 assert_eq!(drain(&mut subscriber), vec![1, 2]);
5997 }
5998
5999 #[test]
6000 fn an_explicit_start_gets_no_exemption_from_the_budget() {
6001 let mut producer = track_producer("test", None);
6002 for second in 0..4 {
6003 append_at(&mut producer, second * 1000);
6004 }
6005
6006 let mut subscriber = producer.subscribe(Subscription::default().with_start(Position::group(0)));
6009 subscriber.start_at(0);
6010 assert_eq!(drain(&mut subscriber), vec![3]);
6011
6012 let mut patient = producer.subscribe(
6013 Subscription::default()
6014 .with_start(Position::group(0))
6015 .with_max_age(Duration::from_secs(10)),
6016 );
6017 patient.start_at(0);
6018 assert_eq!(drain(&mut patient), vec![0, 1, 2, 3]);
6019 }
6020
6021 #[test]
6025 fn a_long_group_is_not_stale_while_its_tail_reaches_the_edge() {
6026 let mut producer = track_producer("test", None);
6027
6028 let mut long = producer.append_group().unwrap();
6030 for ms in [0u64, 500, 1000, 1500, 2000] {
6031 long.write_frame(Timestamp::from_millis(ms).unwrap(), bytes::Bytes::from_static(b"x"))
6032 .unwrap();
6033 }
6034 long.finish().unwrap();
6035 append_at(&mut producer, 2000);
6036
6037 let mut sub = producer.subscribe(Subscription::default().with_max_age(Duration::from_millis(500)));
6039 assert_eq!(drain(&mut sub), vec![0, 1]);
6040 }
6041
6042 #[test]
6048 fn a_long_group_is_stale_once_its_successor_falls_behind() {
6049 let mut producer = track_producer("test", None);
6050
6051 let mut long = producer.append_group().unwrap();
6052 for ms in [0u64, 500, 1000] {
6053 long.write_frame(Timestamp::from_millis(ms).unwrap(), bytes::Bytes::from_static(b"x"))
6054 .unwrap();
6055 }
6056 long.finish().unwrap();
6057 append_at(&mut producer, 3000);
6060 append_at(&mut producer, 4000);
6061
6062 let mut sub = producer.subscribe(Subscription::default().with_max_age(Duration::from_millis(500)));
6063 assert_eq!(drain(&mut sub), vec![1, 2]);
6064 }
6065
6066 #[tokio::test]
6070 async fn an_unstamped_immediate_successor_leaves_reach_unbounded() {
6071 let mut producer = track_producer("test", None);
6072 let mut subscriber = producer.subscribe(None);
6073
6074 append_at(&mut producer, 0); producer.append_group().unwrap(); append_at(&mut producer, 10_000); assert_eq!(drain(&mut subscriber), vec![0, 2]);
6082 }
6083
6084 #[test]
6090 fn reach_follows_the_immediate_successor_not_a_later_rewind() {
6091 let mut producer = track_producer("test", None);
6092
6093 append_at(&mut producer, 0);
6095 append_at(&mut producer, 10_000);
6096 append_at(&mut producer, 1_000);
6097 append_at(&mut producer, 2_000);
6098
6099 let mut sub = producer.subscribe(Subscription::default().with_max_age(Duration::from_millis(500)));
6103 assert!(
6104 drain(&mut sub).contains(&0),
6105 "group 0 is bounded by its successor at 10s, not by a later rewind"
6106 );
6107 }
6108
6109 #[tokio::test]
6112 async fn ordered_carries_datagrams() {
6113 let mut producer = track_producer("test", None);
6114 let mut sub = producer.subscribe(None).ordered();
6115
6116 producer
6117 .insert_datagram(5, Timestamp::from_millis(5).unwrap(), bytes::Bytes::from_static(b"x"))
6118 .unwrap();
6119 producer.create_group(group::Info { sequence: 3 }).unwrap();
6120
6121 let datagram = sub
6122 .recv_datagram()
6123 .now_or_never()
6124 .expect("datagram would have blocked")
6125 .expect("would have errored")
6126 .expect("track was closed");
6127 assert_eq!(datagram.sequence, 5);
6128
6129 let group = sub
6131 .next_group()
6132 .now_or_never()
6133 .expect("group would have blocked")
6134 .expect("would have errored")
6135 .expect("track was closed");
6136 assert_eq!(group.sequence, 3);
6137 }
6138
6139 fn drain_ordered(subscriber: &mut Ordered) -> Vec<u64> {
6141 let mut sequences = Vec::new();
6142 while let Some(Ok(Some(group))) = subscriber.next_group().now_or_never() {
6143 sequences.push(group.sequence);
6144 }
6145 sequences
6146 }
6147
6148 #[test]
6152 fn next_group_sheds_a_stale_backlog() {
6153 let mut producer = track_producer("test", None);
6154 for second in 0..4 {
6155 append_at(&mut producer, second * 1000);
6156 }
6157
6158 let mut subscriber = producer.subscribe(None).ordered();
6159 assert_eq!(drain_ordered(&mut subscriber), vec![3]);
6160
6161 let mut arrival = producer.subscribe(None);
6163 assert_eq!(drain(&mut arrival), vec![3]);
6164 }
6165
6166 #[test]
6170 fn next_group_keeps_a_backlog_inside_the_budget() {
6171 let mut producer = track_producer("test", None);
6172 for second in 0..4 {
6173 append_at(&mut producer, second * 1000);
6174 }
6175
6176 let mut subscriber = producer
6179 .subscribe(Subscription::default().with_max_age(Duration::from_millis(1500)))
6180 .ordered();
6181 assert_eq!(drain_ordered(&mut subscriber), vec![1, 2, 3]);
6182
6183 let mut replay = producer.subscribe(replay()).ordered();
6185 assert_eq!(drain_ordered(&mut replay), vec![0, 1, 2, 3]);
6186 }
6187
6188 #[test]
6191 fn next_group_keeps_a_group_with_no_proven_reach() {
6192 let mut producer = track_producer("test", None);
6193 append_at(&mut producer, 0); producer.append_group().unwrap(); append_at(&mut producer, 10_000); let mut subscriber = producer.subscribe(None).ordered();
6200 assert_eq!(drain_ordered(&mut subscriber), vec![0, 2]);
6201 }
6202
6203 #[tokio::test]
6204 async fn real_time_skips_older_sequences_with_equal_ages() {
6205 let mut producer = track_producer("test", None);
6206 append_at(&mut producer, 0);
6207 append_at(&mut producer, 0);
6208
6209 let mut subscriber = producer.subscribe(None);
6210 assert_eq!(drain(&mut subscriber), vec![1]);
6211 }
6212
6213 #[test]
6214 fn fetch_ignores_the_budget() {
6215 let mut producer = track_producer("test", None);
6216 for second in 0..4 {
6217 append_at(&mut producer, second * 1000);
6218 }
6219
6220 let consumer = producer.consume();
6223 let group = consumer.fetch_group(0, None).now_or_never().unwrap().unwrap();
6224 assert_eq!(group.sequence, 0);
6225 }
6226
6227 #[tokio::test]
6230 async fn fetched_group_is_not_a_live_drift_edge() {
6231 let mut producer = track_producer("test", None);
6232 let dynamic = producer.dynamic();
6233 let consumer = producer.consume();
6234 append_at(&mut producer, 0);
6235
6236 let pending = consumer.fetch_group(100, None);
6237 let req = dynamic
6238 .requested_group()
6239 .now_or_never()
6240 .expect("fetch request is ready")
6241 .unwrap();
6242 let mut fetched = req.accept(None).unwrap();
6243 fetched
6244 .write_frame(
6245 Timestamp::from_millis(100_000).unwrap(),
6246 bytes::Bytes::from_static(b"fetched"),
6247 )
6248 .unwrap();
6249 fetched.finish().unwrap();
6250 pending.await.unwrap();
6251
6252 let mut groups = producer.subscribe(None);
6253 assert_eq!(groups.assert_group().sequence, 0);
6254 groups.assert_no_group();
6255 }
6256
6257 #[test]
6262 fn an_evicted_live_edge_convicts_nothing() {
6263 let mut producer = track_producer("test", None);
6264 append_at(&mut producer, 0);
6265 let edge = append_at(&mut producer, 30_000);
6266
6267 let state = producer.state.read();
6268 let drift = Drift {
6269 budget: Duration::ZERO,
6270 edge: state.drift_edge(None, None, None),
6271 outer: None,
6272 successor: None,
6273 };
6274 assert!(
6275 state.is_stale(0, &drift.edge, drift.budget),
6276 "stale against a live edge"
6277 );
6278 drop(state);
6279
6280 let slot = producer.modify().unwrap().lookup.remove(&edge).unwrap();
6282 let _ = slot.group.abort(Error::Evicted);
6283
6284 let state = producer.state.read();
6285 assert!(
6286 !state.is_stale(0, &drift.edge, drift.budget),
6287 "a vanished edge is no reason to drop what is left"
6288 );
6289 }
6290
6291 #[tokio::test]
6295 async fn an_evicted_outer_edge_convicts_nothing() {
6296 let mut producer = track_producer("a", None);
6297 let mut subscriber = producer.subscribe(Subscription::default().with_max_age(Duration::from_secs(1)));
6298 let mut open = producer.append_group().unwrap();
6299 open.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"a"))
6300 .unwrap();
6301 let mut group = subscriber.recv_group().await.unwrap().expect("group");
6302 assert!(group.read_frame().await.unwrap().is_some());
6303 append_at(&mut producer, 10);
6305
6306 let mut next = track_producer("b", None);
6307 append_at(&mut next, 30_000);
6308 let edge = append_at(&mut next, 30_010);
6309 subscriber.set_anchor(Anchor {
6310 cap: None,
6311 edge: next.consume().live_edge(None),
6312 successor: None,
6313 });
6314
6315 let mut control = group.clone();
6316 assert!(
6317 matches!(control.read_frame().now_or_never(), Some(Ok(None))),
6318 "the open group ends against the outer edge"
6319 );
6320
6321 let slot = next.modify().unwrap().lookup.remove(&edge).unwrap();
6323 let _ = slot.group.abort(Error::Evicted);
6324
6325 assert!(
6326 group.read_frame().now_or_never().is_none(),
6327 "a vanished outer edge is no reason to drop what is left"
6328 );
6329 assert!(!group.latency_expired());
6330 open.finish().unwrap();
6331 }
6332
6333 #[tokio::test]
6337 async fn an_evicted_successor_convicts_nothing() {
6338 let producer = track_producer("a", None);
6339 let mut subscriber = producer.subscribe(Subscription::default().with_max_age(Duration::from_secs(1)));
6340 let mut open = producer.append_group().unwrap();
6341 open.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"a"))
6342 .unwrap();
6343 let mut group = subscriber.recv_group().await.unwrap().expect("group");
6344 assert!(group.read_frame().await.unwrap().is_some());
6345
6346 let mut next = track_producer("b", None);
6347 let successor = append_at(&mut next, 10);
6349 append_at(&mut next, 30_000);
6350 let consumer = next.consume();
6351 subscriber.set_anchor(Anchor {
6352 cap: None,
6353 edge: consumer.live_edge(None),
6354 successor: consumer.served_start(successor, None),
6355 });
6356
6357 let mut control = group.clone();
6358 assert!(
6359 matches!(control.read_frame().now_or_never(), Some(Ok(None))),
6360 "the open group ends against the successor"
6361 );
6362
6363 let slot = next.modify().unwrap().lookup.remove(&successor).unwrap();
6364 let _ = slot.group.abort(Error::Evicted);
6365
6366 assert!(
6367 group.read_frame().now_or_never().is_none(),
6368 "a vanished successor is no reason to drop what is left"
6369 );
6370 assert!(!group.latency_expired());
6371 open.finish().unwrap();
6372 }
6373
6374 #[tokio::test]
6375 async fn a_lower_sequence_is_never_the_live_edge() {
6376 let producer = track_producer("test", None);
6377 let mut straggler = producer.create_group(0u64.into()).unwrap();
6382 straggler
6383 .write_frame(Timestamp::from_millis(60_000).unwrap(), bytes::Bytes::from_static(b"x"))
6384 .unwrap();
6385 straggler.finish().unwrap();
6386
6387 let mut rewound = producer.create_group(1u64.into()).unwrap();
6388 rewound
6389 .write_frame(Timestamp::from_millis(0).unwrap(), bytes::Bytes::from_static(b"x"))
6390 .unwrap();
6391 rewound.finish().unwrap();
6392
6393 let mut subscriber = producer.subscribe(replay());
6394 assert_eq!(drain(&mut subscriber), vec![0, 1]);
6395 }
6396
6397 #[test]
6398 fn a_requested_end_does_not_cap_the_live_edge() {
6399 let mut producer = track_producer("test", None);
6400 for second in 0..4 {
6401 append_at(&mut producer, second * 1000);
6402 }
6403
6404 let mut subscriber = producer.subscribe(Subscription::default().with_end(Position::after_group(1)));
6409 assert_eq!(drain(&mut subscriber), vec![3]);
6410 }
6411
6412 #[test]
6413 fn a_capped_subscriber_measures_drift_against_its_cap() {
6414 let mut producer = track_producer("test", None);
6415 append_at(&mut producer, 0);
6416 append_at(&mut producer, 1000);
6417
6418 let mut subscriber = producer.subscribe(Subscription::default().with_end(Position::after_group(0)));
6423 subscriber.set_groups(..1);
6424 assert_eq!(drain(&mut subscriber), vec![0]);
6425
6426 subscriber.set_groups(..);
6429 assert_eq!(drain(&mut subscriber), vec![1]);
6430 }
6431
6432 #[test]
6433 fn subscriber_control_updates_while_read_future_is_pending() {
6434 let producer = track_producer("test", None);
6435 let mut subscriber = producer.subscribe(None);
6436 let control = subscriber.control();
6437
6438 let mut recv = Box::pin(subscriber.recv_group());
6439 assert!(recv.as_mut().now_or_never().is_none());
6440
6441 control.update(Subscription::default().with_priority(7)).unwrap();
6442
6443 let aggregate = producer.subscription().expect("expected an active subscription");
6444 assert_eq!(aggregate.priority, 7);
6445 }
6446
6447 #[test]
6448 fn dropped_subscriber_leaves_no_ghost_in_aggregate() {
6449 let mut producer = track_producer("test", None);
6454 let a = producer.subscribe(Subscription::default().with_priority(5));
6455
6456 let waiter = kio::Waiter::noop();
6458 assert!(
6459 matches!(producer.poll_subscription_changed(&waiter), Poll::Ready(Ok(Some(_)))),
6460 "one live subscriber should aggregate to Some",
6461 );
6462
6463 drop(a);
6465
6466 assert!(
6468 matches!(producer.poll_subscription_changed(&waiter), Poll::Ready(Ok(None))),
6469 "a dropped subscriber must not linger in the aggregate",
6470 );
6471
6472 assert!(
6474 producer.subscription().is_none(),
6475 "snapshot must exclude a dropped subscriber",
6476 );
6477 }
6478
6479 #[test]
6480 fn dropped_subscriber_wakes_the_aggregate() {
6481 use std::sync::atomic::{AtomicBool, Ordering};
6488
6489 let mut producer = track_producer("test", None);
6490 let a = producer.subscribe(Subscription::default().with_priority(5));
6491
6492 let woken = Arc::new(AtomicBool::new(false));
6493 let waiter = kio::Waiter::new(futures::task::waker(Arc::new(FlagWake(woken.clone()))));
6494
6495 assert!(matches!(
6497 producer.poll_subscription_changed(&waiter),
6498 Poll::Ready(Ok(Some(_)))
6499 ));
6500 assert!(
6501 producer.poll_subscription_changed(&waiter).is_pending(),
6502 "the aggregate is unchanged, so this poll must park",
6503 );
6504 assert!(!woken.load(Ordering::SeqCst), "nothing happened yet");
6505
6506 drop(a);
6507 assert!(
6508 woken.load(Ordering::SeqCst),
6509 "the last subscriber leaving must wake the aggregate watcher",
6510 );
6511 }
6512
6513 #[test]
6514 fn widest_subscriber_update_wakes_the_aggregate() {
6515 use std::sync::atomic::{AtomicBool, Ordering};
6521
6522 let mut producer = track_producer("test", None);
6523 let _narrow = producer.subscribe(Subscription::default().with_end(Position::after_group(3)));
6524 let mut wide = producer.subscribe(Subscription::default());
6525
6526 let woken = Arc::new(AtomicBool::new(false));
6527 let waiter = kio::Waiter::new(futures::task::waker(Arc::new(FlagWake(woken.clone()))));
6528
6529 assert!(matches!(
6530 producer.poll_subscription_changed(&waiter),
6531 Poll::Ready(Ok(Some(_)))
6532 ));
6533 assert!(producer.poll_subscription_changed(&waiter).is_pending());
6534 assert!(!woken.load(Ordering::SeqCst), "nothing happened yet");
6535
6536 wide.update(Subscription::default().with_end(Position::after_group(5)))
6537 .unwrap();
6538 assert!(
6539 woken.load(Ordering::SeqCst),
6540 "the widest subscriber changing must wake the aggregate watcher",
6541 );
6542 match producer.poll_subscription_changed(&waiter) {
6543 Poll::Ready(Ok(Some(sub))) => assert_eq!(sub.end, Position::after_group(5)),
6544 other => panic!("expected the narrowed aggregate, got {other:?}"),
6545 }
6546 }
6547
6548 struct FlagWake(Arc<std::sync::atomic::AtomicBool>);
6550
6551 impl futures::task::ArcWake for FlagWake {
6552 fn wake_by_ref(arc_self: &Arc<Self>) {
6553 arc_self.0.store(true, std::sync::atomic::Ordering::SeqCst);
6554 }
6555 }
6556
6557 #[tokio::test]
6558 async fn out_of_order_max_sequence_at_front() {
6559 let producer = track_producer("test", None);
6560
6561 producer.create_group(group::Info { sequence: 5 }).unwrap();
6563 producer.create_group(group::Info { sequence: 3 }).unwrap();
6564 producer.create_group(group::Info { sequence: 4 }).unwrap();
6565
6566 {
6568 let state = producer.state.read();
6569 assert_eq!(state.max_sequence, Some(5));
6570 }
6571
6572 crate::model::clock::advance(cache::DEFAULT_EXPIRY + Duration::from_secs(1));
6574
6575 producer.append_group().unwrap(); {
6581 let state = producer.state.read();
6582 assert_eq!(live_groups(&state), 1);
6583 assert_eq!(first_live_sequence(&state), 6);
6584 assert!(!state.lookup.contains_key(&3));
6585 assert!(!state.lookup.contains_key(&4));
6586 assert!(!state.lookup.contains_key(&5));
6587 assert!(state.lookup.contains_key(&6));
6588 }
6589 }
6590
6591 #[tokio::test]
6592 async fn max_sequence_at_front_blocks_trim() {
6593 let producer = track_producer("test", None);
6594
6595 producer.create_group(group::Info { sequence: 5 }).unwrap();
6597
6598 crate::model::clock::advance(cache::DEFAULT_EXPIRY + Duration::from_secs(1));
6599
6600 producer.create_group(group::Info { sequence: 3 }).unwrap();
6602
6603 {
6606 let state = producer.state.read();
6607 assert_eq!(live_groups(&state), 2);
6608 assert_eq!(state.offset, 0);
6609 }
6610
6611 crate::model::clock::advance(cache::DEFAULT_EXPIRY + Duration::from_secs(1));
6613
6614 producer.create_group(group::Info { sequence: 2 }).unwrap();
6616
6617 {
6622 let state = producer.state.read();
6623 assert_eq!(live_groups(&state), 2);
6624 assert_eq!(state.offset, 0);
6625 assert!(state.lookup.contains_key(&5));
6626 assert!(!state.lookup.contains_key(&3));
6627 assert!(state.lookup.contains_key(&2));
6628 }
6629
6630 let mut consumer = producer.subscribe(None);
6632 let group = consumer.assert_group();
6633 assert_eq!(group.sequence, 5);
6635 }
6636
6637 #[tokio::test]
6638 async fn abort_clears_cached_groups() {
6639 let producer = track_producer("test", None);
6640 producer.append_group().unwrap();
6641 producer.append_group().unwrap();
6642
6643 let mut consumer = producer.subscribe(None);
6645 assert_eq!(live_groups(&producer.state.read()), 2);
6646
6647 producer.clone().abort(Error::Cancel).unwrap();
6648
6649 {
6650 let state = producer.state.read();
6651 assert!(state.lookup.is_empty(), "cached groups should be dropped on abort");
6652 assert!(state.arrival.is_empty());
6653 assert!(state.evict.is_empty());
6654 }
6655
6656 let result = consumer.recv_group().now_or_never().expect("should not block");
6658 assert!(matches!(result, Err(Error::Cancel)));
6659 }
6660
6661 #[tokio::test]
6662 async fn drop_unfinished_clears_cached_groups() {
6663 let producer = track_producer("test", None);
6664 let writer = producer.clone();
6665 writer.append_group().unwrap();
6666
6667 let mut consumer = producer.subscribe(None);
6669 assert_eq!(live_groups(&producer.state.read()), 1);
6670
6671 drop(writer);
6673 drop(producer);
6674
6675 let result = consumer.recv_group().now_or_never().expect("should not block");
6676 assert!(matches!(result, Err(Error::Dropped)));
6677 }
6678
6679 #[tokio::test]
6680 async fn drop_after_abort_does_not_warn() {
6681 let warns = count_drop_warnings("track::Producer dropped without finish", || {
6684 let producer = track_producer("test", None);
6685 let keep = producer.clone();
6686 let writer = producer.clone();
6687 let group = writer.append_group().unwrap();
6688 group.finish().unwrap();
6689 let _consumer = producer.subscribe(None);
6690 writer.abort(Error::Cancel).unwrap();
6691 drop(keep);
6692 });
6693 assert_eq!(warns, 0, "abort-then-drop must not emit unfinished-producer WARN");
6694 }
6695
6696 #[tokio::test]
6697 async fn drop_unfinished_warns() {
6698 let warns = count_drop_warnings("track::Producer dropped without finish", || {
6699 let producer = track_producer("test", None);
6700 let writer = producer.clone();
6701 writer.append_group().unwrap();
6702 let _consumer = producer.subscribe(None);
6703 drop(writer);
6704 drop(producer);
6705 });
6706 assert!(warns >= 1, "unfinished drop must emit unfinished-producer WARN");
6707 }
6708
6709 #[tokio::test]
6710 async fn drop_finished_keeps_cached_groups() {
6711 let producer = track_producer("test", None);
6712 producer.append_group().unwrap();
6713 producer.finish().unwrap();
6714
6715 let mut consumer = producer.subscribe(None);
6716 drop(producer);
6717
6718 assert_eq!(consumer.assert_group().sequence, 0);
6720 let done = consumer.recv_group().now_or_never().expect("should not block").unwrap();
6721 assert!(done.is_none(), "consumer should drain then see clean finish");
6722 }
6723
6724 #[tokio::test]
6725 async fn cached_groups_preserve_arrival_order() {
6726 let producer = track_producer("test", None);
6727 producer.create_group(group::Info { sequence: 5 }).unwrap();
6728 producer.create_group(group::Info { sequence: 3 }).unwrap();
6729
6730 let groups = producer.consume().cached_groups();
6731 let sequences: Vec<u64> = groups.iter().map(|(group, _)| group.sequence).collect();
6732 assert_eq!(
6733 sequences,
6734 vec![5, 3],
6735 "warm snapshot must follow arrival, not sequence order"
6736 );
6737 }
6738
6739 #[test]
6740 fn append_finish_cannot_be_rewritten() {
6741 let producer = track_producer("test", None);
6742
6743 assert!(producer.finish().is_ok());
6745 assert!(producer.finish().is_err());
6746 assert!(producer.append_group().is_err());
6747 }
6748
6749 #[test]
6750 fn finish_after_groups() {
6751 let producer = track_producer("test", None);
6752
6753 producer.append_group().unwrap();
6754 assert!(producer.finish().is_ok());
6755 assert!(producer.finish().is_err());
6756 assert!(producer.append_group().is_err());
6757 }
6758
6759 #[test]
6760 fn finish_at_rejects_a_boundary_at_or_below_the_live_edge() {
6761 let mut producer = track_producer("test", None);
6762 producer.create_group(group::Info { sequence: 5 }).unwrap();
6763
6764 assert!(producer.finish_at(4).is_err());
6767 assert!(producer.finish_at(5).is_err());
6768 assert!(producer.finish_at(6).is_ok());
6769
6770 {
6771 let state = producer.state.read();
6772 assert_eq!(state.final_sequence, Some(6));
6773 }
6774
6775 assert!(producer.finish_at(6).is_err());
6777 assert!(producer.create_group(group::Info { sequence: 4 }).is_ok());
6778 assert!(producer.create_group(group::Info { sequence: 6 }).is_err());
6779 }
6780
6781 #[test]
6782 fn final_sequence_reports_the_declared_boundary() {
6783 let mut producer = track_producer("test", None);
6784 assert_eq!(producer.final_sequence(), None);
6785
6786 producer.create_group(group::Info { sequence: 5 }).unwrap();
6787 assert_eq!(producer.final_sequence(), None, "a group does not declare a boundary");
6788
6789 producer.finish_at(9).unwrap();
6790 assert_eq!(producer.final_sequence(), Some(9));
6791
6792 assert!(producer.finish().is_err());
6794 }
6795
6796 #[test]
6797 fn final_sequence_reports_the_live_edge_after_finish() {
6798 let producer = track_producer("test", None);
6799 producer.create_group(group::Info { sequence: 5 }).unwrap();
6800 producer.finish().unwrap();
6801 assert_eq!(producer.final_sequence(), Some(6));
6802 }
6803
6804 #[tokio::test]
6805 async fn finish_at_declares_a_future_boundary() {
6806 let mut producer = track_producer("test", None);
6807 producer.create_group(group::Info { sequence: 5 }).unwrap();
6808
6809 producer.finish_at(7).unwrap();
6811
6812 let mut consumer = producer.subscribe(None);
6813 assert_eq!(consumer.assert_group().sequence, 5);
6814
6815 let boundary = consumer
6818 .finished()
6819 .now_or_never()
6820 .expect("boundary is known immediately")
6821 .expect("would have errored");
6822 assert_eq!(boundary, 7);
6823 assert!(
6824 consumer.recv_group().now_or_never().is_none(),
6825 "should wait for the outstanding group"
6826 );
6827
6828 producer.create_group(group::Info { sequence: 6 }).unwrap();
6830 assert_eq!(consumer.assert_group().sequence, 6);
6831 let done = consumer
6832 .recv_group()
6833 .now_or_never()
6834 .expect("should not block")
6835 .expect("would have errored");
6836 assert!(done.is_none(), "track completes once the boundary is reached");
6837 }
6838
6839 #[tokio::test]
6840 async fn recv_group_finishes_without_waiting_for_gaps() {
6841 let producer = track_producer("test", None);
6842 producer.create_group(group::Info { sequence: 1 }).unwrap();
6843 producer.finish().unwrap();
6844
6845 let mut consumer = producer.subscribe(None);
6846 assert_eq!(consumer.assert_group().sequence, 1);
6847
6848 let done = consumer
6849 .recv_group()
6850 .now_or_never()
6851 .expect("should not block")
6852 .expect("would have errored");
6853 assert!(done.is_none(), "track should finish without waiting for gaps");
6854 }
6855
6856 #[tokio::test]
6857 async fn next_group_skips_late_arrivals() {
6858 let producer = track_producer("test", None);
6859 let mut consumer = producer.subscribe(None).ordered();
6860
6861 producer.create_group(group::Info { sequence: 5 }).unwrap();
6863 let group = consumer
6864 .next_group()
6865 .now_or_never()
6866 .expect("should not block")
6867 .expect("would have errored")
6868 .expect("track should not be closed");
6869 assert_eq!(group.sequence, 5);
6870
6871 producer.create_group(group::Info { sequence: 3 }).unwrap();
6873 producer.create_group(group::Info { sequence: 4 }).unwrap();
6875 producer.create_group(group::Info { sequence: 7 }).unwrap();
6877
6878 let group = consumer
6879 .next_group()
6880 .now_or_never()
6881 .expect("should not block")
6882 .expect("would have errored")
6883 .expect("track should not be closed");
6884 assert_eq!(group.sequence, 7);
6885
6886 assert!(
6888 consumer.next_group().now_or_never().is_none(),
6889 "should block waiting for a higher sequence"
6890 );
6891 }
6892
6893 #[tokio::test]
6894 async fn next_group_returns_arrivals_in_order() {
6895 let producer = track_producer("test", None);
6896 let mut consumer = producer.subscribe(replay()).ordered();
6897
6898 producer.create_group(group::Info { sequence: 3 }).unwrap();
6900 producer.create_group(group::Info { sequence: 5 }).unwrap();
6901
6902 let group = consumer
6903 .next_group()
6904 .now_or_never()
6905 .expect("should not block")
6906 .expect("would have errored")
6907 .expect("track should not be closed");
6908 assert_eq!(group.sequence, 3);
6909
6910 let group = consumer
6911 .next_group()
6912 .now_or_never()
6913 .expect("should not block")
6914 .expect("would have errored")
6915 .expect("track should not be closed");
6916 assert_eq!(group.sequence, 5);
6917 }
6918
6919 #[tokio::test]
6920 async fn ordered_and_arrival_cursors_are_independent() {
6921 let producer = track_producer("test", None);
6922 let mut ordered = producer.subscribe(replay()).ordered();
6923 let mut arrival = producer.subscribe(replay());
6924
6925 producer.create_group(group::Info { sequence: 5 }).unwrap();
6927 producer.create_group(group::Info { sequence: 3 }).unwrap();
6928
6929 let group = ordered
6932 .next_group()
6933 .now_or_never()
6934 .expect("should not block")
6935 .expect("would have errored")
6936 .expect("track should not be closed");
6937 assert_eq!(group.sequence, 3);
6938
6939 assert_eq!(arrival.assert_group().sequence, 5);
6941 }
6942
6943 #[tokio::test]
6944 async fn end_at_caps_next_group() {
6945 let producer = track_producer("test", None);
6946 let mut consumer = producer.subscribe(replay()).ordered();
6947
6948 for s in 0..6 {
6949 producer.create_group(group::Info { sequence: s }).unwrap();
6950 }
6951
6952 consumer.set_groups(..3);
6953
6954 assert_eq!(
6956 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6957 0
6958 );
6959 assert_eq!(
6960 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6961 1
6962 );
6963 assert_eq!(
6964 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6965 2
6966 );
6967
6968 assert!(
6970 consumer.next_group().now_or_never().is_none(),
6971 "capped consumer must block instead of returning out-of-range groups"
6972 );
6973 }
6974
6975 #[tokio::test]
6976 async fn end_at_release_drains_cached_groups() {
6977 let producer = track_producer("test", None);
6978 let mut consumer = producer.subscribe(replay()).ordered();
6979
6980 for s in 0..6 {
6981 producer.create_group(group::Info { sequence: s }).unwrap();
6982 }
6983
6984 consumer.set_groups(..2);
6985 assert_eq!(
6986 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6987 0
6988 );
6989 assert_eq!(
6990 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6991 1
6992 );
6993 assert!(consumer.next_group().now_or_never().is_none(), "capped at 2");
6994
6995 consumer.set_groups(..5);
6997 assert_eq!(
6998 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6999 2
7000 );
7001 assert_eq!(
7002 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
7003 3
7004 );
7005 assert_eq!(
7006 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
7007 4
7008 );
7009 assert!(consumer.next_group().now_or_never().is_none(), "capped at 5");
7010
7011 consumer.set_groups(..);
7013 assert_eq!(
7014 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
7015 5
7016 );
7017 assert!(consumer.next_group().now_or_never().is_none(), "no more groups");
7018 }
7019
7020 #[tokio::test]
7021 async fn end_at_lower_than_cursor_parks_consumer() {
7022 let producer = track_producer("test", None);
7023 let mut consumer = producer.subscribe(replay()).ordered();
7024
7025 for s in 0..3 {
7026 producer.create_group(group::Info { sequence: s }).unwrap();
7027 }
7028
7029 assert_eq!(
7031 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
7032 0
7033 );
7034 assert_eq!(
7035 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
7036 1
7037 );
7038 assert_eq!(
7039 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
7040 2
7041 );
7042
7043 consumer.set_groups(..2);
7045 producer.create_group(group::Info { sequence: 3 }).unwrap();
7046 producer.create_group(group::Info { sequence: 4 }).unwrap();
7047 assert!(
7048 consumer.next_group().now_or_never().is_none(),
7049 "cap is below cursor; nothing returnable until cap rises"
7050 );
7051
7052 consumer.set_groups(..);
7054 assert_eq!(
7055 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
7056 3
7057 );
7058 assert_eq!(
7059 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
7060 4
7061 );
7062 }
7063
7064 #[tokio::test]
7065 async fn end_at_toggling_around_late_arrivals() {
7066 let producer = track_producer("test", None);
7067 let mut consumer = producer.subscribe(replay()).ordered();
7068
7069 consumer.set_groups(..6);
7070
7071 producer.create_group(group::Info { sequence: 2 }).unwrap();
7073 producer.create_group(group::Info { sequence: 5 }).unwrap();
7074 producer.create_group(group::Info { sequence: 3 }).unwrap();
7075 producer.create_group(group::Info { sequence: 8 }).unwrap();
7077 producer.create_group(group::Info { sequence: 4 }).unwrap();
7078
7079 assert_eq!(
7081 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
7082 2
7083 );
7084 assert_eq!(
7085 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
7086 3
7087 );
7088 assert_eq!(
7089 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
7090 4
7091 );
7092 assert_eq!(
7093 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
7094 5
7095 );
7096 assert!(consumer.next_group().now_or_never().is_none());
7098
7099 consumer.set_groups(..11);
7101 assert_eq!(
7102 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
7103 8
7104 );
7105 }
7106
7107 #[tokio::test]
7111 async fn end_at_parks_recv_group() {
7112 let producer = track_producer("test", None);
7113 let mut consumer = producer.subscribe(replay());
7114
7115 for s in 0..3 {
7116 producer.create_group(group::Info { sequence: s }).unwrap();
7117 }
7118
7119 consumer.set_groups(..2);
7120 assert_eq!(
7121 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
7122 0
7123 );
7124 assert_eq!(
7125 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
7126 1
7127 );
7128 assert!(consumer.recv_group().now_or_never().is_none(), "capped at 2");
7129
7130 producer.finish().unwrap();
7132 assert!(
7133 consumer.recv_group().now_or_never().is_none(),
7134 "still parked after finish"
7135 );
7136
7137 consumer.set_groups(..);
7138 assert_eq!(
7139 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
7140 2
7141 );
7142 assert!(
7143 matches!(consumer.recv_group().now_or_never(), Some(Ok(None))),
7144 "finished once the parked group drains"
7145 );
7146 }
7147
7148 #[tokio::test]
7151 async fn recv_group_serves_arrivals_behind_the_cap() {
7152 let producer = track_producer("test", None);
7153 let mut consumer = producer.subscribe(replay());
7154
7155 consumer.set_groups(..2);
7156
7157 producer.create_group(group::Info { sequence: 2 }).unwrap();
7159 producer.create_group(group::Info { sequence: 0 }).unwrap();
7160 producer.create_group(group::Info { sequence: 1 }).unwrap();
7161
7162 assert_eq!(
7163 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
7164 0
7165 );
7166 assert_eq!(
7167 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
7168 1
7169 );
7170 assert!(consumer.recv_group().now_or_never().is_none(), "capped at 2");
7171
7172 consumer.set_groups(..3);
7173 assert_eq!(
7174 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
7175 2
7176 );
7177 }
7178
7179 #[tokio::test]
7180 async fn group_ranges_preserve_the_floor_when_the_cap_changes() {
7181 let producer = track_producer("test", None);
7182 let mut consumer = producer.subscribe(None);
7183 consumer.set_groups(2..=2);
7184 producer.create_group(group::Info { sequence: 2 }).unwrap();
7185 assert_eq!(consumer.recv_group().await.unwrap().unwrap().sequence, 2);
7186 consumer.set_groups(..4);
7187 producer.create_group(group::Info { sequence: 1 }).unwrap();
7188 producer.create_group(group::Info { sequence: 3 }).unwrap();
7189 assert_eq!(consumer.recv_group().await.unwrap().unwrap().sequence, 3);
7190 consumer.set_groups(0..=4);
7191 producer.create_group(group::Info { sequence: 0 }).unwrap();
7192 producer.create_group(group::Info { sequence: 4 }).unwrap();
7193 assert_eq!(consumer.recv_group().await.unwrap().unwrap().sequence, 4);
7194 }
7195
7196 #[tokio::test]
7199 async fn start_at_drops_parked_recv_groups() {
7200 let producer = track_producer("test", None);
7201 let mut consumer = producer.subscribe(None);
7202
7203 consumer.set_groups(..1);
7204 producer.create_group(group::Info { sequence: 1 }).unwrap();
7205 assert!(
7206 consumer.recv_group().now_or_never().is_none(),
7207 "group 1 parked at the cap"
7208 );
7209
7210 consumer.start_at(2);
7211 consumer.set_groups(..);
7212 producer.create_group(group::Info { sequence: 2 }).unwrap();
7213 assert_eq!(
7214 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
7215 2,
7216 "the overtaken parked group is dropped, not re-offered"
7217 );
7218 }
7219
7220 #[tokio::test]
7224 async fn evicted_parked_recv_groups_are_dropped() {
7225 let producer = track_producer("test", None);
7226 let mut consumer = producer.subscribe(None);
7227
7228 producer.create_group(group::Info { sequence: 0 }).unwrap();
7229 assert_eq!(
7230 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
7231 0
7232 );
7233
7234 consumer.set_groups(..1);
7235 let straggler = producer.create_group(group::Info { sequence: 1 }).unwrap();
7236 assert!(
7237 consumer.recv_group().now_or_never().is_none(),
7238 "group 1 parked at the cap"
7239 );
7240
7241 straggler.abort(Error::Old).unwrap();
7243 producer.finish().unwrap();
7244
7245 consumer.set_groups(..);
7246 assert!(
7247 matches!(consumer.recv_group().now_or_never(), Some(Ok(None))),
7248 "a dead parked group must not be delivered or hold the stream open"
7249 );
7250 }
7251
7252 #[tokio::test]
7256 async fn evicted_parked_group_wakes_the_clean_end() {
7257 use std::sync::atomic::{AtomicUsize, Ordering};
7258 use std::task::{Context, Wake};
7259
7260 struct CountWaker(AtomicUsize);
7263 impl Wake for CountWaker {
7264 fn wake(self: std::sync::Arc<Self>) {
7265 self.0.fetch_add(1, Ordering::SeqCst);
7266 }
7267 }
7268
7269 let producer = track_producer("test", None);
7270 let mut consumer = producer.subscribe(None);
7271
7272 producer.create_group(group::Info { sequence: 0 }).unwrap();
7273 assert_eq!(
7274 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
7275 0
7276 );
7277
7278 consumer.set_groups(..1);
7279 let straggler = producer.create_group(group::Info { sequence: 1 }).unwrap();
7280 assert!(consumer.recv_group().now_or_never().is_none(), "parked at the cap");
7281 producer.finish().unwrap();
7282
7283 let counter = std::sync::Arc::new(CountWaker(AtomicUsize::new(0)));
7284 let waker = std::task::Waker::from(counter.clone());
7285 let mut cx = Context::from_waker(&waker);
7286 let mut fut = std::pin::pin!(consumer.recv_group());
7287 assert!(
7288 fut.as_mut().poll(&mut cx).is_pending(),
7289 "the parked group holds it open"
7290 );
7291
7292 straggler.abort(Error::Old).unwrap();
7293 assert!(counter.0.load(Ordering::SeqCst) > 0, "the eviction wakeup was lost");
7294 assert!(matches!(fut.as_mut().poll(&mut cx), Poll::Ready(Ok(None))));
7295 }
7296
7297 #[tokio::test]
7299 async fn end_at_zero_is_the_empty_range() {
7300 let producer = track_producer("test", None);
7301 let mut consumer = producer.subscribe(replay());
7302 producer.create_group(group::Info { sequence: 0 }).unwrap();
7303 producer.create_group(group::Info { sequence: 1 }).unwrap();
7304
7305 consumer.set_groups(..0);
7306 assert!(
7307 consumer.recv_group().now_or_never().is_none(),
7308 "empty cap delivers nothing"
7309 );
7310
7311 consumer.set_groups(..1);
7312 assert_eq!(
7313 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
7314 0
7315 );
7316 assert!(consumer.recv_group().now_or_never().is_none(), "group 1 stays parked");
7317 }
7318
7319 #[tokio::test]
7322 async fn empty_local_cap_holds_while_another_subscriber_requests_everything() {
7323 let producer = track_producer("test", None);
7324 let mut everything = producer.subscribe(replay());
7325 let mut empty = producer.subscribe(Subscription::default().with_end(Position::group(0)));
7326 empty.set_groups((Bound::Unbounded, Position::group(0).group_end()));
7327
7328 for s in 0..3 {
7329 producer.create_group(group::Info { sequence: s }).unwrap();
7330 }
7331
7332 assert_eq!(
7333 everything
7334 .recv_group()
7335 .now_or_never()
7336 .unwrap()
7337 .unwrap()
7338 .unwrap()
7339 .sequence,
7340 0
7341 );
7342 assert!(
7343 empty.recv_group().now_or_never().is_none(),
7344 "local empty cap must not ride the unbounded aggregate"
7345 );
7346
7347 empty.set_groups(..2);
7348 assert_eq!(empty.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence, 0);
7349 assert_eq!(empty.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence, 1);
7350 assert!(empty.recv_group().now_or_never().is_none(), "still capped at 2");
7351
7352 empty.set_groups(..);
7353 assert_eq!(empty.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence, 2);
7354 }
7355
7356 #[tokio::test]
7358 async fn end_at_frame_limited_last_group() {
7359 let producer = track_producer("test", None);
7360 let mut consumer = producer.subscribe(replay()).ordered();
7361 let end = Position::after(1, 1).unwrap();
7362 consumer.set_groups((Bound::Unbounded, end.group_end()));
7363
7364 for s in 0..3u64 {
7365 let mut group = producer.create_group(group::Info { sequence: s }).unwrap();
7366 for i in 0..3u8 {
7367 group.write_frame(Timestamp::ZERO, vec![i]).unwrap();
7368 }
7369 group.finish().unwrap();
7370 }
7371
7372 assert_eq!(
7373 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
7374 0
7375 );
7376 let mut last = consumer.next_group().now_or_never().unwrap().unwrap().unwrap();
7377 assert_eq!(last.sequence, 1);
7378 last.set_frames(..end.frame);
7379 assert_eq!(
7380 last.read_frame().now_or_never().unwrap().unwrap().unwrap().payload[0],
7381 0
7382 );
7383 assert_eq!(
7384 last.read_frame().now_or_never().unwrap().unwrap().unwrap().payload[0],
7385 1
7386 );
7387 assert!(
7388 last.read_frame().now_or_never().unwrap().unwrap().is_none(),
7389 "frame cap is exclusive"
7390 );
7391 assert!(
7392 consumer.next_group().now_or_never().is_none(),
7393 "group 2 is past the exclusive group cap"
7394 );
7395 }
7396
7397 #[tokio::test]
7399 async fn end_at_maximum_group_is_unbounded() {
7400 let producer = track_producer("test", None);
7401 let mut consumer = producer.subscribe(replay()).ordered();
7402 consumer.set_groups(..=u64::MAX);
7403
7404 producer.create_group(group::Info { sequence: 0 }).unwrap();
7405 producer.create_group(group::Info { sequence: u64::MAX }).unwrap();
7406
7407 assert_eq!(
7408 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
7409 0
7410 );
7411 assert_eq!(
7412 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
7413 u64::MAX
7414 );
7415 }
7416
7417 #[test]
7418 fn write_frame_rejects_an_oversized_frame_before_appending_its_group() {
7419 let mut producer = track_producer("test", None);
7420 let frame = bytes::Bytes::from(vec![0; group::MAX_CACHE_BYTES as usize + 1]);
7421
7422 assert!(matches!(
7423 producer.write_frame(Timestamp::ZERO, frame),
7424 Err(Error::FrameTooLarge)
7425 ));
7426 assert_eq!(producer.latest(), None, "the rejected frame did not publish a group");
7427 }
7428
7429 #[test]
7430 fn append_group_returns_bounds_exceeded_on_sequence_overflow() {
7431 let producer = track_producer("test", None);
7432 {
7433 let mut state = producer.state.write().ok().unwrap();
7434 state.max_sequence = Some(u64::MAX);
7435 }
7436
7437 assert!(matches!(producer.append_group(), Err(Error::BoundsExceeded(_))));
7438 }
7439
7440 #[tokio::test]
7441 async fn fetch_cache_hit() {
7442 let producer = track_producer("test", None);
7443
7444 let mut group = producer.append_group().unwrap(); group
7447 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hello"))
7448 .unwrap();
7449 group.finish().unwrap();
7450
7451 let dynamic = producer.dynamic();
7454 let consumer = producer.consume();
7455 assert!(consumer.peek_group(0).is_some());
7456 let mut g = consumer.fetch_group(0, None).await.unwrap();
7457 assert_eq!(g.sequence, 0);
7458 assert_eq!(&g.read_frame().await.unwrap().unwrap().payload[..], b"hello");
7459
7460 assert!(dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending());
7462 }
7463
7464 #[tokio::test]
7465 async fn fetch_miss_signals_dynamic() {
7466 let producer = track_producer("test", None);
7467 let dynamic = producer.dynamic();
7468 let consumer = producer.consume();
7469
7470 assert!(consumer.peek_group(5).is_none());
7474 let pending = consumer.fetch_group(5, group::Fetch::default().with_priority(7));
7475 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
7476
7477 let req = dynamic
7478 .requested_group()
7479 .now_or_never()
7480 .expect("should not block")
7481 .unwrap();
7482 assert_eq!(req.sequence(), 5);
7483 assert_eq!(req.priority(), 7);
7484
7485 let mut group = req.accept(None).unwrap();
7487 group
7488 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hi"))
7489 .unwrap();
7490 group.finish().unwrap();
7491
7492 let mut g = pending.await.unwrap();
7493 assert_eq!(g.sequence, 5);
7494 assert_eq!(&g.read_frame().await.unwrap().unwrap().payload[..], b"hi");
7495 }
7496
7497 #[tokio::test]
7498 async fn fetch_miss_rejects() {
7499 let producer = track_producer("test", None);
7500 let dynamic = producer.dynamic();
7501 let consumer = producer.consume();
7502
7503 let pending = consumer.fetch_group(5, None);
7504 let req = dynamic
7505 .requested_group()
7506 .now_or_never()
7507 .expect("should not block")
7508 .unwrap();
7509
7510 req.reject(Error::Cancel);
7511 assert!(matches!(pending.await, Err(Error::Cancel)));
7512 let fetch = producer.state.read().fetch.clone();
7513 assert!(fetch.read().is_empty());
7514 }
7515
7516 #[tokio::test]
7517 async fn fetch_miss_drop_rejects() {
7518 let producer = track_producer("test", None);
7519 let dynamic = producer.dynamic();
7520 let consumer = producer.consume();
7521
7522 let pending = consumer.fetch_group(5, None);
7523 let req = dynamic
7524 .requested_group()
7525 .now_or_never()
7526 .expect("should not block")
7527 .unwrap();
7528
7529 drop(req);
7530 assert!(matches!(pending.await, Err(Error::Dropped)));
7531 }
7532
7533 #[tokio::test]
7534 async fn fetch_reject_does_not_poison_retry() {
7535 let producer = track_producer("test", None);
7536 let dynamic = producer.dynamic();
7537 let consumer = producer.consume();
7538
7539 let pending = consumer.fetch_group(5, None);
7540 let req = dynamic
7541 .requested_group()
7542 .now_or_never()
7543 .expect("should not block")
7544 .unwrap();
7545 req.reject(Error::Cancel);
7546 assert!(matches!(pending.await, Err(Error::Cancel)));
7547
7548 let retry = consumer.fetch_group(5, None);
7549 let req = dynamic
7550 .requested_group()
7551 .now_or_never()
7552 .expect("should not block")
7553 .unwrap();
7554 let mut group = req.accept(None).unwrap();
7555 group
7556 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"retry"))
7557 .unwrap();
7558 group.finish().unwrap();
7559
7560 let mut group = retry.await.unwrap();
7561 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"retry");
7562 }
7563
7564 #[tokio::test]
7568 async fn fetch_ignores_a_group_that_starts_too_late() {
7569 let producer = track_producer("test", None);
7570 let dynamic = producer.dynamic();
7571 let consumer = producer.consume();
7572
7573 let mut group = producer.create_group(group::Info { sequence: 0 }).unwrap();
7575 group.start_at(3).unwrap();
7576 group
7577 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"tail"))
7578 .unwrap();
7579 group.finish().unwrap();
7580
7581 let fetch = consumer.fetch_group(0, group::Fetch::default().with_frame_start(3));
7583 let cached = fetch.now_or_never().expect("covered by the cache").unwrap();
7584 assert_eq!(cached.index(), 3);
7585
7586 let mut fetch = std::pin::pin!(consumer.fetch_group(0, None));
7589 assert!(
7590 futures::poll!(fetch.as_mut()).is_pending(),
7591 "must not answer from the tail"
7592 );
7593
7594 let request = dynamic.requested_group().await.unwrap();
7595 assert_eq!((request.sequence(), request.frame_start()), (0, 0));
7596
7597 let mut whole = request.accept(None).unwrap();
7599 whole
7600 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"head"))
7601 .unwrap();
7602 whole.finish().unwrap();
7603
7604 let mut served = fetch.await.unwrap();
7605 assert_eq!(served.index(), 0);
7606 assert_eq!(
7607 served.read_frame().await.unwrap().unwrap().payload,
7608 bytes::Bytes::from_static(b"head")
7609 );
7610 }
7611
7612 #[tokio::test]
7615 async fn fetch_widens_or_fails_cleanly() {
7616 let producer = track_producer("test", None);
7617 let dynamic = producer.dynamic();
7618 let consumer = producer.consume();
7619
7620 let _narrow = consumer.fetch_group(0, group::Fetch::default().with_frame_start(5));
7621 let mut narrow = std::pin::pin!(_narrow);
7622 assert!(futures::poll!(narrow.as_mut()).is_pending());
7623
7624 let _wider = consumer.fetch_group(0, group::Fetch::default().with_frame_start(2));
7626 let mut wider = std::pin::pin!(_wider);
7627 assert!(futures::poll!(wider.as_mut()).is_pending());
7628
7629 let request = dynamic.requested_group().await.unwrap();
7630 assert_eq!(request.frame_start(), 2, "widened while queued");
7631
7632 let _widest = consumer.fetch_group(0, group::Fetch::default().with_frame_start(0));
7635 let mut widest = std::pin::pin!(_widest);
7636 assert!(futures::poll!(widest.as_mut()).is_pending());
7637 assert_eq!(request.frame_start(), 2, "the in-flight range is already on the wire");
7638
7639 let mut group = request.accept(None).unwrap();
7640 group.start_at(2).unwrap();
7643 group
7644 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"from2"))
7645 .unwrap();
7646 group.finish().unwrap();
7647
7648 assert_eq!(narrow.await.unwrap().index(), 5);
7651 assert_eq!(wider.await.unwrap().index(), 2);
7652 assert!(matches!(widest.await, Err(Error::NotFound)));
7653 }
7654
7655 #[tokio::test]
7656 async fn fetch_coalesces_concurrent() {
7657 let producer = track_producer("test", None);
7658 let dynamic = producer.dynamic();
7659 let consumer = producer.consume();
7660
7661 let first = consumer.fetch_group(5, group::Fetch::default().with_priority(1));
7664 let second = consumer.fetch_group(5, group::Fetch::default().with_priority(7));
7665 assert!(kio::Pollable::poll(&*first, &kio::Waiter::noop()).is_pending());
7666
7667 let req = dynamic
7668 .requested_group()
7669 .now_or_never()
7670 .expect("should not block")
7671 .unwrap();
7672 assert_eq!(req.sequence(), 5);
7673 assert_eq!(req.priority(), 7);
7674 assert!(
7675 dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending(),
7676 "the second fetch queued a duplicate request"
7677 );
7678
7679 let third = consumer.fetch_group(5, None);
7681
7682 let mut group = req.accept(None).unwrap();
7684 group
7685 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hi"))
7686 .unwrap();
7687 group.finish().unwrap();
7688
7689 assert_eq!(first.await.unwrap().sequence, 5);
7690 assert_eq!(second.await.unwrap().sequence, 5);
7691 assert_eq!(third.await.unwrap().sequence, 5);
7692 }
7693
7694 #[tokio::test]
7695 async fn fetch_coalesced_reject_fails_all() {
7696 let producer = track_producer("test", None);
7697 let dynamic = producer.dynamic();
7698 let consumer = producer.consume();
7699
7700 let first = consumer.fetch_group(5, None);
7701 let second = consumer.fetch_group(5, None);
7702 let req = dynamic
7703 .requested_group()
7704 .now_or_never()
7705 .expect("should not block")
7706 .unwrap();
7707 req.reject(Error::Cancel);
7708
7709 assert!(matches!(first.await, Err(Error::Cancel)));
7710 assert!(matches!(second.await, Err(Error::Cancel)));
7711
7712 let retry = consumer.fetch_group(5, None);
7714 assert!(kio::Pollable::poll(&*retry, &kio::Waiter::noop()).is_pending());
7715 let req = dynamic
7716 .requested_group()
7717 .now_or_never()
7718 .expect("should not block")
7719 .unwrap();
7720 assert_eq!(req.sequence(), 5);
7721 }
7722
7723 #[tokio::test]
7724 async fn fetch_queued_fails_when_handlers_leave() {
7725 let producer = track_producer("test", None);
7726 let dynamic = producer.dynamic();
7727 let consumer = producer.consume();
7728
7729 let pending = consumer.fetch_group(5, None);
7731 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
7732 drop(dynamic);
7733 assert!(matches!(pending.await, Err(Error::NotFound)));
7734
7735 let fetch = producer.state.read().fetch.clone();
7737 assert!(fetch.read().is_empty());
7738 }
7739
7740 #[tokio::test]
7741 async fn fetch_miss_no_dynamic_not_found() {
7742 let producer = track_producer("test", None);
7745 producer.append_group().unwrap(); let consumer = producer.consume();
7747 assert!(matches!(consumer.fetch_group(5, None).await, Err(Error::NotFound)));
7748 }
7749
7750 #[tokio::test]
7751 async fn fetch_past_final_not_found() {
7752 let producer = track_producer("test", None);
7753 producer.append_group().unwrap(); producer.finish().unwrap(); let dynamic = producer.dynamic();
7759 let consumer = producer.consume();
7760 assert!(matches!(consumer.fetch_group(5, None).await, Err(Error::NotFound)));
7761
7762 assert!(dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending());
7764 }
7765
7766 fn pooled_producer(capacity: u64) -> (Producer, cache::Pool) {
7768 let config = cache::Config::default()
7769 .with_capacity(capacity)
7770 .with_expiry(cache::DEFAULT_EXPIRY);
7771 let pool = cache::Pool::new(config);
7772 let broadcast = broadcast::Info {
7773 pool: pool.clone(),
7774 ..Default::default()
7775 };
7776 let producer = Producer::new(Arc::new(broadcast), "test", None);
7777 (producer, pool)
7778 }
7779
7780 fn finished_group(producer: &mut Producer, size: usize) -> u64 {
7781 let mut group = producer.append_group().unwrap();
7782 group
7783 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; size]))
7784 .unwrap();
7785 group.finish().unwrap();
7786 group.sequence
7787 }
7788
7789 #[tokio::test]
7792 async fn debt_evicts_oldest_group() {
7793 let (mut producer, pool) = pooled_producer(10_000);
7795
7796 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let consumer = producer.consume();
7801 assert!(consumer.peek_group(0).is_none(), "oldest group is evicted");
7802 assert!(consumer.peek_group(2).is_some(), "latest group survives");
7803 assert!(
7806 pool.used() <= 2 * (10_000 + cache::ENTRY_OVERHEAD),
7807 "usage hovers near capacity: {}",
7808 pool.used()
7809 );
7810
7811 let mut subscriber = producer.subscribe(replay());
7813 assert!(subscriber.assert_group().sequence > 0, "evicted group is not delivered");
7814 }
7815
7816 #[tokio::test]
7818 async fn latest_group_never_evicted() {
7819 let (mut producer, pool) = pooled_producer(100);
7821 finished_group(&mut producer, 1000); assert!(pool.used() > 100, "the latest may exceed the budget");
7823
7824 finished_group(&mut producer, 1000); finished_group(&mut producer, 1000); let consumer = producer.consume();
7829 assert!(consumer.peek_group(0).is_none());
7830 let mut group = consumer.peek_group(2).expect("latest survives");
7831 assert_eq!(group.read_frame().await.unwrap().unwrap().payload.len(), 1000);
7832 }
7833
7834 #[tokio::test]
7838 async fn fetch_refresh_survives_eviction() {
7839 let (mut producer, _pool) = pooled_producer(10_000);
7840 let consumer = producer.consume();
7841
7842 finished_group(&mut producer, 3_000); crate::model::clock::advance(Duration::from_secs(1));
7844 finished_group(&mut producer, 3_000); crate::model::clock::advance(Duration::from_secs(1));
7846 finished_group(&mut producer, 3_000); crate::model::clock::advance(Duration::from_millis(500));
7848
7849 let mut fetched = consumer.fetch_group(0, None).await.unwrap();
7851 assert_eq!(fetched.read_frame().await.unwrap().unwrap().payload.len(), 3_000);
7852 crate::model::clock::advance(Duration::from_millis(500));
7853
7854 finished_group(&mut producer, 3_000); crate::model::clock::advance(Duration::from_secs(1));
7858 finished_group(&mut producer, 3_000); assert!(consumer.peek_group(0).is_some(), "refreshed group survives");
7861 assert!(consumer.peek_group(1).is_none(), "unread group is evicted instead");
7862 }
7863
7864 #[tokio::test]
7867 async fn eviction_aborts_readers() {
7868 let (mut producer, _pool) = pooled_producer(10_000);
7869 let mut subscriber = producer.subscribe(None);
7870
7871 finished_group(&mut producer, 10_000); let mut group0 = subscriber.assert_group();
7873
7874 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let read = group0.read_frame().await;
7878 assert!(matches!(read, Err(Error::Evicted)), "expected Evicted, got {read:?}");
7879 }
7880
7881 #[tokio::test]
7885 async fn small_writes_carry_debt() {
7886 let unit = 100 * cache::ENTRY_OVERHEAD;
7889 let (mut producer, pool) = pooled_producer(22 * unit);
7890 let consumer = producer.consume();
7891
7892 finished_group(&mut producer, 20 * unit as usize); for _ in 0..3 {
7897 finished_group(&mut producer, unit as usize);
7898 }
7899 assert!(consumer.peek_group(0).is_some(), "debt smaller than the victim carries");
7900
7901 for _ in 0..20 {
7903 finished_group(&mut producer, unit as usize);
7904 }
7905 assert!(
7906 consumer.peek_group(0).is_none(),
7907 "accumulated debt evicts the large group"
7908 );
7909 assert!(pool.used() <= 24 * unit, "usage hovers near capacity: {}", pool.used());
7912 }
7913
7914 #[tokio::test]
7918 async fn payment_capped_per_write() {
7919 let (mut producer, pool) = pooled_producer(1 << 40);
7920 for _ in 0..10 {
7921 finished_group(&mut producer, 1_000);
7922 }
7923
7924 pool.resize(100);
7926 let before = pool.used();
7927
7928 finished_group(&mut producer, 1_000);
7930
7931 let consumer = producer.consume();
7932 assert!(consumer.peek_group(0).is_none(), "the oldest groups are evicted");
7933 assert!(consumer.peek_group(1).is_none());
7934 assert!(consumer.peek_group(2).is_some(), "the backlog drains gradually");
7935 assert!(pool.used() > before - 4_000, "one write must not dump the backlog");
7936 }
7937
7938 #[tokio::test]
7942 async fn accept_preserves_write_accounting() {
7943 let config = cache::Config::default()
7944 .with_capacity(12_000)
7945 .with_expiry(cache::DEFAULT_EXPIRY);
7946 let pool = cache::Pool::new(config);
7947 let broadcast = broadcast::Info {
7948 pool: pool.clone(),
7949 ..Default::default()
7950 };
7951 let request = Request::new(Arc::new(broadcast), "test");
7952 let dynamic = request.dynamic();
7953 let consumer = request.consume();
7954
7955 let pending = consumer.fetch_group(0, None);
7957 let req = dynamic
7958 .requested_group()
7959 .now_or_never()
7960 .expect("should not block")
7961 .unwrap();
7962 let mut backfill = req.accept(None).unwrap();
7963 pending.await.unwrap();
7964 backfill
7965 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 30_000]))
7966 .unwrap();
7967
7968 let producer = request.accept(None);
7971 producer.append_group().unwrap().finish().unwrap();
7972 producer.append_group().unwrap().finish().unwrap();
7973
7974 assert!(
7975 producer.consume().peek_group(0).is_none(),
7976 "pre-accept backfill growth is reclaimed after accept"
7977 );
7978 assert!(pool.used() <= 13_000, "usage converges: {}", pool.used());
7979 }
7980
7981 #[tokio::test]
7984 async fn recreated_sequence_bounds_eviction_hints() {
7985 let (producer, _pool) = pooled_producer(1 << 40);
7986 producer.create_group(5u64.into()).unwrap().finish().unwrap();
7987
7988 for _ in 0..200 {
7989 let group = producer.create_group(1u64.into()).unwrap();
7990 group.abort(Error::Cancel).unwrap();
7991 }
7992
7993 let state = producer.state.read();
7994 assert!(
7995 state.evict.len() <= 2 * state.lookup.len() + EVICT_SLACK,
7996 "stale hints are compacted: {} entries for {} slots",
7997 state.evict.len(),
7998 state.lookup.len()
7999 );
8000 }
8001
8002 #[tokio::test]
8005 async fn same_tick_write_outranks_inserted() {
8006 let unit = 100 * cache::ENTRY_OVERHEAD;
8009 let (mut producer, _pool) = pooled_producer(10 * unit);
8011
8012 producer.append_group().unwrap().finish().unwrap(); finished_group(&mut producer, 3 * unit as usize); finished_group(&mut producer, 3 * unit as usize); finished_group(&mut producer, 3 * unit as usize); finished_group(&mut producer, 3 * unit as usize); let consumer = producer.consume();
8019 assert!(consumer.peek_group(0).is_none(), "insert-only content pays first");
8020 assert!(consumer.peek_group(1).is_some(), "same-tick written content survives");
8021 }
8022
8023 #[tokio::test]
8026 async fn frame_only_writer_pays() {
8027 let (producer, pool) = pooled_producer(2_000);
8028 let mut demoted = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); demoted
8034 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 300_000]))
8035 .unwrap();
8036
8037 assert!(
8038 pool.used() <= 5_000,
8039 "the frame write settled the debt: {}",
8040 pool.used()
8041 );
8042 assert!(matches!(demoted.finish(), Err(Error::Evicted)));
8043 }
8044
8045 #[tokio::test]
8048 async fn each_track_owns_its_account() {
8049 let broadcast = Arc::new(broadcast::Info::default());
8050 let info = Info::default();
8051 let a = Producer::new(broadcast.clone(), "a", info.clone());
8052 let b = Producer::new(broadcast, "b", info);
8053
8054 let a = a.state.read().cache.clone();
8055 let b = b.state.read().cache.clone();
8056 assert!(!Arc::ptr_eq(&a, &b), "each track owns its account");
8057 }
8058
8059 #[tokio::test]
8062 async fn a_dynamic_defers_teardown() {
8063 let (mut producer, pool) = pooled_producer(1 << 40);
8064 let dynamic = producer.dynamic();
8065 finished_group(&mut producer, 100);
8066
8067 drop(producer);
8068 assert!(pool.used() > 0, "the handler still serves the cache");
8069
8070 drop(dynamic);
8071 assert_eq!(pool.used(), 0, "the last handle tears it down");
8072 }
8073
8074 #[tokio::test]
8080 async fn finished_track_frees_its_cache() {
8081 let (mut producer, pool) = pooled_producer(1 << 40);
8082 finished_group(&mut producer, 100);
8083 producer.finish().unwrap();
8084
8085 let state = producer.state.downgrade();
8086 drop(producer);
8087
8088 assert!(state.upgrade().is_none(), "the track state is freed");
8089 assert_eq!(pool.used(), 0, "so are its cached bytes");
8090 }
8091
8092 #[tokio::test]
8096 async fn teardown_ignores_a_settling_group() {
8097 let (mut producer, pool) = pooled_producer(1 << 40);
8098 finished_group(&mut producer, 100);
8099
8100 let settling = producer.state.downgrade().upgrade().expect("open");
8102 drop(producer);
8103
8104 assert_eq!(pool.used(), 0, "the abrupt teardown still released the cache");
8105 drop(settling);
8106 }
8107
8108 #[tokio::test]
8111 async fn cached_group_outlives_its_track() {
8112 let (mut producer, pool) = pooled_producer(1 << 40);
8113 let sequence = finished_group(&mut producer, 100);
8114 let group = producer.consume().peek_group(sequence).expect("cached");
8115 producer.finish().unwrap();
8116
8117 let state = producer.state.downgrade();
8118 drop(producer);
8119 assert!(state.upgrade().is_none(), "the track state is freed");
8120 assert!(pool.used() > 0, "the retained group keeps its own bytes");
8121
8122 drop(group);
8123 assert_eq!(pool.used(), 0, "which it releases when dropped");
8124 }
8125
8126 #[tokio::test]
8130 async fn pre_accept_backfill_settles_late_writes() {
8131 let config = cache::Config::default()
8132 .with_capacity(2_000)
8133 .with_expiry(cache::DEFAULT_EXPIRY);
8134 let pool = cache::Pool::new(config);
8135 let broadcast = broadcast::Info {
8136 pool: pool.clone(),
8137 ..Default::default()
8138 };
8139 let request = Request::new(Arc::new(broadcast), "test");
8140 let dynamic = request.dynamic();
8141 let consumer = request.consume();
8142
8143 let pending = consumer.fetch_group(0, None);
8145 let req = dynamic
8146 .requested_group()
8147 .now_or_never()
8148 .expect("should not block")
8149 .unwrap();
8150 let mut backfill = req.accept(None).unwrap();
8151 pending.await.unwrap();
8152
8153 let producer = request.accept(None);
8155 producer.append_group().unwrap().finish().unwrap();
8156
8157 backfill
8160 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 300_000]))
8161 .unwrap();
8162
8163 assert!(
8164 pool.used() <= 5_000,
8165 "the frame write settled the debt: {}",
8166 pool.used()
8167 );
8168 }
8169
8170 #[tokio::test]
8174 async fn write_restarts_retention_clock() {
8175 let (producer, _pool) = pooled_producer(1 << 40);
8176 let mut straggler = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); crate::model::clock::advance(cache::DEFAULT_EXPIRY + Duration::from_secs(1));
8181 straggler
8182 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 100]))
8183 .unwrap();
8184 producer.append_group().unwrap().finish().unwrap(); let consumer = producer.consume();
8187 assert!(consumer.peek_group(0).is_some(), "the write restarted the clock");
8188
8189 crate::model::clock::advance(cache::DEFAULT_EXPIRY + Duration::from_secs(1));
8191 producer.append_group().unwrap().finish().unwrap(); assert!(consumer.peek_group(0).is_none(), "idle content still expires");
8193 }
8194
8195 #[tokio::test]
8198 async fn refreshed_front_does_not_starve_expiry() {
8199 let (producer, _pool) = pooled_producer(1 << 40);
8200 let dynamic = producer.dynamic();
8201 let consumer = producer.consume();
8202
8203 producer.create_group(10u64.into()).unwrap().finish().unwrap();
8204 for sequence in 1..=5u64 {
8205 let pending = consumer.fetch_group(sequence, None);
8206 let req = dynamic
8207 .requested_group()
8208 .now_or_never()
8209 .expect("should not block")
8210 .unwrap();
8211 let mut group = req.accept(None).unwrap();
8212 group
8213 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 100]))
8214 .unwrap();
8215 group.finish().unwrap();
8216 pending.await.unwrap();
8217 }
8218
8219 crate::model::clock::advance(cache::DEFAULT_EXPIRY + Duration::from_secs(1));
8222 for sequence in 1..=4u64 {
8223 consumer.fetch_group(sequence, None).await.unwrap();
8224 }
8225
8226 for _ in 0..3 {
8228 producer.append_group().unwrap().finish().unwrap();
8229 }
8230 assert!(consumer.peek_group(5).is_none(), "expired backfill is reclaimed");
8231 assert!(consumer.peek_group(1).is_some(), "refreshed backfill survives");
8232 }
8233
8234 #[tokio::test]
8237 async fn recreated_sequence_delivered_once() {
8238 let (producer, _pool) = pooled_producer(1 << 40);
8239
8240 producer.create_group(0u64.into()).unwrap().finish().unwrap();
8241 let aborted = producer.create_group(1u64.into()).unwrap();
8242 aborted.abort(Error::Cancel).unwrap();
8243 producer.create_group(2u64.into()).unwrap().finish().unwrap();
8244 producer.create_group(1u64.into()).unwrap().finish().unwrap();
8245
8246 let mut subscriber = producer.subscribe(replay());
8247 assert_eq!(subscriber.assert_group().sequence, 0);
8248 assert_eq!(subscriber.assert_group().sequence, 2);
8249 assert_eq!(
8250 subscriber.assert_group().sequence,
8251 1,
8252 "replacement arrives at its own position"
8253 );
8254 subscriber.assert_no_group();
8255 }
8256
8257 #[tokio::test]
8261 async fn datagrams_do_not_block_eviction() {
8262 let (mut producer, pool) = pooled_producer(1_000);
8263 for _ in 0..10 {
8264 finished_group(&mut producer, 1_000);
8265 producer.append_datagram(Timestamp::ZERO, &b"beat"[..]).unwrap();
8266 }
8267
8268 let consumer = producer.consume();
8269 assert!(consumer.peek_group(0).is_none(), "old groups still evict");
8270 assert!(
8271 pool.used() < 4 * 1_256,
8272 "interleaved datagrams must not bypass the budget: {}",
8273 pool.used()
8274 );
8275 }
8276
8277 #[tokio::test]
8281 async fn aborted_group_leaves_no_ghost_sample() {
8282 let (producer, pool) = pooled_producer(1 << 40);
8283 let group0 = producer.append_group().unwrap();
8284 producer.append_group().unwrap(); assert!(pool.average().is_some(), "demoted group is sampled");
8287 group0.abort(Error::Cancel).unwrap();
8288 assert_eq!(pool.average(), None, "the abort must remove the sample");
8289 }
8290
8291 #[tokio::test]
8294 async fn empty_groups_repay_overhead() {
8295 let (producer, pool) = pooled_producer(1_000);
8296 for _ in 0..100 {
8297 let group = producer.append_group().unwrap();
8298 group.finish().unwrap();
8299 }
8300
8301 assert!(
8302 pool.used() <= 3_000,
8303 "empty-group overhead must stay near the budget: {}",
8304 pool.used()
8305 );
8306 }
8307
8308 #[tokio::test]
8311 async fn growth_on_demoted_group_is_billed() {
8312 let (producer, pool) = pooled_producer(2_000);
8313 let mut straggler = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); straggler
8318 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 10_000]))
8319 .unwrap();
8320
8321 producer.append_group().unwrap().finish().unwrap(); let consumer = producer.consume();
8325 assert!(consumer.peek_group(0).is_none(), "the ballooned group is evicted");
8326 assert!(pool.used() <= 3_000, "growth is reclaimed: {}", pool.used());
8327 }
8328
8329 #[tokio::test]
8332 async fn refilled_sequence_stays_out_of_subscriptions() {
8333 let (producer, _pool) = pooled_producer(1 << 40);
8334 let dynamic = producer.dynamic();
8335 let consumer = producer.consume();
8336
8337 producer.create_group(0u64.into()).unwrap().finish().unwrap();
8338 let aborted = producer.create_group(1u64.into()).unwrap();
8339 aborted.abort(Error::Cancel).unwrap();
8340 producer.create_group(2u64.into()).unwrap().finish().unwrap();
8341
8342 let pending = consumer.fetch_group(1, None);
8345 let req = dynamic
8346 .requested_group()
8347 .now_or_never()
8348 .expect("should not block")
8349 .unwrap();
8350 let mut group = req.accept(None).unwrap();
8351 group
8352 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"backfill"))
8353 .unwrap();
8354 group.finish().unwrap();
8355 pending.await.unwrap();
8356
8357 assert!(consumer.peek_group(1).is_some());
8359 let mut subscriber = producer.subscribe(replay());
8360 assert_eq!(subscriber.assert_group().sequence, 0);
8361 assert_eq!(subscriber.assert_group().sequence, 2);
8362 subscriber.assert_no_group();
8363 }
8364
8365 #[tokio::test]
8368 async fn expired_backfill_behind_refreshed_reclaimed() {
8369 let (producer, _pool) = pooled_producer(1 << 40);
8370 let dynamic = producer.dynamic();
8371 let consumer = producer.consume();
8372
8373 producer.create_group(5u64.into()).unwrap().finish().unwrap();
8374 for sequence in [2u64, 3u64] {
8375 let pending = consumer.fetch_group(sequence, None);
8376 let req = dynamic
8377 .requested_group()
8378 .now_or_never()
8379 .expect("should not block")
8380 .unwrap();
8381 let mut group = req.accept(None).unwrap();
8382 group
8383 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 100]))
8384 .unwrap();
8385 group.finish().unwrap();
8386 pending.await.unwrap();
8387 }
8388
8389 crate::model::clock::advance(cache::DEFAULT_EXPIRY / 2 + Duration::from_secs(1));
8391 consumer.fetch_group(2, None).await.unwrap();
8392 crate::model::clock::advance(cache::DEFAULT_EXPIRY / 2 + Duration::from_secs(1));
8393 producer.create_group(6u64.into()).unwrap().finish().unwrap();
8394
8395 let consumer = producer.consume();
8396 assert!(consumer.peek_group(2).is_some(), "refreshed backfill survives");
8397 assert!(consumer.peek_group(3).is_none(), "expired backfill is reclaimed");
8398 }
8399
8400 #[tokio::test]
8403 async fn same_tick_fetch_protects() {
8404 let (mut producer, _pool) = pooled_producer(10_000);
8406 let consumer = producer.consume();
8407
8408 finished_group(&mut producer, 3_000); finished_group(&mut producer, 3_000); finished_group(&mut producer, 3_000); consumer.fetch_group(0, None).await.unwrap();
8413
8414 finished_group(&mut producer, 3_000); finished_group(&mut producer, 3_000); assert!(consumer.peek_group(0).is_some(), "same-tick refresh protects");
8418 assert!(consumer.peek_group(1).is_none(), "the unread group dies instead");
8419 }
8420
8421 #[tokio::test]
8425 async fn refetched_latest_stays_protected() {
8426 let (producer, _pool) = pooled_producer(10_000);
8427 let dynamic = producer.dynamic();
8428 let consumer = producer.consume();
8429
8430 let straggler = producer.append_group().unwrap(); let latest = producer.append_group().unwrap(); latest.abort(Error::Cancel).unwrap();
8435
8436 let pending = consumer.fetch_group(1, None);
8438 let req = dynamic
8439 .requested_group()
8440 .now_or_never()
8441 .expect("should not block")
8442 .unwrap();
8443 let mut group = req.accept(None).unwrap();
8444 group
8445 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 1000]))
8446 .unwrap();
8447 group.finish().unwrap();
8448 pending.await.unwrap();
8449
8450 {
8453 let state = producer.state.read();
8454 assert!(state.lookup.contains_key(&1), "refetched group is cached");
8455 assert!(
8456 state.evict.iter().all(|(sequence, _)| *sequence != 1),
8457 "the live edge must not be an eviction candidate"
8458 );
8459 }
8460 drop(straggler);
8461 }
8462
8463 #[tokio::test]
8466 async fn eviction_allows_refetch() {
8467 let (mut producer, _pool) = pooled_producer(10_000);
8468 let dynamic = producer.dynamic();
8469
8470 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let consumer = producer.consume();
8475 assert!(consumer.peek_group(0).is_none());
8476 let pending = consumer.fetch_group(0, None);
8477
8478 let req = dynamic
8479 .requested_group()
8480 .now_or_never()
8481 .expect("should not block")
8482 .unwrap();
8483 assert_eq!(req.sequence(), 0);
8484
8485 let mut group = req.accept(None).unwrap();
8486 group
8487 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"refetched"))
8488 .unwrap();
8489 group.finish().unwrap();
8490
8491 let mut group = pending.await.unwrap();
8492 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"refetched");
8493 }
8494
8495 #[test]
8505 fn an_aborted_group_releases_its_sequence() {
8506 let producer = track_producer("test", None);
8507 let consumer = producer.consume();
8508
8509 let mut group = producer.create_group(group::Info { sequence: 3 }).unwrap();
8510 group
8511 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"head"))
8512 .unwrap();
8513
8514 assert!(matches!(
8516 producer.create_group(group::Info { sequence: 3 }),
8517 Err(Error::Duplicate)
8518 ));
8519 assert!(consumer.peek_group(3).is_some());
8520
8521 group.abort(Error::Cancel).unwrap();
8522
8523 assert!(consumer.peek_group(3).is_none(), "an aborted slot is a cache miss");
8524 producer
8525 .create_group(group::Info { sequence: 3 })
8526 .expect("an aborted slot releases its sequence")
8527 .finish()
8528 .unwrap();
8529 }
8530
8531 #[tokio::test]
8534 async fn fetched_backfill_not_subscribed() {
8535 let (producer, _pool) = pooled_producer(1 << 40);
8536 let dynamic = producer.dynamic();
8537 let consumer = producer.consume();
8538
8539 producer.create_group(5u64.into()).unwrap().finish().unwrap();
8541 producer.create_group(6u64.into()).unwrap().finish().unwrap();
8542
8543 let pending = consumer.fetch_group(2, None);
8545 let req = dynamic
8546 .requested_group()
8547 .now_or_never()
8548 .expect("should not block")
8549 .unwrap();
8550 let mut group = req.accept(None).unwrap();
8551 group
8552 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"backfill"))
8553 .unwrap();
8554 group.finish().unwrap();
8555 let mut fetched = pending.await.unwrap();
8556 assert_eq!(&fetched.read_frame().await.unwrap().unwrap().payload[..], b"backfill");
8557 assert!(consumer.peek_group(2).is_some(), "backfill is cached for later fetches");
8558
8559 let mut subscriber = producer.subscribe(replay());
8561 assert_eq!(subscriber.assert_group().sequence, 5);
8562 assert_eq!(subscriber.assert_group().sequence, 6);
8563 subscriber.assert_no_group();
8564 }
8565
8566 #[tokio::test]
8569 async fn expired_backfill_reclaimed() {
8570 let (producer, pool) = pooled_producer(1 << 40);
8571 let dynamic = producer.dynamic();
8572 let consumer = producer.consume();
8573
8574 producer.create_group(5u64.into()).unwrap().finish().unwrap();
8575
8576 let pending = consumer.fetch_group(2, None);
8578 let req = dynamic
8579 .requested_group()
8580 .now_or_never()
8581 .expect("should not block")
8582 .unwrap();
8583 let mut group = req.accept(None).unwrap();
8584 group
8585 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 1000]))
8586 .unwrap();
8587 group.finish().unwrap();
8588 pending.await.unwrap();
8589 let used = pool.used();
8590
8591 crate::model::clock::advance(cache::DEFAULT_EXPIRY + Duration::from_secs(1));
8593 producer.create_group(6u64.into()).unwrap().finish().unwrap();
8594
8595 assert!(consumer.peek_group(2).is_none(), "expired backfill is reclaimed");
8596 assert!(pool.used() < used, "its bytes are released");
8597 }
8598
8599 #[tokio::test]
8600 async fn fetch_aborts_with_track() {
8601 let producer = track_producer("test", None);
8602 let dynamic = producer.dynamic();
8603 let consumer = producer.consume();
8604
8605 let pending = consumer.fetch_group(3, None);
8606 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
8607
8608 producer.abort(Error::Cancel).unwrap();
8609 assert!(pending.await.is_err());
8610 drop(dynamic);
8611 }
8612}