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
153 broadcast: Arc<broadcast::Info>,
156
157 cache: Arc<cache::Track>,
161
162 lookup: BTreeMap<u64, Slot>,
168
169 arrival: VecDeque<(u64, u32)>,
174
175 evict: VecDeque<(u64, u32)>,
183
184 debt: u64,
189
190 datagrams: VecDeque<Datagram>,
193
194 datagram_offset: usize,
197
198 offset: usize,
201
202 max_sequence: Option<u64>,
205
206 latest_group: Option<u64>,
211
212 next_stamp: u32,
214
215 final_sequence: Option<u64>,
217
218 start_sequence: Option<u64>,
222
223 resume: Option<Position>,
227
228 abort: Option<Error>,
230
231 subscriptions: kio::Shared<Subscriptions>,
235
236 fetch: kio::Shared<FetchState>,
239}
240
241struct Slot {
247 group: group::Producer,
248
249 stamp: u32,
254
255 visible: bool,
258}
259
260pub(crate) const CACHE_OVERHEAD: u64 = 2 * (size_of::<u64>() + size_of::<Slot>() + 2 * size_of::<(u64, u32)>()) as u64;
269
270type Subscriptions = Vec<kio::Consumer<Subscription>>;
272
273pub(crate) type FetchState = Requests<u64, PendingFetch>;
278
279pub(crate) struct PendingFetch {
281 priority: u8,
283
284 frame_start: u64,
288
289 result: kio::Producer<FetchOutcome>,
294}
295
296#[derive(Default)]
299pub(crate) struct FetchOutcome {
300 pub(crate) rejected: Option<Error>,
301}
302
303impl TrackState {
304 fn normalize_info(broadcast: &broadcast::Info, mut info: Info) -> Info {
305 info.max_age = info.max_age.min(broadcast.cache_duration);
306 info
307 }
308
309 fn accept(&mut self, info: Info) {
310 self.published = true;
311 self.install(info);
312 }
313
314 fn poll_info(&self) -> Poll<Result<Info>> {
315 if let Some(info) = &self.info {
316 Poll::Ready(Ok(info.clone()))
317 } else if let Some(err) = &self.abort {
318 Poll::Ready(Err(err.clone()))
321 } else {
322 Poll::Pending
323 }
324 }
325
326 fn poll_recv_group(&self, index: usize, min_sequence: u64) -> Poll<Result<Option<(group::Producer, usize)>>> {
332 let start = index.saturating_sub(self.offset);
333 for (i, (sequence, stamp)) in self.arrival.iter().enumerate().skip(start) {
334 if *sequence >= min_sequence
335 && let Some(slot) = self.lookup.get(sequence)
336 && slot.stamp == *stamp
337 && !slot.group.is_aborted()
338 {
339 return Poll::Ready(Ok(Some((slot.group.clone(), self.offset + i))));
340 }
341 }
342
343 if self.is_complete() {
345 Poll::Ready(Ok(None))
346 } else if let Some(err) = &self.abort {
347 Poll::Ready(Err(err.clone()))
348 } else {
349 Poll::Pending
350 }
351 }
352
353 fn poll_recv_datagram(&self, index: usize) -> Poll<Result<Option<(Datagram, usize)>>> {
359 let start = index.saturating_sub(self.datagram_offset);
360 if let Some(datagram) = self.datagrams.get(start) {
361 return Poll::Ready(Ok(Some((datagram.clone(), self.datagram_offset + start))));
362 }
363
364 if self.is_complete() {
366 Poll::Ready(Ok(None))
367 } else if let Some(err) = &self.abort {
368 Poll::Ready(Err(err.clone()))
369 } else {
370 Poll::Pending
371 }
372 }
373
374 fn push_datagram(&mut self, datagram: Datagram) {
376 if self.datagrams.len() == MAX_DATAGRAMS {
377 self.datagrams.pop_front();
378 self.datagram_offset += 1;
379 }
380 self.datagrams.push_back(datagram);
381 }
382
383 fn poll_next_in_range(
393 &self,
394 next_sequence: u64,
395 end_sequence: Option<u64>,
396 ) -> Poll<Result<Option<group::Producer>>> {
397 if let Some(end) = end_sequence
402 && end <= next_sequence
403 {
404 if let Some(err) = &self.abort {
405 return Poll::Ready(Err(err.clone()));
406 }
407 return Poll::Pending;
408 }
409
410 let best = self
411 .lookup
412 .range(next_sequence..)
413 .map(|(_, slot)| &slot.group)
414 .take_while(|group| super::subscription::before_end(group.sequence, end_sequence))
415 .find(|group| !group.is_aborted());
416
417 if let Some(group) = best {
418 return Poll::Ready(Ok(Some(group.clone())));
423 }
424
425 if let Some(err) = &self.abort {
427 return Poll::Ready(Err(err.clone()));
428 }
429 if let Some(fin) = self.final_sequence
432 && next_sequence >= fin
433 {
434 return Poll::Ready(Ok(None));
435 }
436 Poll::Pending
437 }
438
439 fn covering_group(&self, sequence: u64, frame_start: u64) -> Option<&group::Producer> {
446 let slot = self.lookup.get(&sequence)?;
447 let first = slot.group.live_first_frame()?;
448 (first as u64 <= frame_start).then_some(&slot.group)
449 }
450
451 fn max_age_bound(&self) -> Option<Duration> {
454 self.info.as_ref().map(|info| info.max_age)
455 }
456
457 fn live_edge(&self, cap: Option<u64>) -> Option<Edge> {
472 let presentation = self
473 .lookup
474 .range(..)
475 .rev()
476 .filter(|(seq, _)| super::subscription::before_end(**seq, cap))
477 .find_map(|(_, slot)| {
478 if !slot.visible || slot.group.is_aborted() {
479 return None;
480 }
481 let timestamp = slot.group.timestamp()?;
484 Some(PresentationEdge {
485 sequence: slot.group.sequence,
486 stamp: slot.stamp,
487 timestamp: slot.group.latest().unwrap_or(timestamp),
488 })
489 });
490
491 presentation.map(|presentation| Edge { presentation, cap })
492 }
493
494 fn reach(&self, sequence: u64, cap: Option<u64>) -> Option<Timestamp> {
505 let successor = self
506 .lookup
507 .range(sequence.saturating_add(1)..)
508 .map(|(_, slot)| slot)
509 .take_while(|slot| super::subscription::before_end(slot.group.sequence, cap))
510 .find(|slot| slot.visible && !slot.group.is_aborted())?;
511 successor.group.timestamp()
512 }
513
514 fn is_stale(&self, sequence: u64, edge: Option<&Edge>, budget: Duration) -> bool {
539 let Some(edge) = edge else {
540 return false;
541 };
542 if !self.lookup.contains_key(&sequence) {
543 return false;
544 }
545
546 let live_edge = &edge.presentation;
550 let reach = self.reach(sequence, edge.cap);
551 live_edge.sequence > sequence
552 && self
553 .lookup
554 .get(&live_edge.sequence)
555 .is_some_and(|live| live.stamp == live_edge.stamp && !live.group.is_aborted())
556 && reach.is_some_and(
557 |reach| matches!(live_edge.timestamp.checked_sub(reach), Ok(age) if Duration::from(age) >= budget),
558 )
559 }
560
561 fn poll_fetch_cached(&self, sequence: u64, frame_start: u64) -> Poll<Result<group::Consumer>> {
566 if let Some(group) = self.covering_group(sequence, frame_start) {
567 group.cache_refresh();
571 return Poll::Ready(Ok(group.consume()));
572 }
573
574 if let Some(err) = &self.abort {
575 return Poll::Ready(Err(err.clone()));
576 }
577
578 if self.final_sequence.is_some_and(|fin| sequence >= fin) {
580 return Poll::Ready(Err(Error::NotFound));
581 }
582
583 Poll::Pending
584 }
585
586 pub(super) fn evict_expired(&mut self) {
601 let scan = self.expiry_scan();
602 self.evict_expired_scan(scan);
603 }
604
605 pub(super) fn expiry_scan(&self) -> ExpiryScan {
607 ExpiryScan {
608 start: self.cache.next_expiry_scan(EVICT_SCAN),
609 width: EVICT_SCAN,
610 now: self.cache.pool().now(),
611 max_ticks: self.cache.pool().expiry_ticks(),
612 gc: false,
613 }
614 }
615
616 pub(super) fn expiry_scan_drain(&self) -> ExpiryScan {
619 ExpiryScan {
620 start: 0,
621 width: self.evict.len(),
622 now: self.cache.pool().now(),
623 max_ticks: self.cache.pool().expiry_ticks(),
624 gc: true,
625 }
626 }
627
628 #[cfg(test)]
629 pub(super) fn date_cache_accesses(&self, now: u64) {
630 for slot in self.lookup.values() {
631 slot.group.cache_accessed_tick(Some(now));
632 }
633 }
634
635 pub(super) fn expiry_mutation_due(&self, scan: ExpiryScan) -> bool {
640 let len = self.evict.len();
641 if len > 0 {
642 let start = scan.start % len;
643 let mut retained = 0;
644 for step in 0..len.min(scan.width) {
645 let (sequence, stamp) = self.evict[(start + step) % len];
646 let Some(slot) = self.lookup.get(&sequence) else {
647 continue;
648 };
649 if slot.stamp != stamp {
650 continue;
651 }
652 if slot.group.is_aborted()
653 || (Some(sequence) != self.latest_group
654 && slot
655 .group
656 .cache_accessed_tick(scan.gc.then_some(scan.now))
657 .is_some_and(|tick| scan.now.saturating_sub(tick) > scan.max_ticks))
658 {
659 return true;
660 }
661 retained += 1;
662 if !scan.gc && retained >= EVICT_SCAN {
663 break;
664 }
665 }
666 }
667
668 self.arrival
669 .front()
670 .is_some_and(|(sequence, stamp)| !self.is_current(*sequence, *stamp))
671 || self
672 .evict
673 .front()
674 .is_some_and(|(sequence, stamp)| !self.is_current(*sequence, *stamp))
675 || self.evict.len() > 2 * self.lookup.len() + EVICT_SLACK
676 }
677
678 pub(super) fn evict_expired_scan(&mut self, scan: ExpiryScan) {
681 let len = self.evict.len();
682 if len > 0 {
683 let start = scan.start % len;
684 let mut retained = 0;
685 for step in 0..len.min(scan.width) {
686 let (sequence, stamp) = self.evict[(start + step) % len];
687 let Some(slot) = self.lookup.get(&sequence) else {
688 continue;
689 };
690 if slot.stamp != stamp {
691 continue;
693 }
694 if slot.group.is_aborted() {
697 self.lookup.remove(&sequence);
698 continue;
699 }
700 if Some(sequence) == self.latest_group
701 || slot
702 .group
703 .cache_accessed_tick(scan.gc.then_some(scan.now))
704 .is_none_or(|tick| scan.now.saturating_sub(tick) <= scan.max_ticks)
705 {
706 retained += 1;
709 if !scan.gc && retained >= EVICT_SCAN {
710 break;
711 }
712 continue;
713 }
714 let slot = self.lookup.remove(&sequence).unwrap();
718 let _ = slot.group.abort(Error::Old);
719 }
720 }
721
722 while let Some((sequence, stamp)) = self.arrival.front() {
725 if self.lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp) {
726 break;
727 }
728 self.arrival.pop_front();
729 self.offset += 1;
730 }
731
732 while let Some((sequence, stamp)) = self.evict.front() {
734 if self.lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp) {
735 break;
736 }
737 self.evict.pop_front();
738 }
739
740 if self.evict.len() > 2 * self.lookup.len() + EVICT_SLACK {
743 let lookup = &self.lookup;
744 self.evict
745 .retain(|(sequence, stamp)| lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp));
746 }
747 }
748
749 fn is_current(&self, sequence: u64, stamp: u32) -> bool {
751 self.lookup.get(&sequence).is_some_and(|slot| slot.stamp == stamp)
752 }
753
754 fn clear_cache(&mut self) {
757 self.lookup.clear();
758 self.arrival.clear();
759 self.evict.clear();
760 self.latest_group = None;
761 self.debt = 0;
762 }
763
764 fn install(&mut self, info: Info) {
770 let info = Self::normalize_info(&self.broadcast, info);
771 self.info = Some(info);
772 }
773
774 fn spawn(broadcast: Arc<broadcast::Info>) -> kio::Producer<Self> {
781 let state = kio::Producer::new(Self {
782 broadcast: broadcast.clone(),
783 ..Default::default()
784 });
785 let cache = cache::Track::new(broadcast.pool.clone(), state.downgrade());
786 state.write().ok().expect("a new track is open").cache = cache;
787 state
788 }
789
790 fn claim_sequence(&mut self, sequence: u64, frame_start: u64) -> Result<()> {
799 if let Some(slot) = self.lookup.get(&sequence) {
800 if slot
804 .group
805 .live_first_frame()
806 .is_some_and(|first| first as u64 <= frame_start)
807 {
808 return Err(Error::Duplicate);
809 }
810 self.lookup.remove(&sequence);
811 }
812 Ok(())
813 }
814
815 fn insert_group(&mut self, group: &group::Producer, visible: bool) {
822 let sequence = group.sequence;
823 self.next_stamp = self.next_stamp.wrapping_add(1);
824 let stamp = self.next_stamp;
825
826 if self.latest_group.is_none_or(|latest| sequence >= latest) {
830 if let Some(latest) = self.latest_group
833 && sequence > latest
834 && let Some(prev) = self.lookup.get(&latest)
835 {
836 prev.group.cache_demote();
837 self.evict.push_back((latest, prev.stamp));
838 }
839 self.latest_group = Some(sequence);
840 } else {
841 group.cache_demote();
842 self.evict.push_back((sequence, stamp));
843 }
844
845 self.max_sequence = Some(self.max_sequence.map_or(sequence, |max| max.max(sequence)));
846 self.lookup.insert(
847 sequence,
848 Slot {
849 group: group.clone(),
850 stamp,
851 visible,
852 },
853 );
854 if visible {
855 self.arrival.push_back((sequence, stamp));
856 }
857 }
858
859 fn commit_group(&mut self, group: &group::Producer, visible: bool) {
863 self.charge_debt();
864 self.insert_group(group, visible);
865 self.evict_expired();
866 }
867
868 pub(super) fn charge_debt(&mut self) {
881 let written = self.cache.take_written();
882 let pool = self.cache.pool().clone();
883 match pool.accrue(written) {
884 Some(mut accrued) => {
885 if self.oldest_is_stale(&pool) {
886 accrued = accrued.saturating_mul(2);
887 }
888 self.debt = self.debt.saturating_add(accrued).min(pool.used());
891 self.pay_debt(&pool, written.saturating_mul(2));
894 }
895 None => self.debt = 0,
898 }
899 }
900
901 fn oldest_is_stale(&self, pool: &cache::Pool) -> bool {
905 let Some(average) = pool.average() else {
906 return false;
907 };
908 let Some((sequence, stamp)) = self.evict.front() else {
909 return false;
910 };
911 let Some(slot) = self.lookup.get(sequence) else {
912 return false;
913 };
914 slot.stamp == *stamp && !slot.group.is_aborted() && slot.group.cache_accessed() <= average
915 }
916
917 fn pay_debt(&mut self, pool: &cache::Pool, cap: u64) {
930 let average = pool.average().unwrap_or(0);
931 let mut paid = 0u64;
932 let mut scanned = 0usize;
933 for _ in 0..self.evict.len() {
934 if self.debt == 0 || paid >= cap || scanned >= EVICT_SCAN {
935 return;
936 }
937 let Some((sequence, stamp)) = self.evict.pop_front() else {
938 return;
939 };
940 let Some(slot) = self.lookup.get(&sequence) else {
941 continue;
943 };
944 if slot.stamp != stamp {
945 continue;
947 }
948 if slot.group.is_aborted() {
949 self.lookup.remove(&sequence);
951 continue;
952 }
953 if Some(sequence) == self.latest_group {
954 self.evict.push_back((sequence, stamp));
956 continue;
957 }
958
959 scanned += 1;
960 if slot.group.cache_accessed() > average {
964 self.evict.push_back((sequence, stamp));
965 continue;
966 }
967 let size = slot.group.cache_size();
970 if size > self.debt {
971 self.evict.push_front((sequence, stamp));
972 return;
973 }
974
975 self.debt -= size;
976 paid = paid.saturating_add(size);
977 let slot = self.lookup.remove(&sequence).unwrap();
978 let _ = slot.group.abort(Error::Evicted);
979 }
980 }
981
982 fn set_start(&mut self, start_sequence: Option<u64>) {
987 self.start_sequence = start_sequence;
988 }
989
990 fn set_final(&mut self, final_sequence: u64) -> Result<()> {
993 if self.final_sequence.is_some() {
994 return Err(Error::Closed);
995 }
996 if let Some(max) = self.max_sequence
997 && final_sequence <= max
998 {
999 return Err(Error::ProtocolViolation);
1000 }
1001 self.final_sequence = Some(final_sequence);
1002 Ok(())
1003 }
1004
1005 fn is_complete(&self) -> bool {
1011 self.final_sequence
1012 .is_some_and(|fin| self.max_sequence.map_or(0, |max| max.saturating_add(1)) >= fin)
1013 }
1014
1015 fn resume_position(&self) -> Option<Position> {
1021 if self.resume.is_some() {
1024 return self.resume;
1025 }
1026
1027 let max = self.latest_group?;
1028 match self.lookup.get(&max).and_then(|slot| slot.group.resume_frame()) {
1029 Some(frame) => Some(Position {
1032 group: max,
1033 frame: frame as u64,
1034 }),
1035 None => Some(Position::group(max.saturating_add(1))),
1040 }
1041 }
1042
1043 fn poll_finished(&self) -> Poll<Result<u64>> {
1044 if let Some(fin) = self.final_sequence {
1045 Poll::Ready(Ok(fin))
1046 } else if let Some(err) = &self.abort {
1047 Poll::Ready(Err(err.clone()))
1048 } else {
1049 Poll::Pending
1050 }
1051 }
1052
1053 fn modify(producer: &kio::Producer<Self>) -> Result<kio::Mut<'_, Self>> {
1054 producer.write().map_err(|r| r.abort.clone().unwrap_or(Error::Dropped))
1055 }
1056
1057 pub(crate) fn insert_group_request(
1063 &mut self,
1064 sequence: u64,
1065 frame_start: u64,
1066 info: Option<Info>,
1067 ) -> Result<group::Producer> {
1068 if let Some(err) = &self.abort {
1069 return Err(err.clone());
1070 }
1071 if let Some(fin) = self.final_sequence
1072 && sequence >= fin
1073 {
1074 return Err(Error::Closed);
1075 }
1076
1077 if self.info.is_none() {
1081 self.install(info.unwrap_or_default());
1082 }
1083 let info = self.info.clone().unwrap();
1084
1085 self.claim_sequence(sequence, frame_start)?;
1087
1088 let group = group::Producer::new(group::Info { sequence }, info, self.cache.clone());
1089 group.cache_refresh();
1094 self.commit_group(&group, false);
1095 Ok(group)
1096 }
1097}
1098
1099fn commit_abort(mut state: kio::Mut<'_, TrackState>, err: Error) {
1102 state.resume = state.resume_position();
1105 state.abort = Some(err);
1106 state.clear_cache();
1107 state.datagrams.clear();
1108 state.close();
1109}
1110
1111#[derive(Clone)]
1113pub struct Producer {
1114 name: Arc<str>,
1115 info: Info,
1116 broadcast: Arc<broadcast::Info>,
1119 state: kio::Producer<TrackState>,
1120 prev_subscription: Option<Subscription>,
1121 alive: Arc<Alive>,
1123 stats: stats::Scope,
1127}
1128
1129impl Producer {
1130 pub(crate) fn new(
1138 broadcast: Arc<broadcast::Info>,
1139 name: impl Into<Arc<str>>,
1140 info: impl Into<Option<Info>>,
1141 ) -> Self {
1142 let name = name.into();
1143 let info = TrackState::normalize_info(&broadcast, info.into().unwrap_or_default());
1144 let state = TrackState::spawn(broadcast.clone());
1145 state.write().ok().expect("a new track is open").accept(info.clone());
1146 let alive = Alive::new(name.clone(), state.clone());
1147 alive.publish(None);
1148 Self {
1149 name,
1150 info,
1151 state,
1152 broadcast,
1153 prev_subscription: None,
1154 alive,
1155 stats: stats::Scope::default(),
1156 }
1157 }
1158
1159 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
1163 self.alive.publish(Some(&scope));
1164 self.stats = scope;
1165 self
1166 }
1167
1168 pub fn name(&self) -> &str {
1170 &self.name
1171 }
1172
1173 pub fn broadcast(&self) -> &broadcast::Info {
1175 &self.broadcast
1176 }
1177
1178 pub(crate) fn adopt_group(&mut self, group: group::Producer, visible: bool) -> Result<()> {
1183 let mut state = self.modify()?;
1184 if let Some(fin) = state.final_sequence
1185 && group.sequence >= fin
1186 {
1187 return Err(Error::Closed);
1188 }
1189 if state.lookup.contains_key(&group.sequence) {
1190 return Err(Error::Duplicate);
1191 }
1192 state.insert_group(&group, visible);
1193 Ok(())
1194 }
1195
1196 pub fn create_group(&self, group: group::Info) -> Result<group::Producer> {
1198 let mut state = self.modify()?;
1199 if let Some(fin) = state.final_sequence
1200 && group.sequence >= fin
1201 {
1202 return Err(Error::Closed);
1203 }
1204 let track = state.info.clone().unwrap();
1205
1206 state.claim_sequence(group.sequence, 0)?;
1208
1209 let group = group::Producer::new(group, track, state.cache.clone()).with_meter(self.stats.meter());
1210 state.commit_group(&group, true);
1211
1212 Ok(group)
1213 }
1214
1215 pub fn append_group(&self) -> Result<group::Producer> {
1217 let mut state = self.modify()?;
1218 let sequence = match state.max_sequence {
1219 Some(s) => s.checked_add(1).ok_or(coding::BoundsExceeded)?,
1220 None => 0,
1221 };
1222 if let Some(fin) = state.final_sequence
1223 && sequence >= fin
1224 {
1225 return Err(Error::Closed);
1226 }
1227
1228 let track = state.info.clone().unwrap();
1229
1230 let group =
1231 group::Producer::new(group::Info { sequence }, track, state.cache.clone()).with_meter(self.stats.meter());
1232 state.commit_group(&group, true);
1233
1234 Ok(group)
1235 }
1236
1237 pub fn append_datagram<B: crate::IntoBytes>(&mut self, timestamp: Timestamp, payload: B) -> Result<u64> {
1249 let payload = payload.into_bytes();
1250 if payload.len() > super::datagram::MAX_DATAGRAM_PAYLOAD {
1251 return Err(Error::FrameTooLarge);
1252 }
1253 let meter = self.stats.meter();
1255 let mut state = self.modify()?;
1256 let timescale = state.info.as_ref().unwrap().timescale;
1258 let timestamp = timestamp.convert(timescale).map_err(|_| Error::TimestampMismatch)?;
1259 let sequence = match state.max_sequence {
1260 Some(s) => s.checked_add(1).ok_or(coding::BoundsExceeded)?,
1261 None => 0,
1262 };
1263 if let Some(fin) = state.final_sequence
1264 && sequence >= fin
1265 {
1266 return Err(Error::Closed);
1267 }
1268 state.max_sequence = Some(sequence);
1269 meter.datagram(payload.len() as u64);
1270 state.push_datagram(Datagram {
1271 sequence,
1272 timestamp,
1273 payload,
1274 });
1275 let cache = state.cache.clone();
1276 drop(state);
1277 cache.settle(None);
1278 Ok(sequence)
1279 }
1280
1281 pub fn insert_datagram<B: crate::IntoBytes>(
1287 &mut self,
1288 sequence: u64,
1289 timestamp: Timestamp,
1290 payload: B,
1291 ) -> Result<()> {
1292 let payload = payload.into_bytes();
1293 if payload.len() > super::datagram::MAX_DATAGRAM_PAYLOAD {
1294 return Err(Error::FrameTooLarge);
1295 }
1296 let meter = self.stats.meter();
1298 let mut state = self.modify()?;
1299 let timescale = state.info.as_ref().unwrap().timescale;
1301 let timestamp = timestamp.convert(timescale).map_err(|_| Error::TimestampMismatch)?;
1302 if let Some(fin) = state.final_sequence
1303 && sequence >= fin
1304 {
1305 return Err(Error::Closed);
1306 }
1307 state.max_sequence = Some(state.max_sequence.unwrap_or(0).max(sequence));
1308 meter.datagram(payload.len() as u64);
1309 state.push_datagram(Datagram {
1310 sequence,
1311 timestamp,
1312 payload,
1313 });
1314 let cache = state.cache.clone();
1315 drop(state);
1316 cache.settle(None);
1317 Ok(())
1318 }
1319
1320 pub fn write_frame<B: crate::IntoBytes>(&mut self, timestamp: Timestamp, frame: B) -> Result<()> {
1325 let frame = crate::IntoBytes::into_bytes(frame);
1326 if frame.len() as u64 > group::MAX_CACHE_BYTES {
1327 return Err(Error::FrameTooLarge);
1328 }
1329 let mut group = self.append_group()?;
1330 group.write_frame(timestamp, frame)?;
1331 group.finish()?;
1332 Ok(())
1333 }
1334
1335 pub fn finish(&self) -> Result<()> {
1341 let mut state = self.modify()?;
1342 let final_sequence = match state.max_sequence {
1343 Some(max) => max.checked_add(1).ok_or(coding::BoundsExceeded)?,
1344 None => 0,
1345 };
1346 state.set_final(final_sequence)
1347 }
1348
1349 pub fn finish_at(&mut self, final_sequence: u64) -> Result<()> {
1362 self.modify()?.set_final(final_sequence)
1363 }
1364
1365 pub fn start_at(&mut self, sequence: impl Into<Option<u64>>) -> Result<()> {
1376 self.modify()?.set_start(sequence.into());
1377 Ok(())
1378 }
1379
1380 #[cfg(test)]
1383 pub(crate) fn start_sequence(&self) -> Option<u64> {
1384 self.state.read().start_sequence
1385 }
1386
1387 pub fn final_sequence(&self) -> Option<u64> {
1392 self.state.read().final_sequence
1393 }
1394
1395 pub fn abort(self, err: Error) -> Result<()> {
1406 commit_abort(self.modify()?, err);
1407 Ok(())
1408 }
1409
1410 #[expect(
1418 clippy::result_large_err,
1419 reason = "return the owned producer without allocating on an idle check"
1420 )]
1421 pub fn abort_unused(self, err: Error) -> std::result::Result<(), Self> {
1422 match self.state.write_unused() {
1423 kio::Unused::Idle(guard) => {
1424 commit_abort(guard, err);
1425 return Ok(());
1426 }
1427 kio::Unused::Closed => return Ok(()),
1428 kio::Unused::Used => {}
1429 }
1430 Err(self)
1431 }
1432
1433 pub fn is_used(&self) -> bool {
1440 !self.is_closed() && self.state.is_used()
1441 }
1442
1443 pub async fn unused(&self) -> Result<()> {
1445 self.state.unused().await.map_err(|_| self.abort_reason())
1446 }
1447
1448 pub async fn used(&self) -> Result<()> {
1450 self.state.used().await.map_err(|_| self.abort_reason())
1451 }
1452
1453 pub async fn closed(&self) -> Error {
1455 kio::wait(|waiter| self.poll_closed(waiter)).await
1456 }
1457
1458 pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<Error> {
1460 self.state.poll_closed(waiter).map(|()| self.abort_reason())
1461 }
1462
1463 fn abort_reason(&self) -> Error {
1465 self.state.read().abort.clone().unwrap_or(Error::Dropped)
1466 }
1467
1468 pub fn is_closed(&self) -> bool {
1470 self.state.read().is_closed()
1471 }
1472
1473 pub fn latest(&self) -> Option<u64> {
1475 self.state.read().max_sequence
1476 }
1477
1478 pub fn is_clone(&self, other: &Self) -> bool {
1480 self.state.same_channel(&other.state)
1481 }
1482
1483 pub(crate) fn weak(&self) -> TrackWeak {
1485 TrackWeak {
1486 name: self.name.clone(),
1487 state: self.state.weak(),
1488 }
1489 }
1490
1491 pub fn demand(&self) -> Demand {
1499 Demand {
1500 name: self.name.clone(),
1501 state: self.state.weak(),
1502 }
1503 }
1504
1505 pub fn consume(&self) -> Consumer {
1510 Consumer::plain(self.name.clone(), self.state.consume())
1511 }
1512
1513 pub fn subscribe(&self, subscription: impl Into<Option<Subscription>>) -> Subscriber {
1523 let preferences = subscription.into().unwrap_or_default();
1524
1525 let info = self.info.clone();
1530 let min_sequence = floor_of(&preferences);
1531 let subscription = kio::Producer::new(preferences);
1532 register_subscription(self.state.read(), &subscription);
1533 let drift_cap = kio::Producer::new(None);
1534
1535 let broadcast = self.state.read().broadcast.clone();
1538 Subscriber {
1539 name: self.name.clone(),
1540 broadcast,
1541 info,
1542 inner: SubscriberKind::Plain(PlainSubscriber {
1543 state: self.state.consume(),
1544 subscription,
1545 min_sequence,
1546 index: 0,
1547 datagram_index: 0,
1548 next_sequence: 0,
1549 end_sequence: None,
1550 parked: BTreeMap::new(),
1551 stale_cap: None,
1552 drift_cap,
1553 stale: stats::Content::default(),
1554 seek_pending: BTreeMap::new(),
1555 }),
1556 stats: stats::Scope::default(),
1558 _stats_sub: stats::Subscription::default(),
1559 }
1560 }
1561
1562 pub async fn subscription_changed(&mut self) -> Result<Option<Subscription>> {
1568 kio::wait(|waiter| self.poll_subscription_changed(waiter)).await
1569 }
1570
1571 pub fn subscription(&self) -> Option<Subscription> {
1579 let state = self.state.read();
1580 let (subs, bound) = (state.subscriptions.clone(), state.max_age_bound());
1581 drop(state);
1582 snapshot_subscription(&subs, bound)
1583 }
1584
1585 pub fn poll_subscription_changed(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Subscription>>> {
1589 if self.state.poll_closed(waiter).is_ready() {
1592 let abort = self.state.read().abort.clone();
1593 return Poll::Ready(Err(abort.unwrap_or(Error::Dropped)));
1594 }
1595
1596 let state = self.state.read();
1598 let (subs, bound) = (state.subscriptions.clone(), state.max_age_bound());
1599 drop(state);
1600
1601 let prev = &self.prev_subscription;
1602 let mut combined = None;
1603 let mut guard = ready!(subs.poll(waiter, |subs| {
1604 let next = combined_subscription(subs, bound, waiter);
1605 if &next == prev {
1606 Poll::Pending
1607 } else {
1608 combined = next;
1609 Poll::Ready(())
1610 }
1611 }));
1612 guard.retain(|sub| !sub.is_closed());
1614 drop(guard);
1615 self.prev_subscription = combined.clone();
1616 Poll::Ready(Ok(combined))
1617 }
1618
1619 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<Result<()>> {
1621 self.state.poll_unused(waiter).map(|used| match used {
1622 Some(()) => Ok(()),
1623 None => Err(self.abort_reason()),
1624 })
1625 }
1626
1627 pub fn dynamic(&self) -> Dynamic {
1631 Dynamic::new(self.name.clone(), self.state.clone(), self.alive.clone())
1632 }
1633
1634 fn modify(&self) -> Result<kio::Mut<'_, TrackState>> {
1635 TrackState::modify(&self.state)
1636 }
1637}
1638
1639fn poll_requested_group(
1643 state: &kio::Producer<TrackState>,
1644 fetch: &kio::Shared<FetchState>,
1645 waiter: &kio::Waiter,
1646) -> Poll<Result<group::Request>> {
1647 if let Poll::Ready(mut guard) = fetch.poll(waiter, |fetch| {
1649 if fetch.has_queued() {
1650 Poll::Ready(())
1651 } else {
1652 Poll::Pending
1653 }
1654 }) {
1655 let sequence = guard.pop().expect("predicate guaranteed a request");
1656 let pending = guard.get(&sequence).expect("popped key must be pending");
1660 let priority = pending.priority;
1661 let frame_start = pending.frame_start;
1662 let result = pending.result.clone();
1663 drop(guard);
1664 return Poll::Ready(Ok(group::Request {
1665 state: state.clone(),
1666 fetch: fetch.clone(),
1667 sequence,
1668 priority,
1669 frame_start,
1670 result,
1671 done: false,
1672 }));
1673 }
1674
1675 match state.poll_ref(waiter, |state| match &state.abort {
1677 Some(err) => Poll::Ready(err.clone()),
1678 None => Poll::Pending,
1679 }) {
1680 Poll::Ready(Ok(err)) => Poll::Ready(Err(err)),
1681 Poll::Ready(Err(closed)) => Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped))),
1682 Poll::Pending => Poll::Pending,
1683 }
1684}
1685
1686pub struct Dynamic {
1696 name: Arc<str>,
1697 state: kio::Producer<TrackState>,
1699 fetch: kio::Shared<FetchState>,
1701 alive: Arc<Alive>,
1704}
1705
1706impl Dynamic {
1707 fn new(name: Arc<str>, state: kio::Producer<TrackState>, alive: Arc<Alive>) -> Self {
1708 let fetch = state.read().fetch.clone();
1709 fetch.lock().add_handler();
1710 Self {
1711 name,
1712 state,
1713 fetch,
1714 alive,
1715 }
1716 }
1717
1718 pub fn name(&self) -> &str {
1720 &self.name
1721 }
1722
1723 pub async fn requested_group(&self) -> Result<group::Request> {
1729 kio::wait(|waiter| self.poll_requested_group(waiter)).await
1730 }
1731
1732 pub fn poll_requested_group(&self, waiter: &kio::Waiter) -> Poll<Result<group::Request>> {
1734 poll_requested_group(&self.state, &self.fetch, waiter)
1735 }
1736
1737 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
1739 self.state.poll_unused(waiter).map(|_| ())
1740 }
1741}
1742
1743impl Clone for Dynamic {
1744 fn clone(&self) -> Self {
1745 self.fetch.lock().add_handler();
1747 Self {
1748 name: self.name.clone(),
1749 state: self.state.clone(),
1750 fetch: self.fetch.clone(),
1751 alive: self.alive.clone(),
1752 }
1753 }
1754}
1755
1756impl Drop for Dynamic {
1757 fn drop(&mut self) {
1758 let mut fetch = self.fetch.lock();
1764 if fetch.remove_handler() {
1765 fetch.drain_queued();
1766 }
1767 }
1768}
1769
1770struct Alive {
1779 name: Arc<str>,
1780 state: kio::Producer<TrackState>,
1781
1782 published: AtomicBool,
1785
1786 stats: OnceLock<stats::Subscription>,
1789}
1790
1791impl Alive {
1792 fn new(name: Arc<str>, state: kio::Producer<TrackState>) -> Arc<Self> {
1793 Arc::new(Self {
1794 name,
1795 state,
1796 published: Default::default(),
1797 stats: Default::default(),
1798 })
1799 }
1800
1801 fn publish(&self, stats: Option<&stats::Scope>) {
1805 self.published.store(true, Ordering::Relaxed);
1806 if let Some(scope) = stats {
1807 let _ = self.stats.set(scope.subscribe());
1810 }
1811 }
1812}
1813
1814impl Drop for Alive {
1815 fn drop(&mut self) {
1816 if !self.published.load(Ordering::Relaxed) {
1818 return;
1819 }
1820 match self.state.write() {
1827 Ok(mut state) => {
1828 if state.final_sequence.is_some() || state.abort.is_some() {
1829 return;
1830 }
1831 tracing::warn!(
1832 track = %self.name,
1833 "track::Producer dropped without finish() or abort()"
1834 );
1835 state.resume = state.resume_position();
1837 state.clear_cache();
1838 state.datagrams.clear();
1839 }
1840 Err(state) => {
1841 if state.final_sequence.is_some() || state.abort.is_some() {
1842 return;
1843 }
1844 tracing::warn!(
1845 track = %self.name,
1846 "track::Producer dropped without finish() or abort()"
1847 );
1848 }
1849 }
1850 }
1851}
1852
1853fn combined_subscription(subs: &Subscriptions, bound: Option<Duration>, waiter: &kio::Waiter) -> Option<Subscription> {
1859 let mut combined = None;
1860 for sub in subs.iter() {
1861 if sub.is_closed() {
1866 continue;
1867 }
1868 let _ = sub.poll_closed(waiter);
1875 let _ = sub.poll(waiter, |_| Poll::<()>::Pending);
1876 if let Poll::Ready(merged) = sub.read().poll_combined(&combined) {
1877 combined = Some(merged);
1878 }
1879 }
1880 clamp_combined(combined, bound)
1881}
1882
1883fn snapshot_subscription(subs: &kio::Shared<Subscriptions>, bound: Option<Duration>) -> Option<Subscription> {
1885 let mut combined: Option<Subscription> = None;
1886 for sub in subs.read().iter() {
1887 if sub.is_closed() {
1889 continue;
1890 }
1891 if let Poll::Ready(merged) = sub.read().poll_combined(&combined) {
1892 combined = Some(merged);
1893 }
1894 }
1895 clamp_combined(combined, bound)
1896}
1897
1898fn servable_cap(cursor: Option<u64>, outer: Option<u64>) -> Option<u64> {
1913 super::subscription::min_some(cursor, outer)
1914}
1915
1916fn floor_of(subscription: &Subscription) -> u64 {
1925 subscription.start.map(|start| start.group).unwrap_or(0)
1926}
1927
1928fn clamp_max_age(mut max_age: Duration, bound: Option<Duration>) -> Duration {
1938 if let Some(bound) = bound {
1939 max_age = max_age.min(bound);
1940 }
1941 max_age
1942}
1943
1944fn clamp_combined(combined: Option<Subscription>, bound: Option<Duration>) -> Option<Subscription> {
1946 let mut combined = combined?;
1947 combined.max_age = clamp_max_age(combined.max_age, bound);
1948 Some(combined)
1949}
1950
1951fn register_subscription(state: kio::Ref<'_, TrackState>, subscription: &kio::Producer<Subscription>) {
1955 if state.is_closed() {
1956 return;
1957 }
1958 let subs = state.subscriptions.clone();
1959 drop(state);
1960 subs.lock().push(subscription.consume());
1961}
1962
1963#[derive(Clone)]
1965pub(crate) struct TrackWeak {
1966 name: Arc<str>,
1967 state: kio::ProducerWeak<TrackState>,
1968}
1969
1970impl TrackWeak {
1971 pub fn try_consume(&self) -> Option<Consumer> {
1978 Some(Consumer::plain(self.name.clone(), self.state.try_consume()?))
1979 }
1980
1981 pub(crate) fn name(&self) -> &Arc<str> {
1984 &self.name
1985 }
1986
1987 pub(crate) fn reject(&self, err: Error) -> bool {
1996 let Some(producer) = self.state.produce() else {
1997 return false;
1998 };
1999 let Ok(mut state) = producer.write() else {
2000 return false;
2001 };
2002 if state.published || state.abort.is_some() {
2003 return false;
2004 }
2005 state.abort = Some(err);
2006 state.close();
2007 true
2008 }
2009
2010 pub(crate) fn is_used(&self) -> bool {
2013 !self.state.is_closed() && self.state.is_used()
2014 }
2015
2016 pub(crate) fn poll_used(&self, waiter: &kio::Waiter) {
2019 let _ = self.state.poll_used(waiter);
2020 }
2021
2022 pub(crate) fn poll_unused(&self, waiter: &kio::Waiter) {
2025 let _ = self.state.poll_unused(waiter);
2026 }
2027}
2028
2029impl super::WeakEntry for TrackWeak {
2030 fn is_closed(&self) -> bool {
2031 self.state.is_closed()
2032 }
2033
2034 fn same_channel(&self, other: &Self) -> bool {
2035 self.state.same_channel(&other.state)
2036 }
2037}
2038
2039#[derive(Clone)]
2048pub struct Demand {
2049 name: Arc<str>,
2050 state: kio::ProducerWeak<TrackState>,
2051}
2052
2053impl Demand {
2054 pub fn name(&self) -> &str {
2056 &self.name
2057 }
2058
2059 pub async fn used(&self) -> Result<()> {
2061 self.state.used().await.map_err(|_| self.abort_reason())
2062 }
2063
2064 pub async fn unused(&self) -> Result<()> {
2066 self.state.unused().await.map_err(|_| self.abort_reason())
2067 }
2068
2069 pub async fn closed(&self) -> Error {
2071 self.state.closed().await;
2072 self.abort_reason()
2073 }
2074
2075 pub(crate) fn priority(&self) -> u8 {
2077 self.state.read().info.as_ref().map_or(0, |info| info.priority)
2079 }
2080
2081 pub fn is_used(&self) -> bool {
2083 self.state.is_used()
2084 }
2085
2086 pub fn poll_used(&self, waiter: &kio::Waiter) -> Poll<Result<()>> {
2088 self.state.poll_used(waiter).map(|used| match used {
2089 Some(()) => Ok(()),
2090 None => Err(self.abort_reason()),
2091 })
2092 }
2093
2094 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<Result<()>> {
2096 self.state.poll_unused(waiter).map(|used| match used {
2097 Some(()) => Ok(()),
2098 None => Err(self.abort_reason()),
2099 })
2100 }
2101
2102 pub(crate) fn is_closed(&self) -> bool {
2104 self.state.is_closed()
2105 }
2106
2107 pub(crate) fn poll_state(&self, waiter: &kio::Waiter) -> DemandState {
2114 loop {
2115 match self.state.poll_used(waiter) {
2116 Poll::Ready(None) => return DemandState::Closed,
2117 Poll::Pending => return DemandState::Idle,
2119 Poll::Ready(Some(())) => match self.state.poll_unused(waiter) {
2121 Poll::Ready(None) => return DemandState::Closed,
2122 Poll::Pending => return DemandState::Active,
2123 Poll::Ready(Some(())) => continue,
2126 },
2127 }
2128 }
2129 }
2130
2131 fn abort_reason(&self) -> Error {
2133 self.state.read().abort.clone().unwrap_or(Error::Dropped)
2134 }
2135}
2136
2137#[derive(Copy, Clone, Debug, Eq, PartialEq)]
2139pub(crate) enum DemandState {
2140 Active,
2142 Idle,
2144 Closed,
2146}
2147
2148#[derive(Clone)]
2159pub struct Consumer {
2160 name: Arc<str>,
2161 broadcast: Arc<broadcast::Info>,
2165 inner: ConsumerKind,
2166 stats: stats::Scope,
2169}
2170
2171#[derive(Clone)]
2172enum ConsumerKind {
2173 Plain(kio::Consumer<TrackState>),
2174 Spliced(super::resume::Consumer),
2175}
2176
2177impl Consumer {
2178 fn plain(name: Arc<str>, state: kio::Consumer<TrackState>) -> Self {
2179 let broadcast = state.read().broadcast.clone();
2180 Self {
2181 name,
2182 broadcast,
2183 inner: ConsumerKind::Plain(state),
2184 stats: stats::Scope::default(),
2185 }
2186 }
2187
2188 pub(crate) fn spliced(name: Arc<str>, broadcast: Arc<broadcast::Info>, resume: super::resume::Consumer) -> Self {
2190 Self {
2191 name,
2192 broadcast,
2193 inner: ConsumerKind::Spliced(resume),
2194 stats: stats::Scope::default(),
2195 }
2196 }
2197
2198 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
2201 self.stats = scope;
2202 self
2203 }
2204
2205 pub(crate) fn with_broadcast(mut self, broadcast: Arc<broadcast::Info>) -> Self {
2209 self.broadcast = broadcast;
2210 self
2211 }
2212
2213 pub(crate) fn cached_groups(&self) -> Vec<(group::Producer, bool)> {
2219 match &self.inner {
2220 ConsumerKind::Plain(state) => {
2221 let state = state.read();
2222 let mut out = Vec::with_capacity(state.lookup.len());
2223 for (sequence, stamp) in state.arrival.iter() {
2224 if let Some(slot) = state.lookup.get(sequence)
2225 && slot.stamp == *stamp
2226 && !slot.group.is_aborted()
2227 {
2228 out.push((slot.group.clone(), slot.visible));
2229 }
2230 }
2231 let mut copied: HashSet<u64> = out.iter().map(|(group, _)| group.sequence).collect();
2233 for (sequence, slot) in state.lookup.iter() {
2234 if !slot.group.is_aborted() && copied.insert(*sequence) {
2235 out.push((slot.group.clone(), slot.visible));
2236 }
2237 }
2238 out
2239 }
2240 ConsumerKind::Spliced(resume) => resume.cached_groups(),
2241 }
2242 }
2243
2244 pub(crate) fn cached_group(&self, sequence: u64, frame_start: u64) -> Option<group::Consumer> {
2249 match &self.inner {
2250 ConsumerKind::Plain(state) => {
2251 let state = state.read();
2252 let group = state.covering_group(sequence, frame_start)?;
2253 group.cache_refresh();
2254 Some(group.consume())
2255 }
2256 ConsumerKind::Spliced(resume) => resume.cached_group(sequence, frame_start),
2257 }
2258 }
2259
2260 pub(crate) fn cached_info(&self) -> Option<Info> {
2262 match &self.inner {
2263 ConsumerKind::Plain(state) => state.read().info.clone(),
2264 ConsumerKind::Spliced(resume) => resume.cached_info(),
2265 }
2266 }
2267
2268 pub fn name(&self) -> &str {
2270 &self.name
2271 }
2272
2273 pub fn broadcast(&self) -> &broadcast::Info {
2277 &self.broadcast
2278 }
2279
2280 pub fn subscribe(&self, subscription: impl Into<Option<Subscription>>) -> kio::Pending<Subscribing> {
2291 let subscription = kio::Producer::new(subscription.into().unwrap_or_default());
2292
2293 let inner = match &self.inner {
2294 ConsumerKind::Plain(state) => {
2295 register_subscription(state.read(), &subscription);
2298 SubscribingKind::Plain(state.clone())
2299 }
2300 ConsumerKind::Spliced(resume) => SubscribingKind::Spliced(resume.clone()),
2302 };
2303
2304 kio::Pending::new(Subscribing {
2305 name: self.name.clone(),
2306 broadcast: self.broadcast.clone(),
2307 inner,
2308 subscription,
2309 stats: self.stats.clone(),
2310 })
2311 }
2312
2313 pub(crate) fn peek_latest(&self) -> Option<group::Consumer> {
2317 match &self.inner {
2318 ConsumerKind::Plain(state) => {
2319 let sequence = state.read().max_sequence?;
2320 self.peek_group(sequence)
2321 }
2322 ConsumerKind::Spliced(resume) => resume.peek_latest(),
2323 }
2324 }
2325
2326 pub(crate) fn peek_before(&self, sequence: u64) -> Option<group::Consumer> {
2330 match &self.inner {
2331 ConsumerKind::Plain(state) => {
2332 let state = state.read();
2333 state
2334 .lookup
2335 .range(..sequence)
2336 .rev()
2337 .map(|(_, slot)| &slot.group)
2338 .find(|group| !group.is_aborted())
2339 .map(|group| group.consume())
2340 }
2341 ConsumerKind::Spliced(resume) => resume.peek_before(sequence),
2342 }
2343 }
2344
2345 pub(crate) fn peek_group(&self, sequence: u64) -> Option<group::Consumer> {
2350 match &self.inner {
2351 ConsumerKind::Plain(state) => {
2352 let state = state.read();
2353 let slot = state.lookup.get(&sequence)?;
2354 if slot.group.is_aborted() {
2355 return None;
2356 }
2357 Some(slot.group.consume())
2358 }
2359 ConsumerKind::Spliced(resume) => resume.peek_group(sequence),
2360 }
2361 }
2362
2363 pub(crate) fn guard_group(
2365 &self,
2366 group: group::Consumer,
2367 subscription: kio::Consumer<Subscription>,
2368 cap: kio::Consumer<Option<u64>>,
2369 bound: Option<u64>,
2370 ) -> group::Consumer {
2371 let ConsumerKind::Plain(state) = &self.inner else {
2372 return group;
2373 };
2374 let sequence = group.sequence;
2375 group.with_expiry(Arc::new(GroupExpiry {
2376 state: state.weak(),
2377 subscription,
2378 cap,
2379 bound,
2380 sequence,
2381 }))
2382 }
2383
2384 pub(crate) fn poll_peek_group(&self, sequence: u64, waiter: &kio::Waiter) -> Poll<Option<group::Consumer>> {
2392 let ConsumerKind::Plain(state) = &self.inner else {
2393 return Poll::Pending;
2395 };
2396
2397 let res = state.poll(waiter, |state| {
2398 match state.lookup.get(&sequence) {
2399 Some(slot) if !slot.group.is_aborted() => Poll::Ready(Some(slot.group.consume())),
2400 Some(_) => Poll::Ready(None),
2402 None if state.final_sequence.is_some_and(|fin| sequence >= fin) => Poll::Ready(None),
2404 None if state.start_sequence.is_some_and(|start| sequence < start) => Poll::Ready(None),
2406 None => Poll::Pending,
2407 }
2408 });
2409
2410 match res {
2411 Poll::Ready(Ok(res)) => Poll::Ready(res),
2412 Poll::Ready(Err(_)) => Poll::Ready(None),
2414 Poll::Pending => Poll::Pending,
2415 }
2416 }
2417
2418 pub(crate) fn poll_serving_group(&self, sequence: u64, index: u64, waiter: &kio::Waiter) -> Poll<()> {
2433 let ConsumerKind::Plain(state) = &self.inner else {
2434 return Poll::Pending;
2436 };
2437 let res = state.poll(waiter, |state| match state.lookup.get(&sequence) {
2438 Some(slot) if !slot.group.is_aborted() => {
2439 let mut group = slot.group.consume();
2442 group.start_at(index);
2443 match group.index() == index {
2444 true => Poll::Ready(()),
2445 false => Poll::Pending,
2446 }
2447 }
2448 _ => Poll::Pending,
2449 });
2450 match res {
2451 Poll::Ready(Ok(())) => Poll::Ready(()),
2452 Poll::Ready(Err(_)) | Poll::Pending => Poll::Pending,
2455 }
2456 }
2457
2458 pub fn fetch_group(&self, sequence: u64, options: impl Into<Option<group::Fetch>>) -> kio::Pending<Fetching> {
2470 let options = options.into().unwrap_or_default();
2471
2472 self.stats.fetch();
2476
2477 let state = match &self.inner {
2478 ConsumerKind::Plain(state) => state,
2479 ConsumerKind::Spliced(resume) => {
2482 return kio::Pending::new(Fetching {
2483 inner: FetchingKind::Spliced(resume.fetch_group(sequence, options)),
2484 stats: self.stats.clone(),
2485 });
2486 }
2487 };
2488
2489 let mut result = None;
2490
2491 let (fetch, unresolved) = {
2495 let state = state.read();
2496 (
2497 state.fetch.clone(),
2498 state.poll_fetch_cached(sequence, options.frame_start).is_pending(),
2499 )
2500 };
2501
2502 if unresolved {
2503 let mut fetch = fetch.lock();
2504 if let Some(pending) = fetch.join(&sequence) {
2505 pending.priority = pending.priority.max(options.priority);
2515 pending.frame_start = pending.frame_start.min(options.frame_start);
2516 result = Some(pending.result.consume());
2517 } else {
2518 let producer = kio::Producer::<FetchOutcome>::default();
2522 let consumer = producer.consume();
2523 let attempt = PendingFetch {
2524 priority: options.priority,
2525 frame_start: options.frame_start,
2526 result: producer,
2527 };
2528 if fetch.insert(sequence, attempt).is_ok() {
2529 result = Some(consumer);
2530 }
2531 }
2532 }
2533
2534 kio::Pending::new(Fetching {
2535 inner: FetchingKind::Plain {
2536 state: state.clone(),
2537 fetch,
2538 sequence,
2539 frame_start: options.frame_start,
2540 result,
2541 },
2542 stats: self.stats.clone(),
2543 })
2544 }
2545
2546 pub fn query(&self) -> kio::Pending<Querying> {
2553 kio::Pending::new(Querying {
2554 inner: match &self.inner {
2555 ConsumerKind::Plain(state) => QueryingKind::Plain(state.clone()),
2556 ConsumerKind::Spliced(resume) => QueryingKind::Spliced(resume.clone()),
2557 },
2558 })
2559 }
2560
2561 pub fn latest(&self) -> Option<u64> {
2563 match &self.inner {
2564 ConsumerKind::Plain(state) => state.read().max_sequence,
2565 ConsumerKind::Spliced(resume) => resume.latest(),
2566 }
2567 }
2568
2569 pub(crate) fn resume_position(&self) -> Option<Position> {
2574 match &self.inner {
2575 ConsumerKind::Plain(state) => state.read().resume_position(),
2576 ConsumerKind::Spliced(resume) => resume.resume_position(),
2577 }
2578 }
2579
2580 pub(crate) fn poll_complete(&self, waiter: &kio::Waiter) -> Poll<Result<()>> {
2585 let ConsumerKind::Plain(state) = &self.inner else {
2586 return Poll::Pending;
2588 };
2589 match ready!(state.poll(waiter, |state| {
2590 if state.is_complete() {
2591 Poll::Ready(())
2592 } else {
2593 Poll::Pending
2594 }
2595 })) {
2596 Ok(_) => Poll::Ready(Ok(())),
2597 Err(closed) => Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped))),
2600 }
2601 }
2602}
2603
2604pub struct Subscribing {
2607 name: Arc<str>,
2608 broadcast: Arc<broadcast::Info>,
2609 inner: SubscribingKind,
2610 subscription: kio::Producer<Subscription>,
2611 stats: stats::Scope,
2612}
2613
2614enum SubscribingKind {
2615 Plain(kio::Consumer<TrackState>),
2616 Spliced(super::resume::Consumer),
2617}
2618
2619impl Subscribing {
2620 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<Subscriber>> {
2623 match &self.inner {
2624 SubscribingKind::Plain(state) => {
2625 let info = ready!(state.poll(waiter, |state| state.poll_info()))
2627 .map_err(|e| e.abort.clone().unwrap_or(Error::Dropped))??;
2628
2629 let drift_cap = kio::Producer::new(None);
2630 let min_sequence = floor_of(&self.subscription.read());
2631 Poll::Ready(Ok(Subscriber {
2632 name: self.name.clone(),
2633 broadcast: self.broadcast.clone(),
2634 info,
2635 inner: SubscriberKind::Plain(PlainSubscriber {
2636 state: state.clone(),
2637 subscription: self.subscription.clone(),
2638 min_sequence,
2639 index: 0,
2640 datagram_index: 0,
2641 next_sequence: 0,
2642 end_sequence: None,
2643 parked: BTreeMap::new(),
2644 stale_cap: None,
2645 drift_cap,
2646 stale: stats::Content::default(),
2647 seek_pending: BTreeMap::new(),
2648 }),
2649 stats: self.stats.clone(),
2650 _stats_sub: self.stats.subscribe(),
2651 }))
2652 }
2653 SubscribingKind::Spliced(resume) => {
2654 let info = ready!(resume.poll_info(waiter))?;
2657
2658 Poll::Ready(Ok(Subscriber {
2659 name: self.name.clone(),
2660 broadcast: self.broadcast.clone(),
2661 info,
2662 inner: SubscriberKind::Spliced(Box::new(resume.subscribe_shared(self.subscription.clone()))),
2663 stats: self.stats.clone(),
2664 _stats_sub: self.stats.subscribe(),
2665 }))
2666 }
2667 }
2668 }
2669
2670 pub fn update(&mut self, subscription: Subscription) -> Result<()> {
2675 let mut state = self.subscription.write().map_err(|_| Error::Closed)?;
2676 *state = subscription;
2677 Ok(())
2678 }
2679}
2680
2681impl kio::Pollable for Subscribing {
2682 type Output = Result<Subscriber>;
2683
2684 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
2685 self.poll_ok(waiter)
2686 }
2687}
2688
2689pub struct Querying {
2692 inner: QueryingKind,
2693}
2694
2695enum QueryingKind {
2696 Plain(kio::Consumer<TrackState>),
2697 Spliced(super::resume::Consumer),
2698}
2699
2700impl Querying {
2701 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<Info>> {
2703 match &self.inner {
2704 QueryingKind::Plain(state) => {
2705 let info = ready!(state.poll(waiter, |state| state.poll_info()))
2707 .map_err(|e| e.abort.clone().unwrap_or(Error::Dropped))??;
2708 Poll::Ready(Ok(info))
2709 }
2710 QueryingKind::Spliced(resume) => resume.poll_info(waiter),
2711 }
2712 }
2713}
2714
2715impl kio::Pollable for Querying {
2716 type Output = Result<Info>;
2717
2718 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
2719 self.poll_ok(waiter)
2720 }
2721}
2722
2723impl group::Request {
2724 pub fn sequence(&self) -> u64 {
2726 self.sequence
2727 }
2728
2729 pub fn priority(&self) -> u8 {
2731 self.priority
2732 }
2733
2734 pub fn frame_start(&self) -> u64 {
2742 self.frame_start
2743 }
2744
2745 pub fn accept(mut self, info: impl Into<Option<Info>>) -> Result<group::Producer> {
2753 self.done = true;
2754 let res = TrackState::modify(&self.state)
2758 .and_then(|mut state| state.insert_group_request(self.sequence, self.frame_start, info.into()));
2759 self.remove();
2760 res
2761 }
2762
2763 pub fn reject(mut self, err: Error) {
2765 self.done = true;
2766 self.remove();
2769 if let Ok(mut outcome) = self.result.write() {
2770 outcome.rejected = Some(err);
2771 }
2772 }
2773
2774 fn remove(&self) {
2777 self.fetch
2778 .lock()
2779 .remove_if(&self.sequence, |pending| pending.result.same_channel(&self.result));
2780 }
2781}
2782
2783impl Drop for group::Request {
2784 fn drop(&mut self) {
2785 if self.done {
2786 return;
2787 }
2788 self.remove();
2789 if let Ok(mut outcome) = self.result.write() {
2790 outcome.rejected = Some(Error::Dropped);
2791 }
2792 }
2793}
2794
2795pub struct Fetching {
2801 inner: FetchingKind,
2802 stats: stats::Scope,
2805}
2806
2807enum FetchingKind {
2808 Plain {
2809 state: kio::Consumer<TrackState>,
2810 fetch: kio::Shared<FetchState>,
2811 sequence: u64,
2812 frame_start: u64,
2815 result: Option<kio::Consumer<FetchOutcome>>,
2817 },
2818 Spliced(kio::Pending<super::resume::Fetching>),
2820}
2821
2822impl kio::Pollable for Fetching {
2823 type Output = Result<group::Consumer>;
2824
2825 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
2826 let (state, fetch, sequence, frame_start, result) = match &self.inner {
2827 FetchingKind::Plain {
2828 state,
2829 fetch,
2830 sequence,
2831 frame_start,
2832 result,
2833 } => (state, fetch, *sequence, *frame_start, result.as_ref()),
2834 FetchingKind::Spliced(spliced) => {
2835 return kio::Pollable::poll(&**spliced, waiter)
2838 .map(|res| res.map(|group| group.with_meter(self.stats.meter())));
2839 }
2840 };
2841
2842 match state.poll(waiter, |state| state.poll_fetch_cached(sequence, frame_start)) {
2845 Poll::Ready(Ok(res)) => {
2846 return Poll::Ready(res.map(|mut group| {
2847 group.start_at(frame_start);
2851 group.with_meter(self.stats.meter())
2852 }));
2853 }
2854 Poll::Ready(Err(closed)) => {
2855 return Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped)));
2856 }
2857 Poll::Pending => {}
2858 }
2859
2860 let Some(result) = result else {
2862 return match fetch.poll(waiter, |fetch| match fetch.has_handlers() {
2865 false => Poll::Ready(()),
2866 true => Poll::Pending,
2867 }) {
2868 Poll::Ready(_guard) => Poll::Ready(Err(Error::NotFound)),
2869 Poll::Pending => Poll::Pending,
2870 };
2871 };
2872
2873 match result.poll(waiter, |outcome| match &outcome.rejected {
2876 Some(err) => Poll::Ready(err.clone()),
2877 None => Poll::Pending,
2878 }) {
2879 Poll::Ready(Ok(err)) => Poll::Ready(Err(err)),
2880 Poll::Ready(Err(_closed)) => Poll::Ready(Err(Error::NotFound)),
2881 Poll::Pending => Poll::Pending,
2882 }
2883 }
2884}
2885
2886pub struct Subscriber {
2914 name: Arc<str>,
2915 broadcast: Arc<broadcast::Info>,
2917 info: Info,
2918 inner: SubscriberKind,
2919 stats: stats::Scope,
2922 _stats_sub: stats::Subscription,
2925}
2926
2927enum SubscriberKind {
2928 Plain(PlainSubscriber),
2929 Spliced(Box<super::resume::Subscriber>),
2931}
2932
2933#[derive(Clone)]
2937struct Drift {
2938 budget: Duration,
2939 edge: Option<Edge>,
2940}
2941
2942struct GroupExpiry {
2944 state: kio::ConsumerWeak<TrackState>,
2950 subscription: kio::Consumer<Subscription>,
2951 cap: kio::Consumer<Option<u64>>,
2952 bound: Option<u64>,
2953 sequence: u64,
2954}
2955
2956impl group::Expiry for GroupExpiry {
2957 fn is_expired(&self, waiter: &kio::Waiter) -> bool {
2958 let mut max_age = Duration::default();
2959 let _ = self.subscription.poll(waiter, |subscription| {
2960 max_age = subscription.max_age;
2961 Poll::<()>::Pending
2962 });
2963
2964 let mut cap = None;
2965 let _ = self.cap.poll(waiter, |current| {
2966 cap = **current;
2967 Poll::<()>::Pending
2968 });
2969 let cap = super::subscription::min_some(cap, self.bound);
2970
2971 let mut expired = false;
2972 let _ = self.state.poll(waiter, |state| {
2973 let budget = clamp_max_age(max_age, state.max_age_bound());
2974 loop {
2975 let edge = state.live_edge(cap);
2976 expired = state.is_stale(self.sequence, edge.as_ref(), budget);
2977 if expired {
2978 break;
2979 }
2980
2981 let mut timestamp_raced = false;
2987 if let Some(slot) = state.lookup.get(&self.sequence) {
2988 let group = &slot.group;
2989 if group.timestamp().is_none()
2990 && group.poll_timestamp(waiter).is_ready()
2991 && group.timestamp().is_some()
2992 {
2993 timestamp_raced = true;
2994 }
2995 }
2996 for (_, slot) in state
2997 .lookup
2998 .range((std::ops::Bound::Excluded(self.sequence), std::ops::Bound::Unbounded))
2999 {
3000 let group = &slot.group;
3001 if !super::subscription::before_end(group.sequence, cap) {
3002 break;
3003 }
3004 if slot.visible
3005 && !group.is_aborted()
3006 && group.timestamp().is_none()
3007 && group.poll_timestamp(waiter).is_ready()
3008 && group.timestamp().is_some()
3009 {
3010 timestamp_raced = true;
3011 break;
3012 }
3013 }
3014 if !timestamp_raced {
3015 break;
3016 }
3017 }
3018
3019 Poll::<()>::Pending
3022 });
3023
3024 expired
3025 }
3026}
3027
3028#[derive(Clone)]
3031struct Edge {
3032 presentation: PresentationEdge,
3033 cap: Option<u64>,
3036}
3037
3038#[derive(Clone, Copy)]
3040struct PresentationEdge {
3041 sequence: u64,
3042 stamp: u32,
3045 timestamp: Timestamp,
3050}
3051
3052struct PlainSubscriber {
3054 state: kio::Consumer<TrackState>,
3055
3056 subscription: kio::Producer<Subscription>,
3057 index: usize,
3059 datagram_index: usize,
3061 min_sequence: u64,
3063 next_sequence: u64,
3066 end_sequence: Option<u64>,
3072 parked: BTreeMap<u64, group::Consumer>,
3077 stale_cap: Option<u64>,
3081 drift_cap: kio::Producer<Option<u64>>,
3083 stale: stats::Content,
3089 seek_pending: BTreeMap<u64, stats::Content>,
3092}
3093
3094impl PlainSubscriber {
3095 fn update_drift_cap(&mut self) {
3096 if let Ok(mut cap) = self.drift_cap.write() {
3097 *cap = servable_cap(self.end_sequence, self.stale_cap);
3098 }
3099 }
3100
3101 fn poll<F, R>(&self, waiter: &kio::Waiter, f: F) -> Poll<Result<R>>
3103 where
3104 F: Fn(&kio::Ref<'_, TrackState>) -> Poll<Result<R>>,
3105 {
3106 Poll::Ready(match ready!(self.state.poll(waiter, f)) {
3107 Ok(res) => res,
3108 Err(state) => Err(state.abort.clone().unwrap_or(Error::Dropped)),
3110 })
3111 }
3112
3113 fn take_stale(&mut self) -> stats::Content {
3115 std::mem::take(&mut self.stale)
3116 }
3117
3118 fn note_stale(&mut self, group: &group::Consumer) {
3121 self.stale.add(group.content());
3122 }
3123
3124 fn poll_drift(&self, cap: Option<u64>, waiter: &kio::Waiter) -> Poll<Result<Drift>> {
3135 let mut max_age = Duration::default();
3136 let _ = self.subscription.poll(waiter, |subscription| {
3137 max_age = subscription.max_age;
3138 Poll::<()>::Pending
3139 });
3140 self.poll(waiter, move |state| {
3141 Poll::Ready(Ok(Drift {
3142 budget: clamp_max_age(max_age, state.max_age_bound()),
3143 edge: state.live_edge(cap),
3144 }))
3145 })
3146 }
3147
3148 fn poll_stale(&self, group: &group::Consumer, drift: &Drift, waiter: &kio::Waiter) -> Poll<Result<bool>> {
3151 self.poll(waiter, move |state| {
3152 Poll::Ready(Ok(state.is_stale(group.sequence, drift.edge.as_ref(), drift.budget)))
3153 })
3154 }
3155
3156 fn with_expiry(&self, group: group::Consumer) -> group::Consumer {
3157 let sequence = group.sequence;
3158 group.with_expiry(Arc::new(GroupExpiry {
3159 state: self.state.weak(),
3160 subscription: self.subscription.consume(),
3161 cap: self.drift_cap.consume(),
3162 bound: None,
3163 sequence,
3164 }))
3165 }
3166
3167 fn poll_recv_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
3168 let watch = |group: &group::Consumer| match group.poll_closed(waiter) {
3176 Poll::Pending => true,
3177 Poll::Ready(()) => !group.is_aborted(),
3178 };
3179
3180 let min_sequence = self.min_sequence;
3185 self.parked
3186 .retain(|sequence, group| *sequence >= min_sequence && watch(group));
3187
3188 let drift = ready!(self.poll_drift(servable_cap(self.end_sequence, self.stale_cap), waiter))?;
3190
3191 loop {
3192 let consumer = match self.parked.keys().next().copied() {
3195 Some(sequence) if super::subscription::before_end(sequence, self.end_sequence) => {
3196 let group = self.parked.remove(&sequence).expect("just looked it up");
3197 group.cache_refresh();
3199 group
3200 }
3201 _ => {
3202 let Some((producer, found_index)) =
3203 ready!(self.poll(waiter, |state| state.poll_recv_group(self.index, self.min_sequence))?)
3204 else {
3205 if self.parked.is_empty() {
3208 return Poll::Ready(Ok(None));
3209 }
3210 return Poll::Pending;
3211 };
3212 let consumer = producer.consume();
3213 consumer.cache_refresh();
3216 self.index = found_index + 1;
3217
3218 if !super::subscription::before_end(consumer.sequence, self.end_sequence) {
3221 if watch(&consumer) {
3225 self.parked.insert(consumer.sequence, consumer);
3226 }
3227 continue;
3228 }
3229 consumer
3230 }
3231 };
3232
3233 if ready!(self.poll_stale(&consumer, &drift, waiter))? {
3236 self.stale.add(consumer.content());
3237 continue;
3238 }
3239 return Poll::Ready(Ok(Some(self.with_expiry(consumer))));
3240 }
3241 }
3242
3243 fn poll_recv_datagram(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Datagram>>> {
3244 let Some((datagram, found_index)) =
3245 ready!(self.poll(waiter, |state| state.poll_recv_datagram(self.datagram_index))?)
3246 else {
3247 return Poll::Ready(Ok(None));
3248 };
3249
3250 self.datagram_index = found_index + 1;
3251 Poll::Ready(Ok(Some(datagram)))
3252 }
3253
3254 fn poll_next_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
3255 let floor = self.next_sequence.max(self.min_sequence);
3256 let Some(group) = ready!(self.poll_seek_group(floor, self.end_sequence, waiter))? else {
3257 return Poll::Ready(Ok(None));
3258 };
3259 self.next_sequence = group.sequence.saturating_add(1);
3260 self.commit_seek_stale(self.next_sequence);
3262 group.cache_refresh();
3264 Poll::Ready(Ok(Some(group)))
3265 }
3266
3267 fn poll_seek_group(
3280 &mut self,
3281 floor: u64,
3282 end: Option<u64>,
3283 waiter: &kio::Waiter,
3284 ) -> Poll<Result<Option<group::Consumer>>> {
3285 let mut floor = floor.max(self.min_sequence);
3286 let end = super::subscription::min_some(end, self.end_sequence);
3287 let drift = ready!(self.poll_drift(servable_cap(end, self.stale_cap), waiter))?;
3289
3290 loop {
3291 let Some(producer) = ready!(self.poll(waiter, |state| state.poll_next_in_range(floor, end))?) else {
3292 return Poll::Ready(Ok(None));
3299 };
3300 let group = producer.consume();
3301
3302 if ready!(self.poll_stale(&group, &drift, waiter))? {
3305 self.seek_pending.insert(group.sequence, group.content());
3310 floor = group.sequence.saturating_add(1);
3311 continue;
3312 }
3313
3314 self.seek_pending.remove(&group.sequence);
3317 return Poll::Ready(Ok(Some(self.with_expiry(group))));
3318 }
3319 }
3320
3321 fn commit_seek_stale(&mut self, committed: u64) {
3332 let keep = self.seek_pending.split_off(&committed);
3333 for (_, content) in std::mem::replace(&mut self.seek_pending, keep) {
3334 self.stale.add(content);
3335 }
3336 }
3337
3338 fn discard_seek_conviction(&mut self, sequence: u64) {
3342 self.seek_pending.remove(&sequence);
3343 }
3344}
3345
3346#[derive(Clone)]
3352pub struct Control {
3353 subscription: kio::Producer<Subscription>,
3354}
3355
3356impl Control {
3357 pub fn subscription(&self) -> Subscription {
3359 self.subscription.read().clone()
3360 }
3361
3362 pub fn update(&self, subscription: Subscription) -> Result<()> {
3367 let mut state = self.subscription.write().map_err(|_| Error::Closed)?;
3368 *state = subscription;
3369 Ok(())
3370 }
3371}
3372
3373impl Subscriber {
3374 pub fn info(&self) -> &Info {
3379 &self.info
3380 }
3381
3382 pub fn name(&self) -> &str {
3384 &self.name
3385 }
3386
3387 pub fn broadcast(&self) -> &broadcast::Info {
3391 &self.broadcast
3392 }
3393
3394 fn count_stale(&mut self, meter: &stats::Meter) {
3400 if meter.is_tracked() {
3405 meter.stale(self.take_stale());
3406 }
3407 }
3408
3409 pub(crate) fn take_stale(&mut self) -> stats::Content {
3412 match &mut self.inner {
3413 SubscriberKind::Plain(plain) => plain.take_stale(),
3414 SubscriberKind::Spliced(spliced) => spliced.take_stale(),
3415 }
3416 }
3417
3418 pub(crate) fn set_stale_cap(&mut self, cap: Option<u64>) {
3428 match &mut self.inner {
3429 SubscriberKind::Plain(plain) => {
3430 plain.stale_cap = cap;
3431 plain.update_drift_cap();
3432 }
3433 SubscriberKind::Spliced(spliced) => spliced.set_stale_cap(cap),
3434 }
3435 }
3436
3437 pub(crate) fn commit_seek_stale(&mut self, committed: u64) {
3444 match &mut self.inner {
3445 SubscriberKind::Plain(plain) => plain.commit_seek_stale(committed),
3446 SubscriberKind::Spliced(spliced) => spliced.commit_seek_stale(committed),
3447 }
3448 }
3449
3450 pub(crate) fn discard_seek_conviction(&mut self, sequence: u64) {
3455 match &mut self.inner {
3456 SubscriberKind::Plain(plain) => plain.discard_seek_conviction(sequence),
3457 SubscriberKind::Spliced(spliced) => spliced.discard_seek_conviction(sequence),
3458 }
3459 }
3460
3461 pub(crate) fn poll_stale(&mut self, group: &group::Consumer, waiter: &kio::Waiter) -> Poll<Result<bool>> {
3467 let plain = match &mut self.inner {
3468 SubscriberKind::Plain(plain) => plain,
3469 SubscriberKind::Spliced(spliced) => return spliced.poll_stale(group, waiter),
3470 };
3471 let drift = ready!(plain.poll_drift(servable_cap(plain.end_sequence, plain.stale_cap), waiter))?;
3472 let stale = ready!(plain.poll_stale(group, &drift, waiter))?;
3473 if stale {
3474 plain.note_stale(group);
3475 }
3476 Poll::Ready(Ok(stale))
3477 }
3478
3479 pub fn control(&self) -> Control {
3481 Control {
3482 subscription: match &self.inner {
3483 SubscriberKind::Plain(plain) => plain.subscription.clone(),
3484 SubscriberKind::Spliced(spliced) => spliced.prefs(),
3485 },
3486 }
3487 }
3488
3489 pub fn poll_recv_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
3518 let meter = self.stats.meter();
3519 let res = match &mut self.inner {
3520 SubscriberKind::Plain(plain) => plain.poll_recv_group(waiter),
3521 SubscriberKind::Spliced(spliced) => spliced.poll_recv_group(waiter),
3522 };
3523 self.count_stale(&meter);
3524 res.map(|res| res.map(|group| group.map(|group| group.with_meter(meter))))
3525 }
3526
3527 pub async fn recv_group(&mut self) -> Result<Option<group::Consumer>> {
3534 kio::wait(|waiter| self.poll_recv_group(waiter)).await
3535 }
3536
3537 pub fn poll_recv_datagram(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Datagram>>> {
3548 let meter = self.stats.meter();
3549 let res = match &mut self.inner {
3550 SubscriberKind::Plain(plain) => plain.poll_recv_datagram(waiter),
3551 SubscriberKind::Spliced(spliced) => spliced.poll_recv_datagram(waiter),
3552 };
3553 if let Poll::Ready(Ok(Some(datagram))) = &res {
3556 meter.datagram(datagram.payload.len() as u64);
3557 }
3558 res
3559 }
3560
3561 pub async fn recv_datagram(&mut self) -> Result<Option<Datagram>> {
3568 kio::wait(|waiter| self.poll_recv_datagram(waiter)).await
3569 }
3570
3571 fn poll_next_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
3573 let meter = self.stats.meter();
3574 let res = match &mut self.inner {
3575 SubscriberKind::Plain(plain) => plain.poll_next_group(waiter),
3576 SubscriberKind::Spliced(spliced) => spliced.poll_next_group(waiter),
3577 };
3578 self.count_stale(&meter);
3579 res.map(|res| res.map(|group| group.map(|group| group.with_meter(meter))))
3580 }
3581
3582 pub(crate) fn poll_seek_group(
3587 &mut self,
3588 floor: u64,
3589 end: Option<u64>,
3590 waiter: &kio::Waiter,
3591 ) -> Poll<Result<Option<group::Consumer>>> {
3592 let meter = self.stats.meter();
3593 let res = match &mut self.inner {
3594 SubscriberKind::Plain(plain) => plain.poll_seek_group(floor, end, waiter),
3595 SubscriberKind::Spliced(spliced) => spliced.poll_seek_group(floor, end, waiter),
3596 };
3597 self.count_stale(&meter);
3598 res.map(|res| res.map(|group| group.map(|group| group.with_meter(meter))))
3599 }
3600
3601 pub fn ordered(self) -> Ordered {
3611 Ordered { inner: self }
3612 }
3613
3614 pub fn is_clone(&self, other: &Self) -> bool {
3616 match (&self.inner, &other.inner) {
3617 (SubscriberKind::Plain(a), SubscriberKind::Plain(b)) => a.state.same_channel(&b.state),
3618 (SubscriberKind::Spliced(a), SubscriberKind::Spliced(b)) => a.is_clone(b),
3619 _ => false,
3620 }
3621 }
3622
3623 pub fn poll_finished(&mut self, waiter: &kio::Waiter) -> Poll<Result<u64>> {
3625 match &mut self.inner {
3626 SubscriberKind::Plain(plain) => plain.poll(waiter, |state| state.poll_finished()),
3627 SubscriberKind::Spliced(spliced) => spliced.poll_finished(waiter),
3628 }
3629 }
3630
3631 pub async fn finished(&mut self) -> Result<u64> {
3639 kio::wait(|waiter| self.poll_finished(waiter)).await
3640 }
3641
3642 pub fn set_groups(&mut self, groups: impl RangeBounds<u64>) {
3650 let (start, end) = super::subscription::sequence_bounds(groups);
3651 self.raise_start_to(start);
3652 self.end_at(end.map_or(Bound::Unbounded, Bound::Excluded));
3653 }
3654
3655 pub(crate) fn start_at(&mut self, sequence: u64) {
3662 match &mut self.inner {
3663 SubscriberKind::Plain(plain) => plain.min_sequence = sequence,
3664 SubscriberKind::Spliced(spliced) => spliced.start_at(sequence),
3665 }
3666 }
3667
3668 pub(crate) fn raise_start_to(&mut self, sequence: u64) {
3674 match &mut self.inner {
3675 SubscriberKind::Plain(plain) => plain.min_sequence = plain.min_sequence.max(sequence),
3676 SubscriberKind::Spliced(spliced) => spliced.raise_start_to(sequence),
3677 }
3678 }
3679
3680 pub(crate) fn end_at(&mut self, end: impl Into<Cap>) {
3694 let end = end.into();
3695 match &mut self.inner {
3696 SubscriberKind::Plain(plain) => {
3697 plain.end_sequence = end.exclusive();
3698 plain.update_drift_cap();
3699 }
3700 SubscriberKind::Spliced(spliced) => spliced.end_at(end),
3701 }
3702 }
3703
3704 pub fn subscription(&self) -> Subscription {
3706 self.control().subscription()
3707 }
3708
3709 pub fn update(&mut self, subscription: Subscription) -> Result<()> {
3715 match &mut self.inner {
3716 SubscriberKind::Plain(plain) => {
3717 let mut state = plain.subscription.write().map_err(|_| Error::Closed)?;
3718 *state = subscription;
3719 }
3720 SubscriberKind::Spliced(spliced) => spliced.update(subscription),
3721 }
3722 Ok(())
3723 }
3724
3725 pub fn latest(&self) -> Option<u64> {
3727 match &self.inner {
3728 SubscriberKind::Plain(plain) => plain.state.read().max_sequence,
3729 SubscriberKind::Spliced(spliced) => spliced.latest(),
3730 }
3731 }
3732}
3733
3734pub struct Ordered {
3761 inner: Subscriber,
3762}
3763
3764impl Ordered {
3765 pub fn info(&self) -> &Info {
3767 self.inner.info()
3768 }
3769
3770 pub fn name(&self) -> &str {
3772 self.inner.name()
3773 }
3774
3775 pub fn broadcast(&self) -> &broadcast::Info {
3777 self.inner.broadcast()
3778 }
3779
3780 pub fn poll_next_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
3792 self.inner.poll_next_group(waiter)
3793 }
3794
3795 pub async fn next_group(&mut self) -> Result<Option<group::Consumer>> {
3797 kio::wait(|waiter| self.poll_next_group(waiter)).await
3798 }
3799
3800 pub fn poll_recv_datagram(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Datagram>>> {
3808 self.inner.poll_recv_datagram(waiter)
3809 }
3810
3811 pub async fn recv_datagram(&mut self) -> Result<Option<Datagram>> {
3818 kio::wait(|waiter| self.poll_recv_datagram(waiter)).await
3819 }
3820
3821 pub fn set_groups(&mut self, groups: impl RangeBounds<u64>) {
3825 self.inner.set_groups(groups);
3826 }
3827
3828 pub fn control(&self) -> Control {
3831 self.inner.control()
3832 }
3833
3834 pub fn subscription(&self) -> Subscription {
3836 self.inner.subscription()
3837 }
3838
3839 pub fn update(&mut self, subscription: Subscription) -> Result<()> {
3843 self.inner.update(subscription)
3844 }
3845
3846 pub fn poll_finished(&mut self, waiter: &kio::Waiter) -> Poll<Result<u64>> {
3848 self.inner.poll_finished(waiter)
3849 }
3850
3851 pub async fn finished(&mut self) -> Result<u64> {
3856 kio::wait(|waiter| self.poll_finished(waiter)).await
3857 }
3858
3859 pub fn latest(&self) -> Option<u64> {
3861 self.inner.latest()
3862 }
3863
3864 pub fn is_clone(&self, other: &Self) -> bool {
3866 self.inner.is_clone(&other.inner)
3867 }
3868}
3869
3870pub struct Request {
3882 name: Arc<str>,
3883 broadcast: Arc<broadcast::Info>,
3885 state: kio::Producer<TrackState>,
3886
3887 prev_subscription: Option<Subscription>,
3889
3890 alive: Arc<Alive>,
3893
3894 _dynamic: Dynamic,
3899
3900 stats: stats::Scope,
3903}
3904
3905impl Request {
3906 pub(crate) fn new(broadcast: Arc<broadcast::Info>, name: impl Into<Arc<str>>) -> Self {
3907 let name = name.into();
3908 let state = TrackState::spawn(broadcast.clone());
3909 let alive = Alive::new(name.clone(), state.clone());
3910 let dynamic = Dynamic::new(name.clone(), state.clone(), alive.clone());
3911 Self {
3912 name,
3913 broadcast,
3914 state,
3915 prev_subscription: None,
3916 alive,
3917 _dynamic: dynamic,
3918 stats: stats::Scope::default(),
3919 }
3920 }
3921
3922 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
3925 self.stats = scope;
3926 self
3927 }
3928
3929 pub fn name(&self) -> &str {
3931 &self.name
3932 }
3933
3934 pub fn consume(&self) -> Consumer {
3936 Consumer::plain(self.name.clone(), self.state.consume())
3937 }
3938
3939 pub fn dynamic(&self) -> Dynamic {
3943 Dynamic::new(self.name.clone(), self.state.clone(), self.alive.clone())
3944 }
3945
3946 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
3949 self.state.poll_unused(waiter).map(|_| ())
3950 }
3951
3952 pub fn accept(self, info: impl Into<Option<Info>>) -> Producer {
3959 let info = TrackState::normalize_info(&self.broadcast, info.into().unwrap_or_default());
3960 if let Ok(mut state) = self.state.write() {
3963 state.accept(info.clone());
3964 }
3965 self.alive.publish(Some(&self.stats));
3968 Producer {
3969 name: self.name,
3970 info,
3971 broadcast: self.broadcast,
3972 state: self.state,
3973 prev_subscription: None,
3974 alive: self.alive,
3975 stats: self.stats,
3976 }
3977 }
3978
3979 pub fn reject(self, err: Error) {
3981 if let Ok(mut state) = self.state.write() {
3982 state.abort = Some(err);
3983 }
3984 }
3985
3986 pub fn subscription(&self) -> Option<Subscription> {
3989 let state = self.state.read();
3990 let (subs, bound) = (state.subscriptions.clone(), state.max_age_bound());
3991 drop(state);
3992 snapshot_subscription(&subs, bound)
3993 }
3994
3995 pub async fn subscription_changed(&mut self) -> Option<Subscription> {
3998 kio::wait(|waiter| self.poll_subscription_changed(waiter)).await
3999 }
4000
4001 pub fn poll_subscription_changed(&mut self, waiter: &kio::Waiter) -> Poll<Option<Subscription>> {
4003 let state = self.state.read();
4004 let (subs, bound) = (state.subscriptions.clone(), state.max_age_bound());
4005 drop(state);
4006
4007 let prev = &self.prev_subscription;
4008 let mut combined = None;
4009 let mut guard = ready!(subs.poll(waiter, |subs| {
4010 let next = combined_subscription(subs, bound, waiter);
4011 if &next == prev {
4012 Poll::Pending
4013 } else {
4014 combined = next;
4015 Poll::Ready(())
4016 }
4017 }));
4018 guard.retain(|sub| !sub.is_closed());
4020 drop(guard);
4021 self.prev_subscription = combined.clone();
4022 Poll::Ready(combined)
4023 }
4024
4025 pub(super) fn weak(&self) -> TrackWeak {
4026 TrackWeak {
4027 name: self.name.clone(),
4028 state: self.state.weak(),
4029 }
4030 }
4031}
4032
4033#[cfg(test)]
4034use futures::FutureExt;
4035
4036#[cfg(test)]
4037#[allow(missing_docs)] impl Subscriber {
4039 pub fn assert_group(&mut self) -> group::Consumer {
4040 self.recv_group()
4041 .now_or_never()
4042 .expect("group would have blocked")
4043 .expect("would have errored")
4044 .expect("track was closed")
4045 }
4046
4047 pub fn assert_no_group(&mut self) {
4048 assert!(
4049 self.recv_group().now_or_never().is_none(),
4050 "recv_group would not have blocked"
4051 );
4052 }
4053
4054 pub fn assert_not_closed(&mut self) {
4055 assert!(self.finished().now_or_never().is_none(), "should not be closed");
4056 }
4057
4058 pub fn assert_closed(&mut self) {
4059 assert!(self.finished().now_or_never().is_some(), "should be closed");
4060 }
4061
4062 pub fn assert_error(&mut self) {
4064 assert!(
4065 self.finished().now_or_never().expect("should not block").is_err(),
4066 "should be error"
4067 );
4068 }
4069
4070 pub fn assert_is_clone(&self, other: &Self) {
4071 assert!(self.is_clone(other), "should be clone");
4072 }
4073
4074 pub fn assert_not_clone(&self, other: &Self) {
4075 assert!(!self.is_clone(other), "should not be clone");
4076 }
4077}
4078
4079#[cfg(test)]
4080mod test {
4081 use super::*;
4082 use crate::frame;
4083 use crate::model::test_tracing::count_drop_warnings;
4084 use std::time::Duration;
4085
4086 fn track_producer(name: impl Into<Arc<str>>, info: impl Into<Option<Info>>) -> Producer {
4089 Producer::new(Arc::new(broadcast::Info::default()), name, info)
4090 }
4091
4092 fn replay() -> Subscription {
4094 Subscription::default().with_max_age(Duration::from_secs(30))
4095 }
4096
4097 fn live_groups(state: &TrackState) -> usize {
4099 state.lookup.len()
4100 }
4101
4102 fn first_live_sequence(state: &TrackState) -> u64 {
4104 state
4105 .arrival
4106 .iter()
4107 .find(|(sequence, stamp)| state.lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp))
4108 .map(|(sequence, _)| *sequence)
4109 .unwrap()
4110 }
4111
4112 fn recv_datagram(dg: &mut Subscriber) -> Datagram {
4114 dg.recv_datagram()
4115 .now_or_never()
4116 .expect("datagram would have blocked")
4117 .expect("would have errored")
4118 .expect("track was closed")
4119 }
4120
4121 #[tokio::test]
4125 async fn peek_resolves_below_the_declared_start() {
4126 let mut producer = track_producer("test", None);
4127 let consumer = producer.consume();
4128
4129 let mut cached = producer.create_group(group::Info { sequence: 1 }).unwrap();
4130 cached.write_frame(Timestamp::ZERO, b"backfill".to_vec()).unwrap();
4131 cached.finish().unwrap();
4132
4133 let waiter = kio::Waiter::noop();
4135 assert!(consumer.poll_peek_group(0, &waiter).is_pending());
4136
4137 producer.start_at(3).unwrap();
4140 assert!(matches!(consumer.poll_peek_group(0, &waiter), Poll::Ready(None)));
4141 assert!(matches!(consumer.poll_peek_group(1, &waiter), Poll::Ready(Some(_))));
4142 assert!(consumer.poll_peek_group(3, &waiter).is_pending());
4143
4144 producer.start_at(4).unwrap();
4147 assert!(matches!(consumer.poll_peek_group(3, &waiter), Poll::Ready(None)));
4148 producer.start_at(0).unwrap();
4149 assert!(consumer.poll_peek_group(0, &waiter).is_pending());
4150 }
4151
4152 #[tokio::test]
4153 async fn append_datagram_shares_group_sequence() {
4154 let mut producer = track_producer("test", None);
4155 let ts = Timestamp::from_millis(10).unwrap();
4156
4157 assert_eq!(producer.append_group().unwrap().sequence, 0);
4159 assert_eq!(producer.append_datagram(ts, &b"a"[..]).unwrap(), 1);
4160 assert_eq!(producer.append_group().unwrap().sequence, 2);
4161 assert_eq!(producer.append_datagram(ts, &b"b"[..]).unwrap(), 3);
4162 assert_eq!(producer.latest(), Some(3));
4163 }
4164
4165 #[tokio::test]
4166 async fn append_datagram_roundtrip() {
4167 let mut producer = track_producer("test", None);
4168 let mut dg = producer.subscribe(None);
4169
4170 let ts = Timestamp::from_millis(42).unwrap();
4171 let seq = producer.append_datagram(ts, &b"hello"[..]).unwrap();
4172
4173 let got = recv_datagram(&mut dg);
4174 assert_eq!(got.sequence, seq);
4175 assert_eq!(got.timestamp, ts);
4176 assert_eq!(&got.payload[..], b"hello");
4177 }
4178
4179 #[tokio::test]
4180 async fn insert_datagram_preserves_sequence() {
4181 let mut producer = track_producer("test", None);
4182 let mut dg = producer.subscribe(None);
4183
4184 let ts = Timestamp::from_millis(5).unwrap();
4185 producer
4187 .insert_datagram(100, ts, bytes::Bytes::from_static(b"x"))
4188 .unwrap();
4189
4190 assert_eq!(recv_datagram(&mut dg).sequence, 100);
4191 assert_eq!(producer.append_group().unwrap().sequence, 101);
4193 }
4194
4195 #[tokio::test]
4196 async fn insert_datagram_leaves_a_gap() {
4197 let mut producer = track_producer("test", None);
4198 let mut dg = producer.subscribe(None);
4199 let ts = Timestamp::from_millis(0).unwrap();
4200
4201 producer
4202 .insert_datagram(10, ts, bytes::Bytes::from_static(b"gap"))
4203 .unwrap();
4204 assert_eq!(recv_datagram(&mut dg).sequence, 10);
4205 assert_eq!(producer.append_datagram(ts, &b"next"[..]).unwrap(), 11);
4206 assert_eq!(producer.append_group().unwrap().sequence, 12);
4207 }
4208
4209 #[tokio::test]
4210 async fn insert_datagram_out_of_order_does_not_rewind() {
4211 let mut producer = track_producer("test", None);
4212 let mut dg = producer.subscribe(None);
4213 let ts = Timestamp::from_millis(0).unwrap();
4214
4215 producer
4216 .insert_datagram(10, ts, bytes::Bytes::from_static(b"high"))
4217 .unwrap();
4218 producer
4219 .insert_datagram(5, ts, bytes::Bytes::from_static(b"low"))
4220 .unwrap();
4221
4222 assert_eq!(recv_datagram(&mut dg).sequence, 10);
4223 assert_eq!(recv_datagram(&mut dg).sequence, 5);
4224 assert_eq!(producer.append_datagram(ts, &b"next"[..]).unwrap(), 11);
4225 }
4226
4227 #[tokio::test]
4228 async fn insert_datagram_duplicate_is_best_effort() {
4229 let mut producer = track_producer("test", None);
4230 let mut dg = producer.subscribe(None);
4231 let ts = Timestamp::from_millis(0).unwrap();
4232
4233 producer
4234 .insert_datagram(3, ts, bytes::Bytes::from_static(b"first"))
4235 .unwrap();
4236 producer
4237 .insert_datagram(3, ts, bytes::Bytes::from_static(b"again"))
4238 .unwrap();
4239
4240 assert_eq!(&recv_datagram(&mut dg).payload[..], b"first");
4241 assert_eq!(&recv_datagram(&mut dg).payload[..], b"again");
4242 assert_eq!(producer.append_datagram(ts, &b"next"[..]).unwrap(), 4);
4243 }
4244
4245 #[tokio::test]
4246 async fn insert_datagram_stale_does_not_rewind_after_append() {
4247 let mut producer = track_producer("test", None);
4248 let mut dg = producer.subscribe(None);
4249 let ts = Timestamp::from_millis(0).unwrap();
4250
4251 assert_eq!(producer.append_datagram(ts, &b"0"[..]).unwrap(), 0);
4252 assert_eq!(producer.append_datagram(ts, &b"1"[..]).unwrap(), 1);
4253 producer
4254 .insert_datagram(0, ts, bytes::Bytes::from_static(b"stale"))
4255 .unwrap();
4256
4257 assert_eq!(recv_datagram(&mut dg).sequence, 0);
4258 assert_eq!(recv_datagram(&mut dg).sequence, 1);
4259 assert_eq!(recv_datagram(&mut dg).sequence, 0);
4260 assert_eq!(producer.append_datagram(ts, &b"2"[..]).unwrap(), 2);
4261 assert_eq!(producer.append_group().unwrap().sequence, 3);
4262 }
4263
4264 #[tokio::test]
4265 async fn insert_datagram_cloned_producers_share_counter() {
4266 let mut producer = track_producer("test", None);
4267 let mut other = producer.clone();
4268 let mut dg = producer.subscribe(None);
4269 let ts = Timestamp::from_millis(0).unwrap();
4270
4271 producer
4272 .insert_datagram(4, ts, bytes::Bytes::from_static(b"a"))
4273 .unwrap();
4274 assert_eq!(other.append_datagram(ts, &b"b"[..]).unwrap(), 5);
4275 other.insert_datagram(8, ts, bytes::Bytes::from_static(b"c")).unwrap();
4276 assert_eq!(producer.append_group().unwrap().sequence, 9);
4277
4278 assert_eq!(recv_datagram(&mut dg).sequence, 4);
4279 assert_eq!(recv_datagram(&mut dg).sequence, 5);
4280 assert_eq!(recv_datagram(&mut dg).sequence, 8);
4281 }
4282
4283 #[test]
4284 fn insert_datagram_after_finish_is_closed() {
4285 let mut producer = track_producer("test", None);
4286 let ts = Timestamp::from_millis(0).unwrap();
4287 producer.finish().unwrap();
4288 assert!(matches!(
4289 producer.insert_datagram(0, ts, bytes::Bytes::from_static(b"x")),
4290 Err(Error::Closed)
4291 ));
4292 assert!(matches!(producer.append_datagram(ts, &b"x"[..]), Err(Error::Closed)));
4293 }
4294
4295 #[test]
4296 fn insert_datagram_after_abort_fails() {
4297 let producer = track_producer("test", None);
4298 let mut other = producer.clone();
4299 let ts = Timestamp::from_millis(0).unwrap();
4300 producer.abort(Error::Cancel).unwrap();
4301 assert!(other.insert_datagram(0, ts, bytes::Bytes::from_static(b"x")).is_err());
4302 }
4303
4304 #[test]
4305 fn insert_datagram_respects_finish_at() {
4306 let mut producer = track_producer("test", None);
4307 let ts = Timestamp::from_millis(0).unwrap();
4308 producer.finish_at(10).unwrap();
4309 producer
4310 .insert_datagram(5, ts, bytes::Bytes::from_static(b"ok"))
4311 .unwrap();
4312 assert!(matches!(
4313 producer.insert_datagram(10, ts, bytes::Bytes::from_static(b"late")),
4314 Err(Error::Closed)
4315 ));
4316 assert_eq!(producer.append_group().unwrap().sequence, 6);
4317 }
4318
4319 #[test]
4321 fn resume_position_uses_the_latest_group() {
4322 let mut datagram_only = track_producer("datagram-only", None);
4323 let datagram_only_consumer = datagram_only.consume();
4324 datagram_only
4325 .insert_datagram(8, Timestamp::ZERO, bytes::Bytes::from_static(b"x"))
4326 .unwrap();
4327 assert_eq!(
4328 datagram_only_consumer.resume_position(),
4329 None,
4330 "a datagram creates no group position to resume"
4331 );
4332
4333 let mut producer = track_producer("mixed", None);
4334 let consumer = producer.consume();
4335 let mut group = producer.create_group(group::Info { sequence: 3 }).unwrap();
4336 group
4337 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"head"))
4338 .unwrap();
4339 producer
4340 .insert_datagram(8, Timestamp::ZERO, bytes::Bytes::from_static(b"x"))
4341 .unwrap();
4342
4343 assert_eq!(
4344 consumer.resume_position(),
4345 Some(Position { group: 3, frame: 1 }),
4346 "the replacement must continue the open group"
4347 );
4348 group.abort(Error::Cancel).unwrap();
4349 }
4350
4351 #[tokio::test]
4354 async fn recv_datagram_leaves_the_ordered_cursor_alone() {
4355 let mut producer = track_producer("test", None);
4356 let mut datagrams = producer.subscribe(None);
4357 let mut subscriber = producer.subscribe(None).ordered();
4358 let ts = Timestamp::from_millis(5).unwrap();
4359
4360 producer
4361 .insert_datagram(5, ts, bytes::Bytes::from_static(b"x"))
4362 .unwrap();
4363 assert_eq!(recv_datagram(&mut datagrams).sequence, 5);
4364
4365 producer.create_group(group::Info { sequence: 3 }).unwrap();
4366 producer.create_group(group::Info { sequence: 6 }).unwrap();
4367
4368 let mut next = || {
4369 subscriber
4370 .next_group()
4371 .now_or_never()
4372 .expect("group would have blocked")
4373 .expect("would have errored")
4374 .expect("track was closed")
4375 .sequence
4376 };
4377 assert_eq!(next(), 3, "the datagram at sequence 5 did not consume group 3");
4378 assert_eq!(next(), 6);
4379 }
4380
4381 #[tokio::test]
4382 async fn datagram_normalized_to_track_timescale() {
4383 let info = Info::default().with_timescale(Timescale::MICRO);
4384 let mut producer = track_producer("test", info);
4385 let mut dg = producer.subscribe(None);
4386
4387 producer
4389 .append_datagram(Timestamp::from_millis(2).unwrap(), &b"z"[..])
4390 .unwrap();
4391 let got = recv_datagram(&mut dg);
4392 assert_eq!(got.timestamp.scale(), Timescale::MICRO);
4393 assert_eq!(got.timestamp.value(), 2_000);
4394 }
4395
4396 #[tokio::test]
4397 async fn datagram_rejects_oversized() {
4398 let mut producer = track_producer("test", None);
4399 let big = bytes::Bytes::from(vec![0u8; crate::model::datagram::MAX_DATAGRAM_PAYLOAD + 1]);
4400 let ts = Timestamp::from_millis(0).unwrap();
4401 assert!(matches!(
4402 producer.append_datagram(ts, big.clone()),
4403 Err(Error::FrameTooLarge)
4404 ));
4405 assert!(matches!(
4406 producer.insert_datagram(0, ts, big),
4407 Err(Error::FrameTooLarge)
4408 ));
4409 }
4410
4411 #[tokio::test]
4412 async fn datagram_fanout_to_subscribers() {
4413 let mut producer = track_producer("test", None);
4414 let mut a = producer.subscribe(None);
4416 let mut b = producer.subscribe(None);
4417 let ts = Timestamp::from_millis(1).unwrap();
4418
4419 producer.append_datagram(ts, &b"first"[..]).unwrap();
4420 producer.append_datagram(ts, &b"second"[..]).unwrap();
4421
4422 assert_eq!(&recv_datagram(&mut a).payload[..], b"first");
4424 assert_eq!(&recv_datagram(&mut a).payload[..], b"second");
4425 assert_eq!(&recv_datagram(&mut b).payload[..], b"first");
4426 assert_eq!(&recv_datagram(&mut b).payload[..], b"second");
4427 }
4428
4429 #[test]
4430 fn datagram_buffer_drops_oldest_at_capacity() {
4431 let mut producer = track_producer("test", None);
4432 let mut slow = producer.subscribe(None);
4433 let mut fast = producer.subscribe(None);
4434 let count = MAX_DATAGRAMS * 3;
4435 for sequence in 0..count {
4436 producer.append_datagram(Timestamp::ZERO, b"x".as_slice()).unwrap();
4437 assert_eq!(recv_datagram(&mut fast).sequence, sequence as u64);
4438 }
4439 assert_eq!(producer.state.read().datagrams.len(), MAX_DATAGRAMS);
4440 for sequence in count - MAX_DATAGRAMS..count {
4441 assert_eq!(recv_datagram(&mut slow).sequence, sequence as u64);
4442 }
4443 assert!(slow.poll_recv_datagram(&kio::Waiter::noop()).is_pending());
4444 }
4445
4446 #[tokio::test]
4447 async fn datagram_recv_pends_until_written() {
4448 let mut producer = track_producer("test", None);
4449 let mut dg = producer.subscribe(None);
4450
4451 assert!(
4452 dg.recv_datagram().now_or_never().is_none(),
4453 "should block with no datagrams"
4454 );
4455
4456 producer
4457 .append_datagram(Timestamp::from_millis(0).unwrap(), &b"go"[..])
4458 .unwrap();
4459 assert_eq!(&recv_datagram(&mut dg).payload[..], b"go");
4460 }
4461
4462 #[tokio::test]
4466 async fn datagram_wire_roundtrip_between_tracks() {
4467 use crate::coding::{Decode, Encode};
4468 use crate::lite;
4469
4470 let version = lite::Version::Lite05;
4471
4472 let mut origin = track_producer("test", None);
4474 let mut origin_dg = origin.subscribe(None);
4475 let ts = Timestamp::from_millis(7).unwrap();
4476 let seq = origin.append_datagram(ts, &b"payload"[..]).unwrap();
4477
4478 let d = recv_datagram(&mut origin_dg);
4479 let body = lite::Datagram {
4480 subscribe: 5,
4481 sequence: d.sequence,
4482 timestamp: d.timestamp.value(),
4483 payload: d.payload.clone(),
4484 }
4485 .encode_bytes(version)
4486 .unwrap();
4487
4488 let mut slice = &body[..];
4490 let wire = lite::Datagram::decode(&mut slice, version).unwrap();
4491 let mut downstream = track_producer("test", None);
4492 let mut downstream_dg = downstream.subscribe(None);
4493 downstream
4494 .insert_datagram(
4495 wire.sequence,
4496 Timestamp::new(wire.timestamp, Timescale::MILLI).unwrap(),
4497 wire.payload,
4498 )
4499 .unwrap();
4500
4501 let got = recv_datagram(&mut downstream_dg);
4502 assert_eq!(got.sequence, seq);
4503 assert_eq!(got.timestamp, ts);
4504 assert_eq!(&got.payload[..], b"payload");
4505 }
4506
4507 #[tokio::test]
4508 async fn evict_expired_groups() {
4509 let producer = track_producer("test", None);
4510
4511 producer.append_group().unwrap(); producer.append_group().unwrap(); producer.append_group().unwrap(); {
4517 let state = producer.state.read();
4518 assert_eq!(live_groups(&state), 3);
4519 assert_eq!(state.offset, 0);
4520 }
4521
4522 crate::model::clock::advance(cache::DEFAULT_EXPIRY + Duration::from_secs(1));
4524
4525 producer.append_group().unwrap(); {
4532 let state = producer.state.read();
4533 assert_eq!(live_groups(&state), 1);
4534 assert_eq!(first_live_sequence(&state), 3);
4535 assert_eq!(state.offset, 3);
4536 assert!(!state.lookup.contains_key(&0));
4537 assert!(!state.lookup.contains_key(&1));
4538 assert!(!state.lookup.contains_key(&2));
4539 assert!(state.lookup.contains_key(&3));
4540 }
4541 }
4542
4543 #[tokio::test]
4547 async fn aging_out_a_finished_group_keeps_the_clean_end() {
4548 let producer = track_producer("test", None);
4549 let mut group = producer.create_group(group::Info { sequence: 0 }).unwrap();
4550 let mut consumer = group.consume();
4551
4552 group
4553 .write_frame(Timestamp::from_millis(0).unwrap(), b"hello".as_slice())
4554 .unwrap();
4555 assert_eq!(consumer.next_frame().await.unwrap().unwrap().size, 5);
4556
4557 crate::model::clock::advance(cache::DEFAULT_EXPIRY * 2);
4559 group.finish().unwrap();
4560 let _next = producer.create_group(group::Info { sequence: 1 }).unwrap();
4561
4562 assert!(consumer.next_frame().await.unwrap().is_none());
4563 }
4564
4565 #[tokio::test]
4569 async fn active_reader_survives_expiry() {
4570 let producer = track_producer("test", None);
4571 let mut subscriber = producer.subscribe(None);
4572
4573 let mut group = producer.create_group(0u64.into()).unwrap();
4575 for _ in 0..10 {
4576 group.write_frame(Timestamp::ZERO, b"x".as_slice()).unwrap();
4577 }
4578 group.finish().unwrap();
4579 let mut reading = subscriber.assert_group();
4580
4581 producer.create_group(1u64.into()).unwrap().finish().unwrap();
4583
4584 for seq in 2..12u64 {
4587 crate::model::clock::advance(cache::DEFAULT_EXPIRY / 2);
4588 let frame = reading.next_frame().await;
4589 assert!(
4590 matches!(frame, Ok(Some(_))),
4591 "an actively-read group must not expire mid-read (step {seq})"
4592 );
4593 producer.create_group(seq.into()).unwrap().finish().unwrap();
4594 }
4595
4596 let state = producer.state.read();
4597 assert!(state.lookup.contains_key(&0), "the read group survived");
4598 assert!(!state.lookup.contains_key(&1), "the unread group still expired");
4599 }
4600
4601 #[tokio::test]
4606 async fn slow_prefetch_reader_survives_expiry() {
4607 let producer = track_producer("test", None);
4608 let mut subscriber = producer.subscribe(None);
4609
4610 let mut group = producer.create_group(0u64.into()).unwrap();
4611 for _ in 0..20 {
4612 group.write_frame(Timestamp::ZERO, b"x".as_slice()).unwrap();
4613 }
4614 group.finish().unwrap();
4615 let mut reading = subscriber.assert_group();
4616
4617 for seq in 1..20u64 {
4620 crate::model::clock::advance(cache::DEFAULT_EXPIRY / 2);
4621 let frame = reading.read_frame().await;
4622 assert!(
4623 matches!(frame, Ok(Some(_))),
4624 "a slow prefetch reader must not expire mid-read (step {seq})"
4625 );
4626 producer.create_group(seq.into()).unwrap().finish().unwrap();
4627 }
4628 }
4629
4630 #[tokio::test]
4634 async fn delivery_restarts_the_expiry_clock() {
4635 let producer = track_producer("test", None);
4636 let mut subscriber = producer.subscribe(replay());
4637
4638 let mut group = producer.create_group(0u64.into()).unwrap();
4639 group.write_frame(Timestamp::ZERO, b"x".as_slice()).unwrap();
4640 group.finish().unwrap();
4641 producer.create_group(1u64.into()).unwrap().finish().unwrap();
4643
4644 crate::model::clock::advance(cache::DEFAULT_EXPIRY - Duration::from_secs(1));
4646 let mut reading = subscriber.assert_group();
4647
4648 crate::model::clock::advance(cache::DEFAULT_EXPIRY - Duration::from_secs(1));
4651 producer.create_group(2u64.into()).unwrap().finish().unwrap();
4652
4653 let frame = reading.read_frame().await.unwrap();
4654 assert!(frame.is_some(), "a just-delivered group must not expire unread");
4655 }
4656
4657 #[tokio::test]
4661 async fn streaming_frame_writes_keep_the_group_alive() {
4662 let producer = track_producer("test", None);
4663 let mut straggler = producer.create_group(0u64.into()).unwrap();
4664 producer.create_group(1u64.into()).unwrap().finish().unwrap();
4666
4667 let mut frame = straggler
4668 .create_frame(frame::Info {
4669 size: 10,
4670 timestamp: Timestamp::ZERO,
4671 })
4672 .unwrap();
4673 for seq in 2..12u64 {
4676 crate::model::clock::advance(cache::DEFAULT_EXPIRY / 2);
4677 frame.write(bytes::Bytes::from_static(b"x")).unwrap();
4678 producer.create_group(seq.into()).unwrap().finish().unwrap();
4679 }
4680 frame.finish().unwrap();
4681 straggler.finish().unwrap();
4682
4683 let state = producer.state.read();
4684 assert!(
4685 state.lookup.contains_key(&0),
4686 "a group streaming a frame survives expiry"
4687 );
4688 }
4689
4690 #[tokio::test]
4695 async fn coalesced_frame_completion_keeps_the_group_alive() {
4696 let producer = track_producer("test", None);
4697 let mut straggler = producer.create_group(0u64.into()).unwrap();
4698 producer.create_group(1u64.into()).unwrap().finish().unwrap();
4700
4701 let mut frame = straggler
4702 .create_frame_owned(frame::Info {
4703 size: 3,
4704 timestamp: Timestamp::ZERO,
4705 })
4706 .unwrap();
4707
4708 crate::model::clock::advance(cache::DEFAULT_EXPIRY + Duration::from_secs(1));
4712 frame.write(bytes::Bytes::from_static(b"abc")).unwrap();
4713 frame.finish().unwrap();
4714 straggler.finish().unwrap();
4715
4716 producer.create_group(2u64.into()).unwrap().finish().unwrap();
4718
4719 let state = producer.state.read();
4720 assert!(
4721 state.lookup.contains_key(&0),
4722 "a group whose frame just completed must not expire"
4723 );
4724 }
4725
4726 #[tokio::test]
4729 async fn parked_reoffer_restarts_the_expiry_clock() {
4730 let producer = track_producer("test", None);
4731 let mut subscriber = producer.subscribe(None);
4732 subscriber.set_groups(..1);
4733
4734 for seq in 0..2u64 {
4735 let mut group = producer.create_group(seq.into()).unwrap();
4736 group.write_frame(Timestamp::ZERO, b"x".as_slice()).unwrap();
4737 group.finish().unwrap();
4738 }
4739
4740 assert_eq!(subscriber.assert_group().sequence, 0);
4742 subscriber.assert_no_group();
4743
4744 crate::model::clock::advance(cache::DEFAULT_EXPIRY - Duration::from_secs(1));
4746 subscriber.set_groups(..2);
4747 let mut reading = subscriber.assert_group();
4748 assert_eq!(reading.sequence, 1);
4749
4750 crate::model::clock::advance(cache::DEFAULT_EXPIRY - Duration::from_secs(1));
4753 producer.create_group(2u64.into()).unwrap().finish().unwrap();
4754
4755 let frame = reading.read_frame().await.unwrap();
4756 assert!(frame.is_some(), "a just-re-offered group must not expire unread");
4757 }
4758
4759 #[tokio::test]
4760 async fn evict_keeps_max_sequence() {
4761 let producer = track_producer("test", None);
4762 producer.append_group().unwrap(); crate::model::clock::advance(cache::DEFAULT_EXPIRY + Duration::from_secs(1));
4766
4767 producer.append_group().unwrap(); {
4771 let state = producer.state.read();
4772 assert_eq!(live_groups(&state), 1);
4773 assert_eq!(first_live_sequence(&state), 1);
4774 assert_eq!(state.offset, 1);
4775 }
4776 }
4777
4778 #[tokio::test]
4779 async fn no_eviction_when_fresh() {
4780 let producer = track_producer("test", None);
4781 producer.append_group().unwrap(); producer.append_group().unwrap(); producer.append_group().unwrap(); {
4786 let state = producer.state.read();
4787 assert_eq!(live_groups(&state), 3);
4788 assert_eq!(state.offset, 0);
4789 }
4790 }
4791
4792 #[tokio::test]
4793 async fn consumer_skips_evicted_groups() {
4794 let producer = track_producer("test", None);
4795 producer.append_group().unwrap(); let mut consumer = producer.subscribe(None);
4798
4799 crate::model::clock::advance(cache::DEFAULT_EXPIRY + Duration::from_secs(1));
4800 producer.append_group().unwrap(); let group = consumer.assert_group();
4804 assert_eq!(group.sequence, 1);
4805 }
4806
4807 fn track_producer_expiring(name: impl Into<Arc<str>>, expiry: impl Into<Option<Duration>>) -> Producer {
4809 track_producer_pooled(name, cache::Pool::new(cache::Config::default().with_expiry(expiry)))
4810 }
4811
4812 fn track_producer_pooled(name: impl Into<Arc<str>>, pool: cache::Pool) -> Producer {
4814 Producer::new(
4815 Arc::new(broadcast::Info {
4816 pool,
4817 ..Default::default()
4818 }),
4819 name,
4820 None,
4821 )
4822 }
4823
4824 #[tokio::test]
4828 async fn pool_sweep_expires_without_a_write() {
4829 let pool = cache::Pool::new(cache::Config::default().with_expiry(Duration::from_secs(1)));
4830 let producer = track_producer_pooled("test", pool.clone());
4831 let mut stalled = producer.append_group().unwrap(); stalled.write_frame(Timestamp::ZERO, b"x".as_slice()).unwrap();
4833 producer.append_group().unwrap(); crate::model::clock::advance(Duration::from_secs(2));
4836 pool.sweep();
4837
4838 assert!(
4839 !producer.state.read().lookup.contains_key(&0),
4840 "the sweep reclaimed an idle open group with no write behind it"
4841 );
4842 }
4843
4844 #[test]
4845 fn cache_gc_dates_activity_before_expiring_it() {
4846 let pool = cache::Pool::new(cache::Config::default().with_expiry(Duration::from_secs(1)));
4847 let now = crate::model::clock::now();
4848 assert_eq!(pool.gc(now), Some(now + Duration::from_millis(500)));
4849 let producer = track_producer_pooled("test", pool.clone());
4850 let mut group = producer.append_group().unwrap();
4851 group.write_frame(Timestamp::ZERO, b"first".as_slice()).unwrap();
4852 producer.append_group().unwrap();
4853 let later = now + Duration::from_secs(60);
4855 pool.gc(later);
4856 assert!(producer.state.read().lookup.contains_key(&0));
4857 pool.gc(later + Duration::from_secs(2));
4858 assert!(!producer.state.read().lookup.contains_key(&0));
4859 }
4860
4861 #[test]
4862 fn cache_gc_reaches_old_entries_behind_a_fresh_front() {
4863 let expiry = Duration::from_secs(1);
4864 let pool = cache::Pool::new(cache::Config::default().with_expiry(expiry));
4865 let producer = track_producer_pooled("test", pool.clone());
4866 let now = crate::model::clock::now();
4867 let groups: Vec<_> = (0..EVICT_SCAN * 3).map(|_| producer.append_group().unwrap()).collect();
4868 producer.append_group().unwrap();
4869 pool.gc(now);
4870 for group in &groups[..EVICT_SCAN * 2] {
4872 group.cache_refresh();
4873 }
4874 pool.gc(now + expiry * 2);
4875 let state = producer.state.read();
4876 for sequence in 0..EVICT_SCAN * 2 {
4877 assert!(state.lookup.contains_key(&(sequence as u64)), "fresh front survives");
4878 }
4879 for sequence in EVICT_SCAN * 2..EVICT_SCAN * 3 {
4880 assert!(!state.lookup.contains_key(&(sequence as u64)), "old tail is reclaimed");
4881 }
4882 }
4883
4884 #[tokio::test]
4888 async fn pool_sweep_drains_a_deep_backlog() {
4889 let pool = cache::Pool::new(cache::Config::default().with_expiry(Duration::from_secs(1)));
4890 let producer = track_producer_pooled("test", pool.clone());
4891
4892 let backlog = 4 * EVICT_SCAN;
4894 for _ in 0..backlog {
4895 let mut group = producer.append_group().unwrap();
4896 group.write_frame(Timestamp::ZERO, b"x".as_slice()).unwrap();
4897 }
4898 producer.append_group().unwrap(); crate::model::clock::advance(Duration::from_secs(2));
4901 pool.sweep();
4902
4903 let state = producer.state.read();
4904 let stale = (0..backlog as u64).filter(|seq| state.lookup.contains_key(seq)).count();
4905 assert_eq!(stale, 0, "one sweep reclaimed the whole idle backlog");
4906 }
4907
4908 #[tokio::test]
4909 async fn pool_expiry_controls_eviction() {
4910 let producer = track_producer_expiring("test", Duration::from_secs(1));
4912 producer.append_group().unwrap(); crate::model::clock::advance(Duration::from_secs(2));
4916 producer.append_group().unwrap(); let state = producer.state.read();
4920 assert_eq!(live_groups(&state), 1);
4921 assert_eq!(first_live_sequence(&state), 1);
4922 }
4923
4924 #[tokio::test]
4925 async fn small_frame_write_expires_idle_siblings() {
4926 let producer = track_producer_expiring("test", Duration::from_secs(1));
4927 producer.append_group().unwrap().finish().unwrap(); let mut live = producer.append_group().unwrap(); crate::model::clock::advance(Duration::from_secs(2));
4931 live.write_frame(Timestamp::ZERO, b"x".as_slice()).unwrap();
4932
4933 let expired = !producer.state.read().lookup.contains_key(&0);
4934 assert!(expired, "a small frame write runs expiry");
4935 }
4936
4937 #[tokio::test]
4938 async fn fresh_expiry_scan_does_not_wake_track_consumers() {
4939 use std::sync::atomic::{AtomicBool, Ordering};
4940
4941 let producer = track_producer_expiring("test", cache::DEFAULT_EXPIRY);
4942 producer.append_group().unwrap().finish().unwrap();
4943 let mut live = producer.append_group().unwrap();
4944 let mut consumer = producer.subscribe(None);
4945 assert_eq!(consumer.assert_group().sequence, 0);
4946 assert_eq!(consumer.assert_group().sequence, 1);
4947
4948 let woken = Arc::new(AtomicBool::new(false));
4949 let waiter = kio::Waiter::new(futures::task::waker(Arc::new(FlagWake(woken.clone()))));
4950 assert!(consumer.poll_recv_group(&waiter).is_pending());
4951
4952 live.write_frame(Timestamp::ZERO, b"x".as_slice()).unwrap();
4953 assert!(
4954 !woken.load(Ordering::SeqCst),
4955 "a no-op expiry scan must not wake track consumers"
4956 );
4957 }
4958
4959 #[tokio::test]
4960 async fn streaming_frame_write_expires_idle_siblings() {
4961 let producer = track_producer_expiring("test", Duration::from_secs(1));
4962 producer.append_group().unwrap().finish().unwrap(); let mut live = producer.append_group().unwrap(); let mut frame = live
4965 .create_frame(frame::Info {
4966 size: 1,
4967 timestamp: Timestamp::ZERO,
4968 })
4969 .unwrap();
4970
4971 crate::model::clock::advance(Duration::from_secs(2));
4972 frame.write(b"x".as_slice()).unwrap();
4973
4974 let expired = !producer.state.read().lookup.contains_key(&0);
4975 assert!(expired, "a streamed chunk runs expiry");
4976 }
4977
4978 #[tokio::test]
4979 async fn appended_datagram_expires_idle_groups() {
4980 let mut producer = track_producer_expiring("test", Duration::from_secs(1));
4981 producer.append_group().unwrap().finish().unwrap(); producer.append_group().unwrap().finish().unwrap(); crate::model::clock::advance(Duration::from_secs(2));
4985 producer.append_datagram(Timestamp::ZERO, b"x".as_slice()).unwrap();
4986
4987 let expired = !producer.state.read().lookup.contains_key(&0);
4988 assert!(expired, "an appended datagram runs expiry");
4989 }
4990
4991 #[tokio::test]
4992 async fn forwarded_datagram_expires_idle_groups() {
4993 let mut producer = track_producer_expiring("test", Duration::from_secs(1));
4994 producer.append_group().unwrap().finish().unwrap(); producer.append_group().unwrap().finish().unwrap(); crate::model::clock::advance(Duration::from_secs(2));
4998 producer
4999 .insert_datagram(2, Timestamp::ZERO, bytes::Bytes::from_static(b"x"))
5000 .unwrap();
5001
5002 let expired = !producer.state.read().lookup.contains_key(&0);
5003 assert!(expired, "a forwarded datagram runs expiry");
5004 }
5005
5006 #[tokio::test]
5010 async fn max_age_does_not_drive_wall_eviction() {
5011 let producer = track_producer("test", Info::default().with_max_age(Duration::from_secs(1)));
5012 producer.append_group().unwrap(); crate::model::clock::advance(Duration::from_secs(10));
5016 producer.append_group().unwrap(); let state = producer.state.read();
5019 assert_eq!(live_groups(&state), 2, "max_age is media time, not a wall clock");
5020 }
5021
5022 #[tokio::test]
5024 async fn disabled_pool_expiry_never_reclaims() {
5025 let producer = track_producer_expiring("test", None);
5026 producer.append_group().unwrap(); crate::model::clock::advance(Duration::from_secs(3600));
5029 producer.append_group().unwrap(); let state = producer.state.read();
5032 assert_eq!(live_groups(&state), 2);
5033 }
5034
5035 #[test]
5036 fn max_age_clamped_to_cache() {
5037 let producer = track_producer("test", Info::default().with_max_age(Duration::from_secs(2)));
5038
5039 let mut subscriber = producer.subscribe(Subscription::default().with_max_age(Duration::from_secs(10)));
5043 assert_eq!(subscriber.subscription().max_age, Duration::from_secs(10));
5044 assert_eq!(producer.subscription().unwrap().max_age, Duration::from_secs(2));
5045
5046 subscriber
5048 .update(Subscription::default().with_max_age(Duration::from_millis(500)))
5049 .unwrap();
5050 assert_eq!(producer.subscription().unwrap().max_age, Duration::from_millis(500));
5051
5052 subscriber
5053 .update(Subscription::default().with_max_age(Duration::ZERO))
5054 .unwrap();
5055 assert_eq!(producer.subscription().unwrap().max_age, Duration::ZERO);
5056 }
5057
5058 fn track_producer_capped(name: impl Into<Arc<str>>, info: Info, cap: Duration) -> Producer {
5061 Producer::new(
5062 Arc::new(broadcast::Info {
5063 cache_duration: cap,
5064 ..Default::default()
5065 }),
5066 name,
5067 info,
5068 )
5069 }
5070
5071 #[test]
5072 fn origin_cache_duration_clamps_max_age() {
5073 let capped = track_producer_capped(
5076 "test",
5077 Info::default().with_max_age(Duration::from_secs(60)),
5078 Duration::from_secs(1),
5079 );
5080 assert_eq!(capped.state.read().max_age_bound(), Some(Duration::from_secs(1)));
5081 assert_eq!(capped.subscribe(None).info().max_age, Duration::from_secs(1));
5082
5083 let under = track_producer_capped(
5084 "test",
5085 Info::default().with_max_age(Duration::from_millis(500)),
5086 Duration::from_secs(1),
5087 );
5088 assert_eq!(under.state.read().max_age_bound(), Some(Duration::from_millis(500)));
5089 }
5090
5091 #[tokio::test]
5094 async fn origin_cache_duration_does_not_wall_evict() {
5095 let producer = track_producer_capped(
5096 "test",
5097 Info::default().with_max_age(Duration::from_secs(60)),
5098 Duration::from_secs(1),
5099 );
5100 producer.append_group().unwrap(); crate::model::clock::advance(Duration::from_secs(2));
5104 producer.append_group().unwrap(); let state = producer.state.read();
5107 assert_eq!(live_groups(&state), 2);
5108 }
5109
5110 #[test]
5111 fn max_age_clamped_via_every_update_path() {
5112 let producer = track_producer("test", Info::default().with_max_age(Duration::from_secs(2)));
5113 let over = Subscription::default().with_max_age(Duration::from_secs(10));
5114
5115 let mut subscriber = producer.subscribe(over.clone());
5118 assert_eq!(producer.subscription().unwrap().max_age, Duration::from_secs(2));
5119
5120 subscriber.control().update(over.clone()).unwrap();
5121 assert_eq!(producer.subscription().unwrap().max_age, Duration::from_secs(2));
5122
5123 subscriber.update(over).unwrap();
5124 assert_eq!(producer.subscription().unwrap().max_age, Duration::from_secs(2));
5125 }
5126
5127 #[test]
5128 fn max_age_aggregate_clamps_across_subscribers() {
5129 let producer = track_producer("test", Info::default().with_max_age(Duration::from_secs(2)));
5130
5131 let _a = producer.subscribe(Subscription::default().with_max_age(Duration::from_millis(500)));
5134 let _b = producer.subscribe(Subscription::default().with_max_age(Duration::from_secs(10)));
5135
5136 assert_eq!(producer.subscription().unwrap().max_age, Duration::from_secs(2));
5137 }
5138
5139 fn append_at(producer: &mut Producer, millis: u64) -> u64 {
5142 let mut group = producer.append_group().unwrap();
5143 group
5144 .write_frame(Timestamp::from_millis(millis).unwrap(), bytes::Bytes::from_static(b"x"))
5145 .unwrap();
5146 group.finish().unwrap();
5147 group.sequence
5148 }
5149
5150 fn drain(subscriber: &mut Subscriber) -> Vec<u64> {
5152 let mut sequences = Vec::new();
5153 while let Some(Ok(Some(group))) = subscriber.recv_group().now_or_never() {
5154 sequences.push(group.sequence);
5155 }
5156 sequences
5157 }
5158
5159 #[test]
5160 fn real_time_skips_a_backlog_to_the_live_edge() {
5161 let mut producer = track_producer("test", None);
5162 for second in 0..5 {
5163 append_at(&mut producer, second * 1000);
5164 }
5165
5166 let mut subscriber = producer.subscribe(None);
5170 assert_eq!(drain(&mut subscriber), vec![4]);
5171
5172 append_at(&mut producer, 5000);
5174 assert_eq!(drain(&mut subscriber), vec![5]);
5175 }
5176
5177 #[test]
5178 fn real_time_skips_a_backlog_after_catching_up() {
5179 let mut producer = track_producer("test", None);
5180 append_at(&mut producer, 0);
5181 let mut subscriber = producer.subscribe(None);
5182 assert_eq!(drain(&mut subscriber), vec![0]);
5183
5184 for second in 1..6 {
5187 append_at(&mut producer, second * 1000);
5188 }
5189 assert_eq!(drain(&mut subscriber), vec![5]);
5190 }
5191
5192 #[test]
5193 fn a_newer_edge_changes_an_active_catch_up() {
5194 let mut producer = track_producer("test", None);
5195 for second in 0..5 {
5196 append_at(&mut producer, second * 1000);
5197 }
5198 let mut subscriber = producer.subscribe(Subscription::default().with_max_age(Duration::from_secs(1)));
5199 assert_eq!(
5200 subscriber
5201 .recv_group()
5202 .now_or_never()
5203 .unwrap()
5204 .unwrap()
5205 .unwrap()
5206 .sequence,
5207 3
5208 );
5209
5210 append_at(&mut producer, 5000);
5213 append_at(&mut producer, 10000);
5214 assert_eq!(drain(&mut subscriber), vec![5, 6]);
5215 }
5216
5217 #[test]
5218 fn a_growing_edge_changes_an_active_catch_up() {
5219 let mut producer = track_producer("test", None);
5220 append_at(&mut producer, 0);
5221 append_at(&mut producer, 1000);
5222 let mut edge = producer.append_group().unwrap();
5223 edge.write_frame(Timestamp::from_millis(2000).unwrap(), bytes::Bytes::from_static(b"a"))
5224 .unwrap();
5225
5226 let mut subscriber = producer.subscribe(Subscription::default().with_max_age(Duration::from_secs(2)));
5227 assert_eq!(
5228 subscriber
5229 .recv_group()
5230 .now_or_never()
5231 .unwrap()
5232 .unwrap()
5233 .unwrap()
5234 .sequence,
5235 0
5236 );
5237
5238 edge.write_frame(Timestamp::from_millis(5000).unwrap(), bytes::Bytes::from_static(b"b"))
5241 .unwrap();
5242 assert_eq!(drain(&mut subscriber), vec![2]);
5243 }
5244
5245 #[test]
5246 fn a_budget_admits_groups_within_it() {
5247 let mut producer = track_producer("test", None);
5248 for second in 0..5 {
5249 append_at(&mut producer, second * 1000);
5250 }
5251
5252 let mut subscriber = producer.subscribe(Subscription::default().with_max_age(Duration::from_secs(2)));
5257 assert_eq!(drain(&mut subscriber), vec![2, 3, 4]);
5258 }
5259
5260 #[test]
5261 fn a_budget_reaches_back_over_the_cache() {
5262 let mut producer = track_producer("test", None);
5263 for second in 0..5 {
5264 append_at(&mut producer, second * 1000);
5265 }
5266
5267 let mut live = producer.subscribe(None);
5270 assert_eq!(drain(&mut live), vec![4]);
5271
5272 let budget = Subscription::default().with_max_age(Duration::from_secs(2));
5277 let mut subscriber = producer.subscribe(budget);
5278 assert_eq!(drain(&mut subscriber), vec![2, 3, 4]);
5279 }
5280
5281 #[test]
5282 fn a_named_start_is_a_floor_not_a_request() {
5283 let mut producer = track_producer("test", None);
5284 for second in 0..5 {
5285 append_at(&mut producer, second * 1000);
5286 }
5287
5288 let named = Subscription::default().with_start(Position::group(1));
5292 let mut subscriber = producer.subscribe(named);
5293 assert_eq!(drain(&mut subscriber), vec![4]);
5294
5295 let floored = Subscription::default()
5297 .with_start(Position::group(3))
5298 .with_max_age(Duration::from_secs(10));
5299 let mut subscriber = producer.subscribe(floored);
5300 assert_eq!(drain(&mut subscriber), vec![3, 4]);
5301
5302 let slack = Subscription::default()
5304 .with_start(Position::group(1))
5305 .with_max_age(Duration::from_secs(2));
5306 let mut subscriber = producer.subscribe(slack);
5307 assert_eq!(drain(&mut subscriber), vec![2, 3, 4]);
5308 }
5309
5310 #[test]
5311 fn a_floor_above_the_live_edge_waits_there() {
5312 let mut producer = track_producer("test", None);
5313 for second in 0..3 {
5314 append_at(&mut producer, second * 1000);
5315 }
5316
5317 let resumed = Subscription::default()
5320 .with_start(Position::group(7))
5321 .with_max_age(Duration::from_secs(10));
5322 let mut subscriber = producer.subscribe(resumed);
5323 assert_eq!(drain(&mut subscriber), Vec::<u64>::new());
5324 append_at(&mut producer, 3000); assert_eq!(drain(&mut subscriber), Vec::<u64>::new());
5326 for second in 4..8 {
5327 append_at(&mut producer, second * 1000);
5328 }
5329 assert_eq!(drain(&mut subscriber), vec![7]);
5330 }
5331
5332 #[test]
5333 fn a_late_lower_group_within_the_budget_is_delivered() {
5334 let producer = track_producer("test", None);
5335 for (sequence, millis) in [(5, 0), (6, 1000), (7, 2000)] {
5336 let mut group = producer.create_group(group::Info { sequence }).unwrap();
5337 group
5338 .write_frame(Timestamp::from_millis(millis).unwrap(), bytes::Bytes::from_static(b"x"))
5339 .unwrap();
5340 group.finish().unwrap();
5341 }
5342
5343 let mut subscriber = producer.subscribe(Subscription::default().with_max_age(Duration::from_secs(5)));
5344 assert_eq!(drain(&mut subscriber), vec![5, 6, 7]);
5345
5346 let mut late = producer.create_group(group::Info { sequence: 4 }).unwrap();
5350 late.write_frame(Timestamp::from_millis(500).unwrap(), bytes::Bytes::from_static(b"late"))
5351 .unwrap();
5352 late.finish().unwrap();
5353 assert_eq!(drain(&mut subscriber), vec![4]);
5354 }
5355
5356 #[test]
5357 fn drift_is_measured_in_presentation_time_not_arrival_time() {
5358 let mut producer = track_producer("test", None);
5359 for second in 0..4 {
5363 append_at(&mut producer, second * 1000);
5364 }
5365
5366 let mut subscriber = producer.subscribe(Subscription::default().with_max_age(Duration::from_millis(1500)));
5367 assert_eq!(drain(&mut subscriber), vec![1, 2, 3]);
5368 }
5369
5370 #[tokio::test]
5371 async fn a_stamped_successor_expires_an_unstamped_group() {
5372 let mut producer = track_producer("test", None);
5373 let mut subscriber = producer.subscribe(None);
5374 producer.append_group().unwrap(); append_at(&mut producer, 1000); assert_eq!(drain(&mut subscriber), vec![1]);
5381 }
5382
5383 #[tokio::test]
5384 async fn a_handed_out_group_expires_while_its_first_frame_is_stalled() {
5385 let mut producer = track_producer("test", None);
5386 let mut subscriber = producer.subscribe(None);
5387 producer.append_group().unwrap();
5388
5389 let mut stalled = subscriber.recv_group().await.unwrap().expect("stalled group");
5390 let pending = tokio::spawn(async move { stalled.read_frame().await });
5391 tokio::task::yield_now().await;
5392 assert!(
5393 !pending.is_finished(),
5394 "the empty live edge still waits for its first frame"
5395 );
5396
5397 crate::model::clock::advance(Duration::from_secs(1));
5398 append_at(&mut producer, 1000);
5399
5400 let result = pending.await.unwrap();
5405 assert!(matches!(result, Ok(None)), "the held group ends: {result:?}");
5406 }
5407
5408 #[tokio::test]
5413 async fn a_handed_out_group_wakes_when_a_newer_group_gets_its_first_timestamp() {
5414 let mut producer = track_producer("test", None);
5415 let mut subscriber = producer.subscribe(Subscription::default().with_max_age(Duration::from_millis(500)));
5416 let mut old = producer.append_group().unwrap();
5417 old.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"old"))
5418 .unwrap();
5419
5420 append_at(&mut producer, 1000);
5422
5423 let mut held = subscriber.recv_group().await.unwrap().expect("old group");
5424 assert!(held.read_frame().await.unwrap().is_some());
5425
5426 let mut live = producer.append_group().unwrap();
5427 let pending = tokio::spawn(async move { held.read_frame().await });
5428 tokio::task::yield_now().await;
5429 assert!(!pending.is_finished(), "the newer group has no timestamp yet");
5430
5431 live.write_frame(
5433 Timestamp::from_millis(2000).unwrap(),
5434 bytes::Bytes::from_static(b"live"),
5435 )
5436 .unwrap();
5437 tokio::task::yield_now().await;
5438
5439 assert!(pending.is_finished(), "the new presentation edge wakes the held reader");
5440 let result = pending.await.unwrap();
5442 assert!(matches!(result, Ok(None)), "the held group ends: {result:?}");
5443 }
5444
5445 #[tokio::test]
5456 async fn real_time_reads_a_live_stream_without_truncating_it() {
5457 let producer = track_producer("test", None);
5458 let mut subscriber = producer.subscribe(None);
5459
5460 let gop = |n: u64| {
5461 [
5462 Timestamp::from_millis(n * 2000).unwrap(),
5463 Timestamp::from_millis(n * 2000 + 1900).unwrap(),
5464 ]
5465 };
5466 let write = |group: &mut group::Producer, timestamp| {
5467 group.write_frame(timestamp, bytes::Bytes::from_static(b"x")).unwrap();
5468 };
5469
5470 let mut open = producer.append_group().unwrap();
5471 write(&mut open, gop(0)[0]);
5472 write(&mut open, gop(0)[1]);
5473 let mut reading = subscriber.recv_group().await.unwrap().expect("the live group");
5474
5475 let mut read = Vec::new();
5476 for n in 1..5u64 {
5477 let sequence = reading.sequence;
5478 let mut frames = 0;
5479 while let Some(res) = reading.read_frame().now_or_never() {
5480 match res.expect("no truncation while draining") {
5481 Some(_) => frames += 1,
5482 None => panic!("group {sequence} ended early"),
5483 }
5484 }
5485 read.push((sequence, frames));
5486
5487 let next = {
5488 let mut end = std::pin::pin!(reading.read_frame());
5491 assert!(futures::poll!(end.as_mut()).is_pending(), "parked on the FIN");
5492
5493 let mut opened = producer.append_group().unwrap();
5496 write(&mut opened, gop(n)[0]);
5497 let verdict = futures::poll!(end.as_mut());
5498
5499 open.finish().unwrap();
5500 let res = match verdict {
5501 Poll::Ready(res) => res,
5502 Poll::Pending => end.await,
5503 };
5504 assert!(
5505 matches!(res, Ok(None)),
5506 "group {sequence} ends at the boundary rather than failing: {res:?}"
5507 );
5508 opened
5509 };
5510 let mut next = next;
5511 write(&mut next, gop(n)[1]);
5512
5513 reading = subscriber.recv_group().await.unwrap().expect("the next live group");
5514 open = next;
5515 }
5516
5517 assert_eq!(read, vec![(0, 2), (1, 2), (2, 2), (3, 2)], "every frame of every group");
5518 }
5519
5520 #[tokio::test]
5530 async fn a_budget_is_measured_from_the_readers_position() {
5531 let producer = track_producer("test", None);
5532 let mut subscriber = producer.subscribe(Subscription::default().with_max_age(Duration::from_secs(1)));
5533
5534 let mut open = producer.append_group().unwrap();
5535 open.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"key"))
5536 .unwrap();
5537 open.write_frame(
5538 Timestamp::from_millis(1900).unwrap(),
5539 bytes::Bytes::from_static(b"tail"),
5540 )
5541 .unwrap();
5542
5543 let mut reading = subscriber.recv_group().await.unwrap().expect("the live group");
5544 assert!(reading.read_frame().await.unwrap().is_some());
5545 assert!(reading.read_frame().await.unwrap().is_some());
5546
5547 let mut end = std::pin::pin!(reading.read_frame());
5548 assert!(futures::poll!(end.as_mut()).is_pending(), "parked at 1900ms");
5549
5550 let mut next = producer.append_group().unwrap();
5552 next.write_frame(Timestamp::from_millis(2000).unwrap(), bytes::Bytes::from_static(b"key"))
5553 .unwrap();
5554 assert!(
5555 futures::poll!(end.as_mut()).is_pending(),
5556 "a reader inside its budget is not expired by the next group opening"
5557 );
5558
5559 open.write_frame(
5561 Timestamp::from_millis(1950).unwrap(),
5562 bytes::Bytes::from_static(b"late"),
5563 )
5564 .unwrap();
5565 let late = end.await.expect("the straggler is not truncated");
5566 assert_eq!(
5567 late.map(|frame| frame.timestamp),
5568 Some(Timestamp::from_millis(1950).unwrap())
5569 );
5570 }
5571
5572 #[tokio::test]
5576 async fn an_ended_group_stays_ended_when_probed_again() {
5577 let mut producer = track_producer("test", None);
5578 let mut subscriber = producer.subscribe(None);
5579 let mut open = producer.append_group().unwrap();
5580 open.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"a"))
5581 .unwrap();
5582
5583 let mut group = subscriber.recv_group().await.unwrap().expect("group");
5584 assert!(group.read_frame().await.unwrap().is_some());
5585
5586 let probes = tokio::spawn(async move {
5587 let first = group.read_frame().await;
5588 let second = group.read_frame().await;
5589 let finished = group.finished().await;
5590 (first, second, finished)
5591 });
5592 tokio::task::yield_now().await;
5593
5594 append_at(&mut producer, 1000);
5595
5596 let (first, second, finished) = probes.await.unwrap();
5597 assert!(matches!(first, Ok(None)), "the group ends: {first:?}");
5598 assert!(matches!(second, Ok(None)), "and stays ended: {second:?}");
5599 assert!(matches!(finished, Ok(1)), "reporting what it delivered: {finished:?}");
5600 }
5601
5602 #[tokio::test]
5603 async fn a_drained_group_finishes_cleanly_after_the_live_edge_advances() {
5604 let mut producer = track_producer("test", None);
5605 let mut subscriber = producer.subscribe(None);
5606 append_at(&mut producer, 0);
5607
5608 let mut group = subscriber.recv_group().await.unwrap().expect("first group");
5609 assert!(group.read_frame().await.unwrap().is_some());
5610
5611 append_at(&mut producer, 1000);
5612
5613 assert!(group.read_frame().await.unwrap().is_none());
5614 assert!(!group.latency_expired());
5615 }
5616
5617 #[tokio::test]
5618 async fn a_handed_out_partial_frame_expires_while_its_payload_is_stalled() {
5619 let mut producer = track_producer("test", None);
5620 let mut subscriber = producer.subscribe(None);
5621 let mut source = producer.append_group().unwrap();
5622 let mut writing = source
5623 .create_frame(frame::Info {
5624 size: 6,
5625 timestamp: Timestamp::ZERO,
5626 })
5627 .unwrap();
5628 writing.write(bytes::Bytes::from_static(b"old")).unwrap();
5629
5630 let mut group = subscriber.recv_group().await.unwrap().expect("partial group");
5631 let mut frame = group.next_frame().await.unwrap().expect("partial frame");
5632 assert_eq!(
5633 frame.read_chunk().await.unwrap(),
5634 Some(bytes::Bytes::from_static(b"old"))
5635 );
5636 let pending = tokio::spawn(async move { frame.read_chunk().await });
5637 tokio::task::yield_now().await;
5638 assert!(!pending.is_finished(), "the partial payload is still stalled");
5639
5640 crate::model::clock::advance(Duration::from_secs(1));
5641 append_at(&mut producer, 1000);
5642
5643 let result = pending.await.unwrap();
5644 assert!(
5645 matches!(result, Err(Error::Old)),
5646 "the in-flight frame expires: {result:?}"
5647 );
5648 writing.abort(Error::Cancel).unwrap();
5649 }
5650
5651 #[test]
5652 fn max_age_bounds_the_budget() {
5653 let mut producer = track_producer("test", Info::default().with_max_age(Duration::from_millis(500)));
5657 append_at(&mut producer, 0);
5658 append_at(&mut producer, 1000);
5659 append_at(&mut producer, 2000);
5660
5661 let mut subscriber = producer.subscribe(Subscription::default().with_max_age(Duration::from_secs(10)));
5664 assert_eq!(drain(&mut subscriber), vec![1, 2]);
5665 }
5666
5667 #[test]
5668 fn an_explicit_start_gets_no_exemption_from_the_budget() {
5669 let mut producer = track_producer("test", None);
5670 for second in 0..4 {
5671 append_at(&mut producer, second * 1000);
5672 }
5673
5674 let mut subscriber = producer.subscribe(Subscription::default().with_start(Position::group(0)));
5677 subscriber.start_at(0);
5678 assert_eq!(drain(&mut subscriber), vec![3]);
5679
5680 let mut patient = producer.subscribe(
5681 Subscription::default()
5682 .with_start(Position::group(0))
5683 .with_max_age(Duration::from_secs(10)),
5684 );
5685 patient.start_at(0);
5686 assert_eq!(drain(&mut patient), vec![0, 1, 2, 3]);
5687 }
5688
5689 #[test]
5693 fn a_long_group_is_not_stale_while_its_tail_reaches_the_edge() {
5694 let mut producer = track_producer("test", None);
5695
5696 let mut long = producer.append_group().unwrap();
5698 for ms in [0u64, 500, 1000, 1500, 2000] {
5699 long.write_frame(Timestamp::from_millis(ms).unwrap(), bytes::Bytes::from_static(b"x"))
5700 .unwrap();
5701 }
5702 long.finish().unwrap();
5703 append_at(&mut producer, 2000);
5704
5705 let mut sub = producer.subscribe(Subscription::default().with_max_age(Duration::from_millis(500)));
5707 assert_eq!(drain(&mut sub), vec![0, 1]);
5708 }
5709
5710 #[test]
5716 fn a_long_group_is_stale_once_its_successor_falls_behind() {
5717 let mut producer = track_producer("test", None);
5718
5719 let mut long = producer.append_group().unwrap();
5720 for ms in [0u64, 500, 1000] {
5721 long.write_frame(Timestamp::from_millis(ms).unwrap(), bytes::Bytes::from_static(b"x"))
5722 .unwrap();
5723 }
5724 long.finish().unwrap();
5725 append_at(&mut producer, 3000);
5728 append_at(&mut producer, 4000);
5729
5730 let mut sub = producer.subscribe(Subscription::default().with_max_age(Duration::from_millis(500)));
5731 assert_eq!(drain(&mut sub), vec![1, 2]);
5732 }
5733
5734 #[tokio::test]
5738 async fn an_unstamped_immediate_successor_leaves_reach_unbounded() {
5739 let mut producer = track_producer("test", None);
5740 let mut subscriber = producer.subscribe(None);
5741
5742 append_at(&mut producer, 0); producer.append_group().unwrap(); append_at(&mut producer, 10_000); assert_eq!(drain(&mut subscriber), vec![0, 2]);
5750 }
5751
5752 #[test]
5758 fn reach_follows_the_immediate_successor_not_a_later_rewind() {
5759 let mut producer = track_producer("test", None);
5760
5761 append_at(&mut producer, 0);
5763 append_at(&mut producer, 10_000);
5764 append_at(&mut producer, 1_000);
5765 append_at(&mut producer, 2_000);
5766
5767 let mut sub = producer.subscribe(Subscription::default().with_max_age(Duration::from_millis(500)));
5771 assert!(
5772 drain(&mut sub).contains(&0),
5773 "group 0 is bounded by its successor at 10s, not by a later rewind"
5774 );
5775 }
5776
5777 #[tokio::test]
5780 async fn ordered_carries_datagrams() {
5781 let mut producer = track_producer("test", None);
5782 let mut sub = producer.subscribe(None).ordered();
5783
5784 producer
5785 .insert_datagram(5, Timestamp::from_millis(5).unwrap(), bytes::Bytes::from_static(b"x"))
5786 .unwrap();
5787 producer.create_group(group::Info { sequence: 3 }).unwrap();
5788
5789 let datagram = sub
5790 .recv_datagram()
5791 .now_or_never()
5792 .expect("datagram would have blocked")
5793 .expect("would have errored")
5794 .expect("track was closed");
5795 assert_eq!(datagram.sequence, 5);
5796
5797 let group = sub
5799 .next_group()
5800 .now_or_never()
5801 .expect("group would have blocked")
5802 .expect("would have errored")
5803 .expect("track was closed");
5804 assert_eq!(group.sequence, 3);
5805 }
5806
5807 fn drain_ordered(subscriber: &mut Ordered) -> Vec<u64> {
5809 let mut sequences = Vec::new();
5810 while let Some(Ok(Some(group))) = subscriber.next_group().now_or_never() {
5811 sequences.push(group.sequence);
5812 }
5813 sequences
5814 }
5815
5816 #[test]
5820 fn next_group_sheds_a_stale_backlog() {
5821 let mut producer = track_producer("test", None);
5822 for second in 0..4 {
5823 append_at(&mut producer, second * 1000);
5824 }
5825
5826 let mut subscriber = producer.subscribe(None).ordered();
5827 assert_eq!(drain_ordered(&mut subscriber), vec![3]);
5828
5829 let mut arrival = producer.subscribe(None);
5831 assert_eq!(drain(&mut arrival), vec![3]);
5832 }
5833
5834 #[test]
5838 fn next_group_keeps_a_backlog_inside_the_budget() {
5839 let mut producer = track_producer("test", None);
5840 for second in 0..4 {
5841 append_at(&mut producer, second * 1000);
5842 }
5843
5844 let mut subscriber = producer
5847 .subscribe(Subscription::default().with_max_age(Duration::from_millis(1500)))
5848 .ordered();
5849 assert_eq!(drain_ordered(&mut subscriber), vec![1, 2, 3]);
5850
5851 let mut replay = producer.subscribe(replay()).ordered();
5853 assert_eq!(drain_ordered(&mut replay), vec![0, 1, 2, 3]);
5854 }
5855
5856 #[test]
5859 fn next_group_keeps_a_group_with_no_proven_reach() {
5860 let mut producer = track_producer("test", None);
5861 append_at(&mut producer, 0); producer.append_group().unwrap(); append_at(&mut producer, 10_000); let mut subscriber = producer.subscribe(None).ordered();
5868 assert_eq!(drain_ordered(&mut subscriber), vec![0, 2]);
5869 }
5870
5871 #[tokio::test]
5872 async fn real_time_skips_older_sequences_with_equal_ages() {
5873 let mut producer = track_producer("test", None);
5874 append_at(&mut producer, 0);
5875 append_at(&mut producer, 0);
5876
5877 let mut subscriber = producer.subscribe(None);
5878 assert_eq!(drain(&mut subscriber), vec![1]);
5879 }
5880
5881 #[test]
5882 fn fetch_ignores_the_budget() {
5883 let mut producer = track_producer("test", None);
5884 for second in 0..4 {
5885 append_at(&mut producer, second * 1000);
5886 }
5887
5888 let consumer = producer.consume();
5891 let group = consumer.fetch_group(0, None).now_or_never().unwrap().unwrap();
5892 assert_eq!(group.sequence, 0);
5893 }
5894
5895 #[tokio::test]
5898 async fn fetched_group_is_not_a_live_drift_edge() {
5899 let mut producer = track_producer("test", None);
5900 let dynamic = producer.dynamic();
5901 let consumer = producer.consume();
5902 append_at(&mut producer, 0);
5903
5904 let pending = consumer.fetch_group(100, None);
5905 let req = dynamic
5906 .requested_group()
5907 .now_or_never()
5908 .expect("fetch request is ready")
5909 .unwrap();
5910 let mut fetched = req.accept(None).unwrap();
5911 fetched
5912 .write_frame(
5913 Timestamp::from_millis(100_000).unwrap(),
5914 bytes::Bytes::from_static(b"fetched"),
5915 )
5916 .unwrap();
5917 fetched.finish().unwrap();
5918 pending.await.unwrap();
5919
5920 let mut groups = producer.subscribe(None);
5921 assert_eq!(groups.assert_group().sequence, 0);
5922 groups.assert_no_group();
5923 }
5924
5925 #[test]
5930 fn an_evicted_live_edge_convicts_nothing() {
5931 let mut producer = track_producer("test", None);
5932 append_at(&mut producer, 0);
5933 let edge = append_at(&mut producer, 30_000);
5934
5935 let state = producer.state.read();
5936 let drift = Drift {
5937 budget: Duration::ZERO,
5938 edge: state.live_edge(None),
5939 };
5940 assert!(
5941 state.is_stale(0, drift.edge.as_ref(), drift.budget),
5942 "stale against a live edge"
5943 );
5944 drop(state);
5945
5946 let slot = producer.modify().unwrap().lookup.remove(&edge).unwrap();
5948 let _ = slot.group.abort(Error::Evicted);
5949
5950 let state = producer.state.read();
5951 assert!(
5952 !state.is_stale(0, drift.edge.as_ref(), drift.budget),
5953 "a vanished edge is no reason to drop what is left"
5954 );
5955 }
5956
5957 #[tokio::test]
5958 async fn a_lower_sequence_is_never_the_live_edge() {
5959 let producer = track_producer("test", None);
5960 let mut straggler = producer.create_group(0u64.into()).unwrap();
5965 straggler
5966 .write_frame(Timestamp::from_millis(60_000).unwrap(), bytes::Bytes::from_static(b"x"))
5967 .unwrap();
5968 straggler.finish().unwrap();
5969
5970 let mut rewound = producer.create_group(1u64.into()).unwrap();
5971 rewound
5972 .write_frame(Timestamp::from_millis(0).unwrap(), bytes::Bytes::from_static(b"x"))
5973 .unwrap();
5974 rewound.finish().unwrap();
5975
5976 let mut subscriber = producer.subscribe(replay());
5977 assert_eq!(drain(&mut subscriber), vec![0, 1]);
5978 }
5979
5980 #[test]
5981 fn a_requested_end_does_not_cap_the_live_edge() {
5982 let mut producer = track_producer("test", None);
5983 for second in 0..4 {
5984 append_at(&mut producer, second * 1000);
5985 }
5986
5987 let mut subscriber = producer.subscribe(Subscription::default().with_end(Position::after_group(1)));
5992 assert_eq!(drain(&mut subscriber), vec![3]);
5993 }
5994
5995 #[test]
5996 fn a_capped_subscriber_measures_drift_against_its_cap() {
5997 let mut producer = track_producer("test", None);
5998 append_at(&mut producer, 0);
5999 append_at(&mut producer, 1000);
6000
6001 let mut subscriber = producer.subscribe(Subscription::default().with_end(Position::after_group(0)));
6006 subscriber.set_groups(..1);
6007 assert_eq!(drain(&mut subscriber), vec![0]);
6008
6009 subscriber.set_groups(..);
6012 assert_eq!(drain(&mut subscriber), vec![1]);
6013 }
6014
6015 #[test]
6016 fn subscriber_control_updates_while_read_future_is_pending() {
6017 let producer = track_producer("test", None);
6018 let mut subscriber = producer.subscribe(None);
6019 let control = subscriber.control();
6020
6021 let mut recv = Box::pin(subscriber.recv_group());
6022 assert!(recv.as_mut().now_or_never().is_none());
6023
6024 control.update(Subscription::default().with_priority(7)).unwrap();
6025
6026 let aggregate = producer.subscription().expect("expected an active subscription");
6027 assert_eq!(aggregate.priority, 7);
6028 }
6029
6030 #[test]
6031 fn dropped_subscriber_leaves_no_ghost_in_aggregate() {
6032 let mut producer = track_producer("test", None);
6037 let a = producer.subscribe(Subscription::default().with_priority(5));
6038
6039 let waiter = kio::Waiter::noop();
6041 assert!(
6042 matches!(producer.poll_subscription_changed(&waiter), Poll::Ready(Ok(Some(_)))),
6043 "one live subscriber should aggregate to Some",
6044 );
6045
6046 drop(a);
6048
6049 assert!(
6051 matches!(producer.poll_subscription_changed(&waiter), Poll::Ready(Ok(None))),
6052 "a dropped subscriber must not linger in the aggregate",
6053 );
6054
6055 assert!(
6057 producer.subscription().is_none(),
6058 "snapshot must exclude a dropped subscriber",
6059 );
6060 }
6061
6062 #[test]
6063 fn dropped_subscriber_wakes_the_aggregate() {
6064 use std::sync::atomic::{AtomicBool, Ordering};
6071
6072 let mut producer = track_producer("test", None);
6073 let a = producer.subscribe(Subscription::default().with_priority(5));
6074
6075 let woken = Arc::new(AtomicBool::new(false));
6076 let waiter = kio::Waiter::new(futures::task::waker(Arc::new(FlagWake(woken.clone()))));
6077
6078 assert!(matches!(
6080 producer.poll_subscription_changed(&waiter),
6081 Poll::Ready(Ok(Some(_)))
6082 ));
6083 assert!(
6084 producer.poll_subscription_changed(&waiter).is_pending(),
6085 "the aggregate is unchanged, so this poll must park",
6086 );
6087 assert!(!woken.load(Ordering::SeqCst), "nothing happened yet");
6088
6089 drop(a);
6090 assert!(
6091 woken.load(Ordering::SeqCst),
6092 "the last subscriber leaving must wake the aggregate watcher",
6093 );
6094 }
6095
6096 #[test]
6097 fn widest_subscriber_update_wakes_the_aggregate() {
6098 use std::sync::atomic::{AtomicBool, Ordering};
6104
6105 let mut producer = track_producer("test", None);
6106 let _narrow = producer.subscribe(Subscription::default().with_end(Position::after_group(3)));
6107 let mut wide = producer.subscribe(Subscription::default());
6108
6109 let woken = Arc::new(AtomicBool::new(false));
6110 let waiter = kio::Waiter::new(futures::task::waker(Arc::new(FlagWake(woken.clone()))));
6111
6112 assert!(matches!(
6113 producer.poll_subscription_changed(&waiter),
6114 Poll::Ready(Ok(Some(_)))
6115 ));
6116 assert!(producer.poll_subscription_changed(&waiter).is_pending());
6117 assert!(!woken.load(Ordering::SeqCst), "nothing happened yet");
6118
6119 wide.update(Subscription::default().with_end(Position::after_group(5)))
6120 .unwrap();
6121 assert!(
6122 woken.load(Ordering::SeqCst),
6123 "the widest subscriber changing must wake the aggregate watcher",
6124 );
6125 match producer.poll_subscription_changed(&waiter) {
6126 Poll::Ready(Ok(Some(sub))) => assert_eq!(sub.end, Position::after_group(5)),
6127 other => panic!("expected the narrowed aggregate, got {other:?}"),
6128 }
6129 }
6130
6131 struct FlagWake(Arc<std::sync::atomic::AtomicBool>);
6133
6134 impl futures::task::ArcWake for FlagWake {
6135 fn wake_by_ref(arc_self: &Arc<Self>) {
6136 arc_self.0.store(true, std::sync::atomic::Ordering::SeqCst);
6137 }
6138 }
6139
6140 #[tokio::test]
6141 async fn out_of_order_max_sequence_at_front() {
6142 let producer = track_producer("test", None);
6143
6144 producer.create_group(group::Info { sequence: 5 }).unwrap();
6146 producer.create_group(group::Info { sequence: 3 }).unwrap();
6147 producer.create_group(group::Info { sequence: 4 }).unwrap();
6148
6149 {
6151 let state = producer.state.read();
6152 assert_eq!(state.max_sequence, Some(5));
6153 }
6154
6155 crate::model::clock::advance(cache::DEFAULT_EXPIRY + Duration::from_secs(1));
6157
6158 producer.append_group().unwrap(); {
6164 let state = producer.state.read();
6165 assert_eq!(live_groups(&state), 1);
6166 assert_eq!(first_live_sequence(&state), 6);
6167 assert!(!state.lookup.contains_key(&3));
6168 assert!(!state.lookup.contains_key(&4));
6169 assert!(!state.lookup.contains_key(&5));
6170 assert!(state.lookup.contains_key(&6));
6171 }
6172 }
6173
6174 #[tokio::test]
6175 async fn max_sequence_at_front_blocks_trim() {
6176 let producer = track_producer("test", None);
6177
6178 producer.create_group(group::Info { sequence: 5 }).unwrap();
6180
6181 crate::model::clock::advance(cache::DEFAULT_EXPIRY + Duration::from_secs(1));
6182
6183 producer.create_group(group::Info { sequence: 3 }).unwrap();
6185
6186 {
6189 let state = producer.state.read();
6190 assert_eq!(live_groups(&state), 2);
6191 assert_eq!(state.offset, 0);
6192 }
6193
6194 crate::model::clock::advance(cache::DEFAULT_EXPIRY + Duration::from_secs(1));
6196
6197 producer.create_group(group::Info { sequence: 2 }).unwrap();
6199
6200 {
6205 let state = producer.state.read();
6206 assert_eq!(live_groups(&state), 2);
6207 assert_eq!(state.offset, 0);
6208 assert!(state.lookup.contains_key(&5));
6209 assert!(!state.lookup.contains_key(&3));
6210 assert!(state.lookup.contains_key(&2));
6211 }
6212
6213 let mut consumer = producer.subscribe(None);
6215 let group = consumer.assert_group();
6216 assert_eq!(group.sequence, 5);
6218 }
6219
6220 #[tokio::test]
6221 async fn abort_clears_cached_groups() {
6222 let producer = track_producer("test", None);
6223 producer.append_group().unwrap();
6224 producer.append_group().unwrap();
6225
6226 let mut consumer = producer.subscribe(None);
6228 assert_eq!(live_groups(&producer.state.read()), 2);
6229
6230 producer.clone().abort(Error::Cancel).unwrap();
6231
6232 {
6233 let state = producer.state.read();
6234 assert!(state.lookup.is_empty(), "cached groups should be dropped on abort");
6235 assert!(state.arrival.is_empty());
6236 assert!(state.evict.is_empty());
6237 }
6238
6239 let result = consumer.recv_group().now_or_never().expect("should not block");
6241 assert!(matches!(result, Err(Error::Cancel)));
6242 }
6243
6244 #[tokio::test]
6245 async fn drop_unfinished_clears_cached_groups() {
6246 let producer = track_producer("test", None);
6247 let writer = producer.clone();
6248 writer.append_group().unwrap();
6249
6250 let mut consumer = producer.subscribe(None);
6252 assert_eq!(live_groups(&producer.state.read()), 1);
6253
6254 drop(writer);
6256 drop(producer);
6257
6258 let result = consumer.recv_group().now_or_never().expect("should not block");
6259 assert!(matches!(result, Err(Error::Dropped)));
6260 }
6261
6262 #[tokio::test]
6263 async fn drop_after_abort_does_not_warn() {
6264 let warns = count_drop_warnings("track::Producer dropped without finish", || {
6267 let producer = track_producer("test", None);
6268 let keep = producer.clone();
6269 let writer = producer.clone();
6270 let group = writer.append_group().unwrap();
6271 group.finish().unwrap();
6272 let _consumer = producer.subscribe(None);
6273 writer.abort(Error::Cancel).unwrap();
6274 drop(keep);
6275 });
6276 assert_eq!(warns, 0, "abort-then-drop must not emit unfinished-producer WARN");
6277 }
6278
6279 #[tokio::test]
6280 async fn drop_unfinished_warns() {
6281 let warns = count_drop_warnings("track::Producer dropped without finish", || {
6282 let producer = track_producer("test", None);
6283 let writer = producer.clone();
6284 writer.append_group().unwrap();
6285 let _consumer = producer.subscribe(None);
6286 drop(writer);
6287 drop(producer);
6288 });
6289 assert!(warns >= 1, "unfinished drop must emit unfinished-producer WARN");
6290 }
6291
6292 #[tokio::test]
6293 async fn drop_finished_keeps_cached_groups() {
6294 let producer = track_producer("test", None);
6295 producer.append_group().unwrap();
6296 producer.finish().unwrap();
6297
6298 let mut consumer = producer.subscribe(None);
6299 drop(producer);
6300
6301 assert_eq!(consumer.assert_group().sequence, 0);
6303 let done = consumer.recv_group().now_or_never().expect("should not block").unwrap();
6304 assert!(done.is_none(), "consumer should drain then see clean finish");
6305 }
6306
6307 #[tokio::test]
6308 async fn cached_groups_preserve_arrival_order() {
6309 let producer = track_producer("test", None);
6310 producer.create_group(group::Info { sequence: 5 }).unwrap();
6311 producer.create_group(group::Info { sequence: 3 }).unwrap();
6312
6313 let groups = producer.consume().cached_groups();
6314 let sequences: Vec<u64> = groups.iter().map(|(group, _)| group.sequence).collect();
6315 assert_eq!(
6316 sequences,
6317 vec![5, 3],
6318 "warm snapshot must follow arrival, not sequence order"
6319 );
6320 }
6321
6322 #[test]
6323 fn append_finish_cannot_be_rewritten() {
6324 let producer = track_producer("test", None);
6325
6326 assert!(producer.finish().is_ok());
6328 assert!(producer.finish().is_err());
6329 assert!(producer.append_group().is_err());
6330 }
6331
6332 #[test]
6333 fn finish_after_groups() {
6334 let producer = track_producer("test", None);
6335
6336 producer.append_group().unwrap();
6337 assert!(producer.finish().is_ok());
6338 assert!(producer.finish().is_err());
6339 assert!(producer.append_group().is_err());
6340 }
6341
6342 #[test]
6343 fn finish_at_rejects_a_boundary_at_or_below_the_live_edge() {
6344 let mut producer = track_producer("test", None);
6345 producer.create_group(group::Info { sequence: 5 }).unwrap();
6346
6347 assert!(producer.finish_at(4).is_err());
6350 assert!(producer.finish_at(5).is_err());
6351 assert!(producer.finish_at(6).is_ok());
6352
6353 {
6354 let state = producer.state.read();
6355 assert_eq!(state.final_sequence, Some(6));
6356 }
6357
6358 assert!(producer.finish_at(6).is_err());
6360 assert!(producer.create_group(group::Info { sequence: 4 }).is_ok());
6361 assert!(producer.create_group(group::Info { sequence: 6 }).is_err());
6362 }
6363
6364 #[test]
6365 fn final_sequence_reports_the_declared_boundary() {
6366 let mut producer = track_producer("test", None);
6367 assert_eq!(producer.final_sequence(), None);
6368
6369 producer.create_group(group::Info { sequence: 5 }).unwrap();
6370 assert_eq!(producer.final_sequence(), None, "a group does not declare a boundary");
6371
6372 producer.finish_at(9).unwrap();
6373 assert_eq!(producer.final_sequence(), Some(9));
6374
6375 assert!(producer.finish().is_err());
6377 }
6378
6379 #[test]
6380 fn final_sequence_reports_the_live_edge_after_finish() {
6381 let producer = track_producer("test", None);
6382 producer.create_group(group::Info { sequence: 5 }).unwrap();
6383 producer.finish().unwrap();
6384 assert_eq!(producer.final_sequence(), Some(6));
6385 }
6386
6387 #[tokio::test]
6388 async fn finish_at_declares_a_future_boundary() {
6389 let mut producer = track_producer("test", None);
6390 producer.create_group(group::Info { sequence: 5 }).unwrap();
6391
6392 producer.finish_at(7).unwrap();
6394
6395 let mut consumer = producer.subscribe(None);
6396 assert_eq!(consumer.assert_group().sequence, 5);
6397
6398 let boundary = consumer
6401 .finished()
6402 .now_or_never()
6403 .expect("boundary is known immediately")
6404 .expect("would have errored");
6405 assert_eq!(boundary, 7);
6406 assert!(
6407 consumer.recv_group().now_or_never().is_none(),
6408 "should wait for the outstanding group"
6409 );
6410
6411 producer.create_group(group::Info { sequence: 6 }).unwrap();
6413 assert_eq!(consumer.assert_group().sequence, 6);
6414 let done = consumer
6415 .recv_group()
6416 .now_or_never()
6417 .expect("should not block")
6418 .expect("would have errored");
6419 assert!(done.is_none(), "track completes once the boundary is reached");
6420 }
6421
6422 #[tokio::test]
6423 async fn recv_group_finishes_without_waiting_for_gaps() {
6424 let producer = track_producer("test", None);
6425 producer.create_group(group::Info { sequence: 1 }).unwrap();
6426 producer.finish().unwrap();
6427
6428 let mut consumer = producer.subscribe(None);
6429 assert_eq!(consumer.assert_group().sequence, 1);
6430
6431 let done = consumer
6432 .recv_group()
6433 .now_or_never()
6434 .expect("should not block")
6435 .expect("would have errored");
6436 assert!(done.is_none(), "track should finish without waiting for gaps");
6437 }
6438
6439 #[tokio::test]
6440 async fn next_group_skips_late_arrivals() {
6441 let producer = track_producer("test", None);
6442 let mut consumer = producer.subscribe(None).ordered();
6443
6444 producer.create_group(group::Info { sequence: 5 }).unwrap();
6446 let group = consumer
6447 .next_group()
6448 .now_or_never()
6449 .expect("should not block")
6450 .expect("would have errored")
6451 .expect("track should not be closed");
6452 assert_eq!(group.sequence, 5);
6453
6454 producer.create_group(group::Info { sequence: 3 }).unwrap();
6456 producer.create_group(group::Info { sequence: 4 }).unwrap();
6458 producer.create_group(group::Info { sequence: 7 }).unwrap();
6460
6461 let group = consumer
6462 .next_group()
6463 .now_or_never()
6464 .expect("should not block")
6465 .expect("would have errored")
6466 .expect("track should not be closed");
6467 assert_eq!(group.sequence, 7);
6468
6469 assert!(
6471 consumer.next_group().now_or_never().is_none(),
6472 "should block waiting for a higher sequence"
6473 );
6474 }
6475
6476 #[tokio::test]
6477 async fn next_group_returns_arrivals_in_order() {
6478 let producer = track_producer("test", None);
6479 let mut consumer = producer.subscribe(replay()).ordered();
6480
6481 producer.create_group(group::Info { sequence: 3 }).unwrap();
6483 producer.create_group(group::Info { sequence: 5 }).unwrap();
6484
6485 let group = consumer
6486 .next_group()
6487 .now_or_never()
6488 .expect("should not block")
6489 .expect("would have errored")
6490 .expect("track should not be closed");
6491 assert_eq!(group.sequence, 3);
6492
6493 let group = consumer
6494 .next_group()
6495 .now_or_never()
6496 .expect("should not block")
6497 .expect("would have errored")
6498 .expect("track should not be closed");
6499 assert_eq!(group.sequence, 5);
6500 }
6501
6502 #[tokio::test]
6503 async fn ordered_and_arrival_cursors_are_independent() {
6504 let producer = track_producer("test", None);
6505 let mut ordered = producer.subscribe(replay()).ordered();
6506 let mut arrival = producer.subscribe(replay());
6507
6508 producer.create_group(group::Info { sequence: 5 }).unwrap();
6510 producer.create_group(group::Info { sequence: 3 }).unwrap();
6511
6512 let group = ordered
6515 .next_group()
6516 .now_or_never()
6517 .expect("should not block")
6518 .expect("would have errored")
6519 .expect("track should not be closed");
6520 assert_eq!(group.sequence, 3);
6521
6522 assert_eq!(arrival.assert_group().sequence, 5);
6524 }
6525
6526 #[tokio::test]
6527 async fn end_at_caps_next_group() {
6528 let producer = track_producer("test", None);
6529 let mut consumer = producer.subscribe(replay()).ordered();
6530
6531 for s in 0..6 {
6532 producer.create_group(group::Info { sequence: s }).unwrap();
6533 }
6534
6535 consumer.set_groups(..3);
6536
6537 assert_eq!(
6539 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6540 0
6541 );
6542 assert_eq!(
6543 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6544 1
6545 );
6546 assert_eq!(
6547 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6548 2
6549 );
6550
6551 assert!(
6553 consumer.next_group().now_or_never().is_none(),
6554 "capped consumer must block instead of returning out-of-range groups"
6555 );
6556 }
6557
6558 #[tokio::test]
6559 async fn end_at_release_drains_cached_groups() {
6560 let producer = track_producer("test", None);
6561 let mut consumer = producer.subscribe(replay()).ordered();
6562
6563 for s in 0..6 {
6564 producer.create_group(group::Info { sequence: s }).unwrap();
6565 }
6566
6567 consumer.set_groups(..2);
6568 assert_eq!(
6569 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6570 0
6571 );
6572 assert_eq!(
6573 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6574 1
6575 );
6576 assert!(consumer.next_group().now_or_never().is_none(), "capped at 2");
6577
6578 consumer.set_groups(..5);
6580 assert_eq!(
6581 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6582 2
6583 );
6584 assert_eq!(
6585 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6586 3
6587 );
6588 assert_eq!(
6589 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6590 4
6591 );
6592 assert!(consumer.next_group().now_or_never().is_none(), "capped at 5");
6593
6594 consumer.set_groups(..);
6596 assert_eq!(
6597 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6598 5
6599 );
6600 assert!(consumer.next_group().now_or_never().is_none(), "no more groups");
6601 }
6602
6603 #[tokio::test]
6604 async fn end_at_lower_than_cursor_parks_consumer() {
6605 let producer = track_producer("test", None);
6606 let mut consumer = producer.subscribe(replay()).ordered();
6607
6608 for s in 0..3 {
6609 producer.create_group(group::Info { sequence: s }).unwrap();
6610 }
6611
6612 assert_eq!(
6614 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6615 0
6616 );
6617 assert_eq!(
6618 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6619 1
6620 );
6621 assert_eq!(
6622 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6623 2
6624 );
6625
6626 consumer.set_groups(..2);
6628 producer.create_group(group::Info { sequence: 3 }).unwrap();
6629 producer.create_group(group::Info { sequence: 4 }).unwrap();
6630 assert!(
6631 consumer.next_group().now_or_never().is_none(),
6632 "cap is below cursor; nothing returnable until cap rises"
6633 );
6634
6635 consumer.set_groups(..);
6637 assert_eq!(
6638 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6639 3
6640 );
6641 assert_eq!(
6642 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6643 4
6644 );
6645 }
6646
6647 #[tokio::test]
6648 async fn end_at_toggling_around_late_arrivals() {
6649 let producer = track_producer("test", None);
6650 let mut consumer = producer.subscribe(replay()).ordered();
6651
6652 consumer.set_groups(..6);
6653
6654 producer.create_group(group::Info { sequence: 2 }).unwrap();
6656 producer.create_group(group::Info { sequence: 5 }).unwrap();
6657 producer.create_group(group::Info { sequence: 3 }).unwrap();
6658 producer.create_group(group::Info { sequence: 8 }).unwrap();
6660 producer.create_group(group::Info { sequence: 4 }).unwrap();
6661
6662 assert_eq!(
6664 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6665 2
6666 );
6667 assert_eq!(
6668 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6669 3
6670 );
6671 assert_eq!(
6672 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6673 4
6674 );
6675 assert_eq!(
6676 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6677 5
6678 );
6679 assert!(consumer.next_group().now_or_never().is_none());
6681
6682 consumer.set_groups(..11);
6684 assert_eq!(
6685 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6686 8
6687 );
6688 }
6689
6690 #[tokio::test]
6694 async fn end_at_parks_recv_group() {
6695 let producer = track_producer("test", None);
6696 let mut consumer = producer.subscribe(replay());
6697
6698 for s in 0..3 {
6699 producer.create_group(group::Info { sequence: s }).unwrap();
6700 }
6701
6702 consumer.set_groups(..2);
6703 assert_eq!(
6704 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6705 0
6706 );
6707 assert_eq!(
6708 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6709 1
6710 );
6711 assert!(consumer.recv_group().now_or_never().is_none(), "capped at 2");
6712
6713 producer.finish().unwrap();
6715 assert!(
6716 consumer.recv_group().now_or_never().is_none(),
6717 "still parked after finish"
6718 );
6719
6720 consumer.set_groups(..);
6721 assert_eq!(
6722 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6723 2
6724 );
6725 assert!(
6726 matches!(consumer.recv_group().now_or_never(), Some(Ok(None))),
6727 "finished once the parked group drains"
6728 );
6729 }
6730
6731 #[tokio::test]
6734 async fn recv_group_serves_arrivals_behind_the_cap() {
6735 let producer = track_producer("test", None);
6736 let mut consumer = producer.subscribe(replay());
6737
6738 consumer.set_groups(..2);
6739
6740 producer.create_group(group::Info { sequence: 2 }).unwrap();
6742 producer.create_group(group::Info { sequence: 0 }).unwrap();
6743 producer.create_group(group::Info { sequence: 1 }).unwrap();
6744
6745 assert_eq!(
6746 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6747 0
6748 );
6749 assert_eq!(
6750 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6751 1
6752 );
6753 assert!(consumer.recv_group().now_or_never().is_none(), "capped at 2");
6754
6755 consumer.set_groups(..3);
6756 assert_eq!(
6757 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6758 2
6759 );
6760 }
6761
6762 #[tokio::test]
6763 async fn group_ranges_preserve_the_floor_when_the_cap_changes() {
6764 let producer = track_producer("test", None);
6765 let mut consumer = producer.subscribe(None);
6766 consumer.set_groups(2..=2);
6767 producer.create_group(group::Info { sequence: 2 }).unwrap();
6768 assert_eq!(consumer.recv_group().await.unwrap().unwrap().sequence, 2);
6769 consumer.set_groups(..4);
6770 producer.create_group(group::Info { sequence: 1 }).unwrap();
6771 producer.create_group(group::Info { sequence: 3 }).unwrap();
6772 assert_eq!(consumer.recv_group().await.unwrap().unwrap().sequence, 3);
6773 consumer.set_groups(0..=4);
6774 producer.create_group(group::Info { sequence: 0 }).unwrap();
6775 producer.create_group(group::Info { sequence: 4 }).unwrap();
6776 assert_eq!(consumer.recv_group().await.unwrap().unwrap().sequence, 4);
6777 }
6778
6779 #[tokio::test]
6782 async fn start_at_drops_parked_recv_groups() {
6783 let producer = track_producer("test", None);
6784 let mut consumer = producer.subscribe(None);
6785
6786 consumer.set_groups(..1);
6787 producer.create_group(group::Info { sequence: 1 }).unwrap();
6788 assert!(
6789 consumer.recv_group().now_or_never().is_none(),
6790 "group 1 parked at the cap"
6791 );
6792
6793 consumer.start_at(2);
6794 consumer.set_groups(..);
6795 producer.create_group(group::Info { sequence: 2 }).unwrap();
6796 assert_eq!(
6797 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6798 2,
6799 "the overtaken parked group is dropped, not re-offered"
6800 );
6801 }
6802
6803 #[tokio::test]
6807 async fn evicted_parked_recv_groups_are_dropped() {
6808 let producer = track_producer("test", None);
6809 let mut consumer = producer.subscribe(None);
6810
6811 producer.create_group(group::Info { sequence: 0 }).unwrap();
6812 assert_eq!(
6813 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6814 0
6815 );
6816
6817 consumer.set_groups(..1);
6818 let straggler = producer.create_group(group::Info { sequence: 1 }).unwrap();
6819 assert!(
6820 consumer.recv_group().now_or_never().is_none(),
6821 "group 1 parked at the cap"
6822 );
6823
6824 straggler.abort(Error::Old).unwrap();
6826 producer.finish().unwrap();
6827
6828 consumer.set_groups(..);
6829 assert!(
6830 matches!(consumer.recv_group().now_or_never(), Some(Ok(None))),
6831 "a dead parked group must not be delivered or hold the stream open"
6832 );
6833 }
6834
6835 #[tokio::test]
6839 async fn evicted_parked_group_wakes_the_clean_end() {
6840 use std::sync::atomic::{AtomicUsize, Ordering};
6841 use std::task::{Context, Wake};
6842
6843 struct CountWaker(AtomicUsize);
6846 impl Wake for CountWaker {
6847 fn wake(self: std::sync::Arc<Self>) {
6848 self.0.fetch_add(1, Ordering::SeqCst);
6849 }
6850 }
6851
6852 let producer = track_producer("test", None);
6853 let mut consumer = producer.subscribe(None);
6854
6855 producer.create_group(group::Info { sequence: 0 }).unwrap();
6856 assert_eq!(
6857 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6858 0
6859 );
6860
6861 consumer.set_groups(..1);
6862 let straggler = producer.create_group(group::Info { sequence: 1 }).unwrap();
6863 assert!(consumer.recv_group().now_or_never().is_none(), "parked at the cap");
6864 producer.finish().unwrap();
6865
6866 let counter = std::sync::Arc::new(CountWaker(AtomicUsize::new(0)));
6867 let waker = std::task::Waker::from(counter.clone());
6868 let mut cx = Context::from_waker(&waker);
6869 let mut fut = std::pin::pin!(consumer.recv_group());
6870 assert!(
6871 fut.as_mut().poll(&mut cx).is_pending(),
6872 "the parked group holds it open"
6873 );
6874
6875 straggler.abort(Error::Old).unwrap();
6876 assert!(counter.0.load(Ordering::SeqCst) > 0, "the eviction wakeup was lost");
6877 assert!(matches!(fut.as_mut().poll(&mut cx), Poll::Ready(Ok(None))));
6878 }
6879
6880 #[tokio::test]
6882 async fn end_at_zero_is_the_empty_range() {
6883 let producer = track_producer("test", None);
6884 let mut consumer = producer.subscribe(replay());
6885 producer.create_group(group::Info { sequence: 0 }).unwrap();
6886 producer.create_group(group::Info { sequence: 1 }).unwrap();
6887
6888 consumer.set_groups(..0);
6889 assert!(
6890 consumer.recv_group().now_or_never().is_none(),
6891 "empty cap delivers nothing"
6892 );
6893
6894 consumer.set_groups(..1);
6895 assert_eq!(
6896 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6897 0
6898 );
6899 assert!(consumer.recv_group().now_or_never().is_none(), "group 1 stays parked");
6900 }
6901
6902 #[tokio::test]
6905 async fn empty_local_cap_holds_while_another_subscriber_requests_everything() {
6906 let producer = track_producer("test", None);
6907 let mut everything = producer.subscribe(replay());
6908 let mut empty = producer.subscribe(Subscription::default().with_end(Position::group(0)));
6909 empty.set_groups((Bound::Unbounded, Position::group(0).group_end()));
6910
6911 for s in 0..3 {
6912 producer.create_group(group::Info { sequence: s }).unwrap();
6913 }
6914
6915 assert_eq!(
6916 everything
6917 .recv_group()
6918 .now_or_never()
6919 .unwrap()
6920 .unwrap()
6921 .unwrap()
6922 .sequence,
6923 0
6924 );
6925 assert!(
6926 empty.recv_group().now_or_never().is_none(),
6927 "local empty cap must not ride the unbounded aggregate"
6928 );
6929
6930 empty.set_groups(..2);
6931 assert_eq!(empty.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence, 0);
6932 assert_eq!(empty.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence, 1);
6933 assert!(empty.recv_group().now_or_never().is_none(), "still capped at 2");
6934
6935 empty.set_groups(..);
6936 assert_eq!(empty.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence, 2);
6937 }
6938
6939 #[tokio::test]
6941 async fn end_at_frame_limited_last_group() {
6942 let producer = track_producer("test", None);
6943 let mut consumer = producer.subscribe(replay()).ordered();
6944 let end = Position::after(1, 1).unwrap();
6945 consumer.set_groups((Bound::Unbounded, end.group_end()));
6946
6947 for s in 0..3u64 {
6948 let mut group = producer.create_group(group::Info { sequence: s }).unwrap();
6949 for i in 0..3u8 {
6950 group.write_frame(Timestamp::ZERO, vec![i]).unwrap();
6951 }
6952 group.finish().unwrap();
6953 }
6954
6955 assert_eq!(
6956 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6957 0
6958 );
6959 let mut last = consumer.next_group().now_or_never().unwrap().unwrap().unwrap();
6960 assert_eq!(last.sequence, 1);
6961 last.set_frames(..end.frame);
6962 assert_eq!(
6963 last.read_frame().now_or_never().unwrap().unwrap().unwrap().payload[0],
6964 0
6965 );
6966 assert_eq!(
6967 last.read_frame().now_or_never().unwrap().unwrap().unwrap().payload[0],
6968 1
6969 );
6970 assert!(
6971 last.read_frame().now_or_never().unwrap().unwrap().is_none(),
6972 "frame cap is exclusive"
6973 );
6974 assert!(
6975 consumer.next_group().now_or_never().is_none(),
6976 "group 2 is past the exclusive group cap"
6977 );
6978 }
6979
6980 #[tokio::test]
6982 async fn end_at_maximum_group_is_unbounded() {
6983 let producer = track_producer("test", None);
6984 let mut consumer = producer.subscribe(replay()).ordered();
6985 consumer.set_groups(..=u64::MAX);
6986
6987 producer.create_group(group::Info { sequence: 0 }).unwrap();
6988 producer.create_group(group::Info { sequence: u64::MAX }).unwrap();
6989
6990 assert_eq!(
6991 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6992 0
6993 );
6994 assert_eq!(
6995 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
6996 u64::MAX
6997 );
6998 }
6999
7000 #[test]
7001 fn write_frame_rejects_an_oversized_frame_before_appending_its_group() {
7002 let mut producer = track_producer("test", None);
7003 let frame = bytes::Bytes::from(vec![0; group::MAX_CACHE_BYTES as usize + 1]);
7004
7005 assert!(matches!(
7006 producer.write_frame(Timestamp::ZERO, frame),
7007 Err(Error::FrameTooLarge)
7008 ));
7009 assert_eq!(producer.latest(), None, "the rejected frame did not publish a group");
7010 }
7011
7012 #[test]
7013 fn append_group_returns_bounds_exceeded_on_sequence_overflow() {
7014 let producer = track_producer("test", None);
7015 {
7016 let mut state = producer.state.write().ok().unwrap();
7017 state.max_sequence = Some(u64::MAX);
7018 }
7019
7020 assert!(matches!(producer.append_group(), Err(Error::BoundsExceeded(_))));
7021 }
7022
7023 #[tokio::test]
7024 async fn fetch_cache_hit() {
7025 let producer = track_producer("test", None);
7026
7027 let mut group = producer.append_group().unwrap(); group
7030 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hello"))
7031 .unwrap();
7032 group.finish().unwrap();
7033
7034 let dynamic = producer.dynamic();
7037 let consumer = producer.consume();
7038 assert!(consumer.peek_group(0).is_some());
7039 let mut g = consumer.fetch_group(0, None).await.unwrap();
7040 assert_eq!(g.sequence, 0);
7041 assert_eq!(&g.read_frame().await.unwrap().unwrap().payload[..], b"hello");
7042
7043 assert!(dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending());
7045 }
7046
7047 #[tokio::test]
7048 async fn fetch_miss_signals_dynamic() {
7049 let producer = track_producer("test", None);
7050 let dynamic = producer.dynamic();
7051 let consumer = producer.consume();
7052
7053 assert!(consumer.peek_group(5).is_none());
7057 let pending = consumer.fetch_group(5, group::Fetch::default().with_priority(7));
7058 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
7059
7060 let req = dynamic
7061 .requested_group()
7062 .now_or_never()
7063 .expect("should not block")
7064 .unwrap();
7065 assert_eq!(req.sequence(), 5);
7066 assert_eq!(req.priority(), 7);
7067
7068 let mut group = req.accept(None).unwrap();
7070 group
7071 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hi"))
7072 .unwrap();
7073 group.finish().unwrap();
7074
7075 let mut g = pending.await.unwrap();
7076 assert_eq!(g.sequence, 5);
7077 assert_eq!(&g.read_frame().await.unwrap().unwrap().payload[..], b"hi");
7078 }
7079
7080 #[tokio::test]
7081 async fn fetch_miss_rejects() {
7082 let producer = track_producer("test", None);
7083 let dynamic = producer.dynamic();
7084 let consumer = producer.consume();
7085
7086 let pending = consumer.fetch_group(5, None);
7087 let req = dynamic
7088 .requested_group()
7089 .now_or_never()
7090 .expect("should not block")
7091 .unwrap();
7092
7093 req.reject(Error::Cancel);
7094 assert!(matches!(pending.await, Err(Error::Cancel)));
7095 let fetch = producer.state.read().fetch.clone();
7096 assert!(fetch.read().is_empty());
7097 }
7098
7099 #[tokio::test]
7100 async fn fetch_miss_drop_rejects() {
7101 let producer = track_producer("test", None);
7102 let dynamic = producer.dynamic();
7103 let consumer = producer.consume();
7104
7105 let pending = consumer.fetch_group(5, None);
7106 let req = dynamic
7107 .requested_group()
7108 .now_or_never()
7109 .expect("should not block")
7110 .unwrap();
7111
7112 drop(req);
7113 assert!(matches!(pending.await, Err(Error::Dropped)));
7114 }
7115
7116 #[tokio::test]
7117 async fn fetch_reject_does_not_poison_retry() {
7118 let producer = track_producer("test", None);
7119 let dynamic = producer.dynamic();
7120 let consumer = producer.consume();
7121
7122 let pending = consumer.fetch_group(5, None);
7123 let req = dynamic
7124 .requested_group()
7125 .now_or_never()
7126 .expect("should not block")
7127 .unwrap();
7128 req.reject(Error::Cancel);
7129 assert!(matches!(pending.await, Err(Error::Cancel)));
7130
7131 let retry = consumer.fetch_group(5, None);
7132 let req = dynamic
7133 .requested_group()
7134 .now_or_never()
7135 .expect("should not block")
7136 .unwrap();
7137 let mut group = req.accept(None).unwrap();
7138 group
7139 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"retry"))
7140 .unwrap();
7141 group.finish().unwrap();
7142
7143 let mut group = retry.await.unwrap();
7144 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"retry");
7145 }
7146
7147 #[tokio::test]
7151 async fn fetch_ignores_a_group_that_starts_too_late() {
7152 let producer = track_producer("test", None);
7153 let dynamic = producer.dynamic();
7154 let consumer = producer.consume();
7155
7156 let mut group = producer.create_group(group::Info { sequence: 0 }).unwrap();
7158 group.start_at(3).unwrap();
7159 group
7160 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"tail"))
7161 .unwrap();
7162 group.finish().unwrap();
7163
7164 let fetch = consumer.fetch_group(0, group::Fetch::default().with_frame_start(3));
7166 let cached = fetch.now_or_never().expect("covered by the cache").unwrap();
7167 assert_eq!(cached.index(), 3);
7168
7169 let mut fetch = std::pin::pin!(consumer.fetch_group(0, None));
7172 assert!(
7173 futures::poll!(fetch.as_mut()).is_pending(),
7174 "must not answer from the tail"
7175 );
7176
7177 let request = dynamic.requested_group().await.unwrap();
7178 assert_eq!((request.sequence(), request.frame_start()), (0, 0));
7179
7180 let mut whole = request.accept(None).unwrap();
7182 whole
7183 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"head"))
7184 .unwrap();
7185 whole.finish().unwrap();
7186
7187 let mut served = fetch.await.unwrap();
7188 assert_eq!(served.index(), 0);
7189 assert_eq!(
7190 served.read_frame().await.unwrap().unwrap().payload,
7191 bytes::Bytes::from_static(b"head")
7192 );
7193 }
7194
7195 #[tokio::test]
7198 async fn fetch_widens_or_fails_cleanly() {
7199 let producer = track_producer("test", None);
7200 let dynamic = producer.dynamic();
7201 let consumer = producer.consume();
7202
7203 let _narrow = consumer.fetch_group(0, group::Fetch::default().with_frame_start(5));
7204 let mut narrow = std::pin::pin!(_narrow);
7205 assert!(futures::poll!(narrow.as_mut()).is_pending());
7206
7207 let _wider = consumer.fetch_group(0, group::Fetch::default().with_frame_start(2));
7209 let mut wider = std::pin::pin!(_wider);
7210 assert!(futures::poll!(wider.as_mut()).is_pending());
7211
7212 let request = dynamic.requested_group().await.unwrap();
7213 assert_eq!(request.frame_start(), 2, "widened while queued");
7214
7215 let _widest = consumer.fetch_group(0, group::Fetch::default().with_frame_start(0));
7218 let mut widest = std::pin::pin!(_widest);
7219 assert!(futures::poll!(widest.as_mut()).is_pending());
7220 assert_eq!(request.frame_start(), 2, "the in-flight range is already on the wire");
7221
7222 let mut group = request.accept(None).unwrap();
7223 group.start_at(2).unwrap();
7226 group
7227 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"from2"))
7228 .unwrap();
7229 group.finish().unwrap();
7230
7231 assert_eq!(narrow.await.unwrap().index(), 5);
7234 assert_eq!(wider.await.unwrap().index(), 2);
7235 assert!(matches!(widest.await, Err(Error::NotFound)));
7236 }
7237
7238 #[tokio::test]
7239 async fn fetch_coalesces_concurrent() {
7240 let producer = track_producer("test", None);
7241 let dynamic = producer.dynamic();
7242 let consumer = producer.consume();
7243
7244 let first = consumer.fetch_group(5, group::Fetch::default().with_priority(1));
7247 let second = consumer.fetch_group(5, group::Fetch::default().with_priority(7));
7248 assert!(kio::Pollable::poll(&*first, &kio::Waiter::noop()).is_pending());
7249
7250 let req = dynamic
7251 .requested_group()
7252 .now_or_never()
7253 .expect("should not block")
7254 .unwrap();
7255 assert_eq!(req.sequence(), 5);
7256 assert_eq!(req.priority(), 7);
7257 assert!(
7258 dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending(),
7259 "the second fetch queued a duplicate request"
7260 );
7261
7262 let third = consumer.fetch_group(5, None);
7264
7265 let mut group = req.accept(None).unwrap();
7267 group
7268 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hi"))
7269 .unwrap();
7270 group.finish().unwrap();
7271
7272 assert_eq!(first.await.unwrap().sequence, 5);
7273 assert_eq!(second.await.unwrap().sequence, 5);
7274 assert_eq!(third.await.unwrap().sequence, 5);
7275 }
7276
7277 #[tokio::test]
7278 async fn fetch_coalesced_reject_fails_all() {
7279 let producer = track_producer("test", None);
7280 let dynamic = producer.dynamic();
7281 let consumer = producer.consume();
7282
7283 let first = consumer.fetch_group(5, None);
7284 let second = consumer.fetch_group(5, None);
7285 let req = dynamic
7286 .requested_group()
7287 .now_or_never()
7288 .expect("should not block")
7289 .unwrap();
7290 req.reject(Error::Cancel);
7291
7292 assert!(matches!(first.await, Err(Error::Cancel)));
7293 assert!(matches!(second.await, Err(Error::Cancel)));
7294
7295 let retry = consumer.fetch_group(5, None);
7297 assert!(kio::Pollable::poll(&*retry, &kio::Waiter::noop()).is_pending());
7298 let req = dynamic
7299 .requested_group()
7300 .now_or_never()
7301 .expect("should not block")
7302 .unwrap();
7303 assert_eq!(req.sequence(), 5);
7304 }
7305
7306 #[tokio::test]
7307 async fn fetch_queued_fails_when_handlers_leave() {
7308 let producer = track_producer("test", None);
7309 let dynamic = producer.dynamic();
7310 let consumer = producer.consume();
7311
7312 let pending = consumer.fetch_group(5, None);
7314 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
7315 drop(dynamic);
7316 assert!(matches!(pending.await, Err(Error::NotFound)));
7317
7318 let fetch = producer.state.read().fetch.clone();
7320 assert!(fetch.read().is_empty());
7321 }
7322
7323 #[tokio::test]
7324 async fn fetch_miss_no_dynamic_not_found() {
7325 let producer = track_producer("test", None);
7328 producer.append_group().unwrap(); let consumer = producer.consume();
7330 assert!(matches!(consumer.fetch_group(5, None).await, Err(Error::NotFound)));
7331 }
7332
7333 #[tokio::test]
7334 async fn fetch_past_final_not_found() {
7335 let producer = track_producer("test", None);
7336 producer.append_group().unwrap(); producer.finish().unwrap(); let dynamic = producer.dynamic();
7342 let consumer = producer.consume();
7343 assert!(matches!(consumer.fetch_group(5, None).await, Err(Error::NotFound)));
7344
7345 assert!(dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending());
7347 }
7348
7349 fn pooled_producer(capacity: u64) -> (Producer, cache::Pool) {
7351 let config = cache::Config::default()
7352 .with_capacity(capacity)
7353 .with_expiry(cache::DEFAULT_EXPIRY);
7354 let pool = cache::Pool::new(config);
7355 let broadcast = broadcast::Info {
7356 pool: pool.clone(),
7357 ..Default::default()
7358 };
7359 let producer = Producer::new(Arc::new(broadcast), "test", None);
7360 (producer, pool)
7361 }
7362
7363 fn finished_group(producer: &mut Producer, size: usize) -> u64 {
7364 let mut group = producer.append_group().unwrap();
7365 group
7366 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; size]))
7367 .unwrap();
7368 group.finish().unwrap();
7369 group.sequence
7370 }
7371
7372 #[tokio::test]
7375 async fn debt_evicts_oldest_group() {
7376 let (mut producer, pool) = pooled_producer(10_000);
7378
7379 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let consumer = producer.consume();
7384 assert!(consumer.peek_group(0).is_none(), "oldest group is evicted");
7385 assert!(consumer.peek_group(2).is_some(), "latest group survives");
7386 assert!(
7389 pool.used() <= 2 * (10_000 + cache::ENTRY_OVERHEAD),
7390 "usage hovers near capacity: {}",
7391 pool.used()
7392 );
7393
7394 let mut subscriber = producer.subscribe(replay());
7396 assert!(subscriber.assert_group().sequence > 0, "evicted group is not delivered");
7397 }
7398
7399 #[tokio::test]
7401 async fn latest_group_never_evicted() {
7402 let (mut producer, pool) = pooled_producer(100);
7404 finished_group(&mut producer, 1000); assert!(pool.used() > 100, "the latest may exceed the budget");
7406
7407 finished_group(&mut producer, 1000); finished_group(&mut producer, 1000); let consumer = producer.consume();
7412 assert!(consumer.peek_group(0).is_none());
7413 let mut group = consumer.peek_group(2).expect("latest survives");
7414 assert_eq!(group.read_frame().await.unwrap().unwrap().payload.len(), 1000);
7415 }
7416
7417 #[tokio::test]
7421 async fn fetch_refresh_survives_eviction() {
7422 let (mut producer, _pool) = pooled_producer(10_000);
7423 let consumer = producer.consume();
7424
7425 finished_group(&mut producer, 3_000); crate::model::clock::advance(Duration::from_secs(1));
7427 finished_group(&mut producer, 3_000); crate::model::clock::advance(Duration::from_secs(1));
7429 finished_group(&mut producer, 3_000); crate::model::clock::advance(Duration::from_millis(500));
7431
7432 let mut fetched = consumer.fetch_group(0, None).await.unwrap();
7434 assert_eq!(fetched.read_frame().await.unwrap().unwrap().payload.len(), 3_000);
7435 crate::model::clock::advance(Duration::from_millis(500));
7436
7437 finished_group(&mut producer, 3_000); crate::model::clock::advance(Duration::from_secs(1));
7441 finished_group(&mut producer, 3_000); assert!(consumer.peek_group(0).is_some(), "refreshed group survives");
7444 assert!(consumer.peek_group(1).is_none(), "unread group is evicted instead");
7445 }
7446
7447 #[tokio::test]
7450 async fn eviction_aborts_readers() {
7451 let (mut producer, _pool) = pooled_producer(10_000);
7452 let mut subscriber = producer.subscribe(None);
7453
7454 finished_group(&mut producer, 10_000); let mut group0 = subscriber.assert_group();
7456
7457 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let read = group0.read_frame().await;
7461 assert!(matches!(read, Err(Error::Evicted)), "expected Evicted, got {read:?}");
7462 }
7463
7464 #[tokio::test]
7468 async fn small_writes_carry_debt() {
7469 let unit = 100 * cache::ENTRY_OVERHEAD;
7472 let (mut producer, pool) = pooled_producer(22 * unit);
7473 let consumer = producer.consume();
7474
7475 finished_group(&mut producer, 20 * unit as usize); for _ in 0..3 {
7480 finished_group(&mut producer, unit as usize);
7481 }
7482 assert!(consumer.peek_group(0).is_some(), "debt smaller than the victim carries");
7483
7484 for _ in 0..20 {
7486 finished_group(&mut producer, unit as usize);
7487 }
7488 assert!(
7489 consumer.peek_group(0).is_none(),
7490 "accumulated debt evicts the large group"
7491 );
7492 assert!(pool.used() <= 24 * unit, "usage hovers near capacity: {}", pool.used());
7495 }
7496
7497 #[tokio::test]
7501 async fn payment_capped_per_write() {
7502 let (mut producer, pool) = pooled_producer(1 << 40);
7503 for _ in 0..10 {
7504 finished_group(&mut producer, 1_000);
7505 }
7506
7507 pool.resize(100);
7509 let before = pool.used();
7510
7511 finished_group(&mut producer, 1_000);
7513
7514 let consumer = producer.consume();
7515 assert!(consumer.peek_group(0).is_none(), "the oldest groups are evicted");
7516 assert!(consumer.peek_group(1).is_none());
7517 assert!(consumer.peek_group(2).is_some(), "the backlog drains gradually");
7518 assert!(pool.used() > before - 4_000, "one write must not dump the backlog");
7519 }
7520
7521 #[tokio::test]
7525 async fn accept_preserves_write_accounting() {
7526 let config = cache::Config::default()
7527 .with_capacity(12_000)
7528 .with_expiry(cache::DEFAULT_EXPIRY);
7529 let pool = cache::Pool::new(config);
7530 let broadcast = broadcast::Info {
7531 pool: pool.clone(),
7532 ..Default::default()
7533 };
7534 let request = Request::new(Arc::new(broadcast), "test");
7535 let dynamic = request.dynamic();
7536 let consumer = request.consume();
7537
7538 let pending = consumer.fetch_group(0, None);
7540 let req = dynamic
7541 .requested_group()
7542 .now_or_never()
7543 .expect("should not block")
7544 .unwrap();
7545 let mut backfill = req.accept(None).unwrap();
7546 pending.await.unwrap();
7547 backfill
7548 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 30_000]))
7549 .unwrap();
7550
7551 let producer = request.accept(None);
7554 producer.append_group().unwrap().finish().unwrap();
7555 producer.append_group().unwrap().finish().unwrap();
7556
7557 assert!(
7558 producer.consume().peek_group(0).is_none(),
7559 "pre-accept backfill growth is reclaimed after accept"
7560 );
7561 assert!(pool.used() <= 13_000, "usage converges: {}", pool.used());
7562 }
7563
7564 #[tokio::test]
7567 async fn recreated_sequence_bounds_eviction_hints() {
7568 let (producer, _pool) = pooled_producer(1 << 40);
7569 producer.create_group(5u64.into()).unwrap().finish().unwrap();
7570
7571 for _ in 0..200 {
7572 let group = producer.create_group(1u64.into()).unwrap();
7573 group.abort(Error::Cancel).unwrap();
7574 }
7575
7576 let state = producer.state.read();
7577 assert!(
7578 state.evict.len() <= 2 * state.lookup.len() + EVICT_SLACK,
7579 "stale hints are compacted: {} entries for {} slots",
7580 state.evict.len(),
7581 state.lookup.len()
7582 );
7583 }
7584
7585 #[tokio::test]
7588 async fn same_tick_write_outranks_inserted() {
7589 let unit = 100 * cache::ENTRY_OVERHEAD;
7592 let (mut producer, _pool) = pooled_producer(10 * unit);
7594
7595 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();
7602 assert!(consumer.peek_group(0).is_none(), "insert-only content pays first");
7603 assert!(consumer.peek_group(1).is_some(), "same-tick written content survives");
7604 }
7605
7606 #[tokio::test]
7609 async fn frame_only_writer_pays() {
7610 let (producer, pool) = pooled_producer(2_000);
7611 let mut demoted = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); demoted
7617 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 300_000]))
7618 .unwrap();
7619
7620 assert!(
7621 pool.used() <= 5_000,
7622 "the frame write settled the debt: {}",
7623 pool.used()
7624 );
7625 assert!(matches!(demoted.finish(), Err(Error::Evicted)));
7626 }
7627
7628 #[tokio::test]
7631 async fn each_track_owns_its_account() {
7632 let broadcast = Arc::new(broadcast::Info::default());
7633 let info = Info::default();
7634 let a = Producer::new(broadcast.clone(), "a", info.clone());
7635 let b = Producer::new(broadcast, "b", info);
7636
7637 let a = a.state.read().cache.clone();
7638 let b = b.state.read().cache.clone();
7639 assert!(!Arc::ptr_eq(&a, &b), "each track owns its account");
7640 }
7641
7642 #[tokio::test]
7645 async fn a_dynamic_defers_teardown() {
7646 let (mut producer, pool) = pooled_producer(1 << 40);
7647 let dynamic = producer.dynamic();
7648 finished_group(&mut producer, 100);
7649
7650 drop(producer);
7651 assert!(pool.used() > 0, "the handler still serves the cache");
7652
7653 drop(dynamic);
7654 assert_eq!(pool.used(), 0, "the last handle tears it down");
7655 }
7656
7657 #[tokio::test]
7663 async fn finished_track_frees_its_cache() {
7664 let (mut producer, pool) = pooled_producer(1 << 40);
7665 finished_group(&mut producer, 100);
7666 producer.finish().unwrap();
7667
7668 let state = producer.state.downgrade();
7669 drop(producer);
7670
7671 assert!(state.upgrade().is_none(), "the track state is freed");
7672 assert_eq!(pool.used(), 0, "so are its cached bytes");
7673 }
7674
7675 #[tokio::test]
7679 async fn teardown_ignores_a_settling_group() {
7680 let (mut producer, pool) = pooled_producer(1 << 40);
7681 finished_group(&mut producer, 100);
7682
7683 let settling = producer.state.downgrade().upgrade().expect("open");
7685 drop(producer);
7686
7687 assert_eq!(pool.used(), 0, "the abrupt teardown still released the cache");
7688 drop(settling);
7689 }
7690
7691 #[tokio::test]
7694 async fn cached_group_outlives_its_track() {
7695 let (mut producer, pool) = pooled_producer(1 << 40);
7696 let sequence = finished_group(&mut producer, 100);
7697 let group = producer.consume().peek_group(sequence).expect("cached");
7698 producer.finish().unwrap();
7699
7700 let state = producer.state.downgrade();
7701 drop(producer);
7702 assert!(state.upgrade().is_none(), "the track state is freed");
7703 assert!(pool.used() > 0, "the retained group keeps its own bytes");
7704
7705 drop(group);
7706 assert_eq!(pool.used(), 0, "which it releases when dropped");
7707 }
7708
7709 #[tokio::test]
7713 async fn pre_accept_backfill_settles_late_writes() {
7714 let config = cache::Config::default()
7715 .with_capacity(2_000)
7716 .with_expiry(cache::DEFAULT_EXPIRY);
7717 let pool = cache::Pool::new(config);
7718 let broadcast = broadcast::Info {
7719 pool: pool.clone(),
7720 ..Default::default()
7721 };
7722 let request = Request::new(Arc::new(broadcast), "test");
7723 let dynamic = request.dynamic();
7724 let consumer = request.consume();
7725
7726 let pending = consumer.fetch_group(0, None);
7728 let req = dynamic
7729 .requested_group()
7730 .now_or_never()
7731 .expect("should not block")
7732 .unwrap();
7733 let mut backfill = req.accept(None).unwrap();
7734 pending.await.unwrap();
7735
7736 let producer = request.accept(None);
7738 producer.append_group().unwrap().finish().unwrap();
7739
7740 backfill
7743 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 300_000]))
7744 .unwrap();
7745
7746 assert!(
7747 pool.used() <= 5_000,
7748 "the frame write settled the debt: {}",
7749 pool.used()
7750 );
7751 }
7752
7753 #[tokio::test]
7757 async fn write_restarts_retention_clock() {
7758 let (producer, _pool) = pooled_producer(1 << 40);
7759 let mut straggler = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); crate::model::clock::advance(cache::DEFAULT_EXPIRY + Duration::from_secs(1));
7764 straggler
7765 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 100]))
7766 .unwrap();
7767 producer.append_group().unwrap().finish().unwrap(); let consumer = producer.consume();
7770 assert!(consumer.peek_group(0).is_some(), "the write restarted the clock");
7771
7772 crate::model::clock::advance(cache::DEFAULT_EXPIRY + Duration::from_secs(1));
7774 producer.append_group().unwrap().finish().unwrap(); assert!(consumer.peek_group(0).is_none(), "idle content still expires");
7776 }
7777
7778 #[tokio::test]
7781 async fn refreshed_front_does_not_starve_expiry() {
7782 let (producer, _pool) = pooled_producer(1 << 40);
7783 let dynamic = producer.dynamic();
7784 let consumer = producer.consume();
7785
7786 producer.create_group(10u64.into()).unwrap().finish().unwrap();
7787 for sequence in 1..=5u64 {
7788 let pending = consumer.fetch_group(sequence, None);
7789 let req = dynamic
7790 .requested_group()
7791 .now_or_never()
7792 .expect("should not block")
7793 .unwrap();
7794 let mut group = req.accept(None).unwrap();
7795 group
7796 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 100]))
7797 .unwrap();
7798 group.finish().unwrap();
7799 pending.await.unwrap();
7800 }
7801
7802 crate::model::clock::advance(cache::DEFAULT_EXPIRY + Duration::from_secs(1));
7805 for sequence in 1..=4u64 {
7806 consumer.fetch_group(sequence, None).await.unwrap();
7807 }
7808
7809 for _ in 0..3 {
7811 producer.append_group().unwrap().finish().unwrap();
7812 }
7813 assert!(consumer.peek_group(5).is_none(), "expired backfill is reclaimed");
7814 assert!(consumer.peek_group(1).is_some(), "refreshed backfill survives");
7815 }
7816
7817 #[tokio::test]
7820 async fn recreated_sequence_delivered_once() {
7821 let (producer, _pool) = pooled_producer(1 << 40);
7822
7823 producer.create_group(0u64.into()).unwrap().finish().unwrap();
7824 let aborted = producer.create_group(1u64.into()).unwrap();
7825 aborted.abort(Error::Cancel).unwrap();
7826 producer.create_group(2u64.into()).unwrap().finish().unwrap();
7827 producer.create_group(1u64.into()).unwrap().finish().unwrap();
7828
7829 let mut subscriber = producer.subscribe(replay());
7830 assert_eq!(subscriber.assert_group().sequence, 0);
7831 assert_eq!(subscriber.assert_group().sequence, 2);
7832 assert_eq!(
7833 subscriber.assert_group().sequence,
7834 1,
7835 "replacement arrives at its own position"
7836 );
7837 subscriber.assert_no_group();
7838 }
7839
7840 #[tokio::test]
7844 async fn datagrams_do_not_block_eviction() {
7845 let (mut producer, pool) = pooled_producer(1_000);
7846 for _ in 0..10 {
7847 finished_group(&mut producer, 1_000);
7848 producer.append_datagram(Timestamp::ZERO, &b"beat"[..]).unwrap();
7849 }
7850
7851 let consumer = producer.consume();
7852 assert!(consumer.peek_group(0).is_none(), "old groups still evict");
7853 assert!(
7854 pool.used() < 4 * 1_256,
7855 "interleaved datagrams must not bypass the budget: {}",
7856 pool.used()
7857 );
7858 }
7859
7860 #[tokio::test]
7864 async fn aborted_group_leaves_no_ghost_sample() {
7865 let (producer, pool) = pooled_producer(1 << 40);
7866 let group0 = producer.append_group().unwrap();
7867 producer.append_group().unwrap(); assert!(pool.average().is_some(), "demoted group is sampled");
7870 group0.abort(Error::Cancel).unwrap();
7871 assert_eq!(pool.average(), None, "the abort must remove the sample");
7872 }
7873
7874 #[tokio::test]
7877 async fn empty_groups_repay_overhead() {
7878 let (producer, pool) = pooled_producer(1_000);
7879 for _ in 0..100 {
7880 let group = producer.append_group().unwrap();
7881 group.finish().unwrap();
7882 }
7883
7884 assert!(
7885 pool.used() <= 3_000,
7886 "empty-group overhead must stay near the budget: {}",
7887 pool.used()
7888 );
7889 }
7890
7891 #[tokio::test]
7894 async fn growth_on_demoted_group_is_billed() {
7895 let (producer, pool) = pooled_producer(2_000);
7896 let mut straggler = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); straggler
7901 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 10_000]))
7902 .unwrap();
7903
7904 producer.append_group().unwrap().finish().unwrap(); let consumer = producer.consume();
7908 assert!(consumer.peek_group(0).is_none(), "the ballooned group is evicted");
7909 assert!(pool.used() <= 3_000, "growth is reclaimed: {}", pool.used());
7910 }
7911
7912 #[tokio::test]
7915 async fn refilled_sequence_stays_out_of_subscriptions() {
7916 let (producer, _pool) = pooled_producer(1 << 40);
7917 let dynamic = producer.dynamic();
7918 let consumer = producer.consume();
7919
7920 producer.create_group(0u64.into()).unwrap().finish().unwrap();
7921 let aborted = producer.create_group(1u64.into()).unwrap();
7922 aborted.abort(Error::Cancel).unwrap();
7923 producer.create_group(2u64.into()).unwrap().finish().unwrap();
7924
7925 let pending = consumer.fetch_group(1, None);
7928 let req = dynamic
7929 .requested_group()
7930 .now_or_never()
7931 .expect("should not block")
7932 .unwrap();
7933 let mut group = req.accept(None).unwrap();
7934 group
7935 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"backfill"))
7936 .unwrap();
7937 group.finish().unwrap();
7938 pending.await.unwrap();
7939
7940 assert!(consumer.peek_group(1).is_some());
7942 let mut subscriber = producer.subscribe(replay());
7943 assert_eq!(subscriber.assert_group().sequence, 0);
7944 assert_eq!(subscriber.assert_group().sequence, 2);
7945 subscriber.assert_no_group();
7946 }
7947
7948 #[tokio::test]
7951 async fn expired_backfill_behind_refreshed_reclaimed() {
7952 let (producer, _pool) = pooled_producer(1 << 40);
7953 let dynamic = producer.dynamic();
7954 let consumer = producer.consume();
7955
7956 producer.create_group(5u64.into()).unwrap().finish().unwrap();
7957 for sequence in [2u64, 3u64] {
7958 let pending = consumer.fetch_group(sequence, None);
7959 let req = dynamic
7960 .requested_group()
7961 .now_or_never()
7962 .expect("should not block")
7963 .unwrap();
7964 let mut group = req.accept(None).unwrap();
7965 group
7966 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 100]))
7967 .unwrap();
7968 group.finish().unwrap();
7969 pending.await.unwrap();
7970 }
7971
7972 crate::model::clock::advance(cache::DEFAULT_EXPIRY / 2 + Duration::from_secs(1));
7974 consumer.fetch_group(2, None).await.unwrap();
7975 crate::model::clock::advance(cache::DEFAULT_EXPIRY / 2 + Duration::from_secs(1));
7976 producer.create_group(6u64.into()).unwrap().finish().unwrap();
7977
7978 let consumer = producer.consume();
7979 assert!(consumer.peek_group(2).is_some(), "refreshed backfill survives");
7980 assert!(consumer.peek_group(3).is_none(), "expired backfill is reclaimed");
7981 }
7982
7983 #[tokio::test]
7986 async fn same_tick_fetch_protects() {
7987 let (mut producer, _pool) = pooled_producer(10_000);
7989 let consumer = producer.consume();
7990
7991 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();
7996
7997 finished_group(&mut producer, 3_000); finished_group(&mut producer, 3_000); assert!(consumer.peek_group(0).is_some(), "same-tick refresh protects");
8001 assert!(consumer.peek_group(1).is_none(), "the unread group dies instead");
8002 }
8003
8004 #[tokio::test]
8008 async fn refetched_latest_stays_protected() {
8009 let (producer, _pool) = pooled_producer(10_000);
8010 let dynamic = producer.dynamic();
8011 let consumer = producer.consume();
8012
8013 let straggler = producer.append_group().unwrap(); let latest = producer.append_group().unwrap(); latest.abort(Error::Cancel).unwrap();
8018
8019 let pending = consumer.fetch_group(1, None);
8021 let req = dynamic
8022 .requested_group()
8023 .now_or_never()
8024 .expect("should not block")
8025 .unwrap();
8026 let mut group = req.accept(None).unwrap();
8027 group
8028 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 1000]))
8029 .unwrap();
8030 group.finish().unwrap();
8031 pending.await.unwrap();
8032
8033 {
8036 let state = producer.state.read();
8037 assert!(state.lookup.contains_key(&1), "refetched group is cached");
8038 assert!(
8039 state.evict.iter().all(|(sequence, _)| *sequence != 1),
8040 "the live edge must not be an eviction candidate"
8041 );
8042 }
8043 drop(straggler);
8044 }
8045
8046 #[tokio::test]
8049 async fn eviction_allows_refetch() {
8050 let (mut producer, _pool) = pooled_producer(10_000);
8051 let dynamic = producer.dynamic();
8052
8053 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let consumer = producer.consume();
8058 assert!(consumer.peek_group(0).is_none());
8059 let pending = consumer.fetch_group(0, None);
8060
8061 let req = dynamic
8062 .requested_group()
8063 .now_or_never()
8064 .expect("should not block")
8065 .unwrap();
8066 assert_eq!(req.sequence(), 0);
8067
8068 let mut group = req.accept(None).unwrap();
8069 group
8070 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"refetched"))
8071 .unwrap();
8072 group.finish().unwrap();
8073
8074 let mut group = pending.await.unwrap();
8075 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"refetched");
8076 }
8077
8078 #[test]
8088 fn an_aborted_group_releases_its_sequence() {
8089 let producer = track_producer("test", None);
8090 let consumer = producer.consume();
8091
8092 let mut group = producer.create_group(group::Info { sequence: 3 }).unwrap();
8093 group
8094 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"head"))
8095 .unwrap();
8096
8097 assert!(matches!(
8099 producer.create_group(group::Info { sequence: 3 }),
8100 Err(Error::Duplicate)
8101 ));
8102 assert!(consumer.peek_group(3).is_some());
8103
8104 group.abort(Error::Cancel).unwrap();
8105
8106 assert!(consumer.peek_group(3).is_none(), "an aborted slot is a cache miss");
8107 producer
8108 .create_group(group::Info { sequence: 3 })
8109 .expect("an aborted slot releases its sequence")
8110 .finish()
8111 .unwrap();
8112 }
8113
8114 #[tokio::test]
8117 async fn fetched_backfill_not_subscribed() {
8118 let (producer, _pool) = pooled_producer(1 << 40);
8119 let dynamic = producer.dynamic();
8120 let consumer = producer.consume();
8121
8122 producer.create_group(5u64.into()).unwrap().finish().unwrap();
8124 producer.create_group(6u64.into()).unwrap().finish().unwrap();
8125
8126 let pending = consumer.fetch_group(2, None);
8128 let req = dynamic
8129 .requested_group()
8130 .now_or_never()
8131 .expect("should not block")
8132 .unwrap();
8133 let mut group = req.accept(None).unwrap();
8134 group
8135 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"backfill"))
8136 .unwrap();
8137 group.finish().unwrap();
8138 let mut fetched = pending.await.unwrap();
8139 assert_eq!(&fetched.read_frame().await.unwrap().unwrap().payload[..], b"backfill");
8140 assert!(consumer.peek_group(2).is_some(), "backfill is cached for later fetches");
8141
8142 let mut subscriber = producer.subscribe(replay());
8144 assert_eq!(subscriber.assert_group().sequence, 5);
8145 assert_eq!(subscriber.assert_group().sequence, 6);
8146 subscriber.assert_no_group();
8147 }
8148
8149 #[tokio::test]
8152 async fn expired_backfill_reclaimed() {
8153 let (producer, pool) = pooled_producer(1 << 40);
8154 let dynamic = producer.dynamic();
8155 let consumer = producer.consume();
8156
8157 producer.create_group(5u64.into()).unwrap().finish().unwrap();
8158
8159 let pending = consumer.fetch_group(2, None);
8161 let req = dynamic
8162 .requested_group()
8163 .now_or_never()
8164 .expect("should not block")
8165 .unwrap();
8166 let mut group = req.accept(None).unwrap();
8167 group
8168 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 1000]))
8169 .unwrap();
8170 group.finish().unwrap();
8171 pending.await.unwrap();
8172 let used = pool.used();
8173
8174 crate::model::clock::advance(cache::DEFAULT_EXPIRY + Duration::from_secs(1));
8176 producer.create_group(6u64.into()).unwrap().finish().unwrap();
8177
8178 assert!(consumer.peek_group(2).is_none(), "expired backfill is reclaimed");
8179 assert!(pool.used() < used, "its bytes are released");
8180 }
8181
8182 #[tokio::test]
8183 async fn fetch_aborts_with_track() {
8184 let producer = track_producer("test", None);
8185 let dynamic = producer.dynamic();
8186 let consumer = producer.consume();
8187
8188 let pending = consumer.fetch_group(3, None);
8189 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
8190
8191 producer.abort(Error::Cancel).unwrap();
8192 assert!(pending.await.is_err());
8193 drop(dynamic);
8194 }
8195}