1use crate::{Error, Result, Timescale, Timestamp, coding};
17use crate::{broadcast, cache, frame, group, stats};
18
19use super::{Datagram, Requests};
20
21pub use super::subscription::Subscription;
22
23use std::{
24 collections::{BTreeMap, VecDeque},
25 sync::Arc,
26 sync::OnceLock,
27 sync::atomic::{AtomicBool, Ordering},
28 task::{Poll, ready},
29 time::Duration,
30};
31
32pub const DEFAULT_LATENCY_MAX: Duration = Duration::from_secs(5);
34
35const MAX_DATAGRAM_AGE: Duration = Duration::from_millis(50);
41
42const EVICT_SLACK: usize = 64;
45
46const EVICT_SCAN: usize = 4;
50
51#[derive(Clone, Debug)]
62#[non_exhaustive]
63pub struct Info {
64 pub timescale: Timescale,
71 pub latency_max: Duration,
80 pub priority: u8,
83 pub ordered: bool,
87}
88
89impl Default for Info {
90 fn default() -> Self {
91 Self {
92 timescale: Timescale::default(),
93 latency_max: DEFAULT_LATENCY_MAX,
94 priority: 0,
95 ordered: false,
96 }
97 }
98}
99
100impl Info {
101 pub fn with_timescale(mut self, timescale: Timescale) -> Self {
106 self.timescale = timescale;
107 self
108 }
109
110 pub fn with_latency_max(mut self, latency_max: Duration) -> Self {
112 self.latency_max = latency_max;
113 self
114 }
115
116 pub fn with_priority(mut self, priority: u8) -> Self {
118 self.priority = priority;
119 self
120 }
121
122 pub fn with_ordered(mut self, ordered: bool) -> Self {
126 self.ordered = ordered;
127 self
128 }
129}
130
131#[derive(Default)]
132pub(crate) struct TrackState {
133 info: Option<Info>,
136
137 broadcast: Arc<broadcast::Info>,
140
141 cache: Arc<cache::Track>,
145
146 lookup: BTreeMap<u64, Slot>,
154
155 arrival: VecDeque<(u64, u32)>,
160
161 evict: VecDeque<(u64, u32)>,
169
170 debt: u64,
175
176 datagrams: VecDeque<(Datagram, web_async::time::Instant)>,
180
181 datagram_offset: usize,
184
185 offset: usize,
188
189 max_sequence: Option<u64>,
192
193 latest_group: Option<u64>,
198
199 next_stamp: u32,
201
202 expire_cursor: usize,
205
206 final_sequence: Option<u64>,
208
209 abort: Option<Error>,
211
212 subscriptions: kio::Shared<Subscriptions>,
216
217 fetch: kio::Shared<FetchState>,
220}
221
222struct Slot {
228 group: group::Producer,
229
230 stamp: u32,
235}
236
237type Subscriptions = Vec<kio::Consumer<Subscription>>;
239
240type FetchState = Requests<u64, PendingFetch>;
245
246struct PendingFetch {
248 priority: u8,
250
251 result: kio::Producer<FetchOutcome>,
256}
257
258#[derive(Default)]
261struct FetchOutcome {
262 rejected: Option<Error>,
263}
264
265impl TrackState {
266 fn poll_info(&self) -> Poll<Result<Info>> {
267 if let Some(info) = &self.info {
268 Poll::Ready(Ok(info.clone()))
269 } else {
270 Poll::Pending
271 }
272 }
273
274 fn poll_recv_group(&self, index: usize, min_sequence: u64) -> Poll<Result<Option<(group::Consumer, usize)>>> {
278 let start = index.saturating_sub(self.offset);
279 for (i, (sequence, stamp)) in self.arrival.iter().enumerate().skip(start) {
280 if *sequence >= min_sequence
281 && let Some(slot) = self.lookup.get(sequence)
282 && slot.stamp == *stamp
283 && !slot.group.is_aborted()
284 {
285 slot.group.cache_refresh();
288 return Poll::Ready(Ok(Some((slot.group.consume(), self.offset + i))));
289 }
290 }
291
292 if self.is_complete() {
294 Poll::Ready(Ok(None))
295 } else if let Some(err) = &self.abort {
296 Poll::Ready(Err(err.clone()))
297 } else {
298 Poll::Pending
299 }
300 }
301
302 fn poll_recv_datagram(&self, index: usize) -> Poll<Result<Option<(Datagram, usize)>>> {
308 let start = index.saturating_sub(self.datagram_offset);
309 if let Some((datagram, _)) = self.datagrams.get(start) {
310 return Poll::Ready(Ok(Some((datagram.clone(), self.datagram_offset + start))));
311 }
312
313 if self.is_complete() {
315 Poll::Ready(Ok(None))
316 } else if let Some(err) = &self.abort {
317 Poll::Ready(Err(err.clone()))
318 } else {
319 Poll::Pending
320 }
321 }
322
323 fn push_datagram(&mut self, datagram: Datagram) {
325 let now = web_async::time::Instant::now();
326 self.datagrams.push_back((datagram, now));
327 while let Some((_, at)) = self.datagrams.front() {
328 if now.duration_since(*at) <= MAX_DATAGRAM_AGE {
329 break;
330 }
331 self.datagrams.pop_front();
332 self.datagram_offset += 1;
333 }
334 }
335
336 fn poll_read_frame(
340 &self,
341 index: usize,
342 next_sequence: u64,
343 waiter: &kio::Waiter,
344 ) -> Poll<Result<Option<(frame::Frame, usize, u64)>>> {
345 let start = index.saturating_sub(self.offset);
346 let mut pending_seen = false;
347 for (i, (sequence, stamp)) in self.arrival.iter().enumerate().skip(start) {
348 if *sequence < next_sequence {
349 continue;
350 }
351 let Some(slot) = self.lookup.get(sequence) else {
352 continue;
353 };
354 if slot.stamp != *stamp {
355 continue;
358 }
359
360 let mut consumer = slot.group.consume();
361 match consumer.poll_read_frame(waiter) {
362 Poll::Ready(Ok(Some(frame))) => {
363 return Poll::Ready(Ok(Some((frame, self.offset + i, *sequence))));
364 }
365 Poll::Ready(Ok(None)) => continue,
366 Poll::Ready(Err(_)) => continue,
369 Poll::Pending => {
370 pending_seen = true;
371 continue;
372 }
373 }
374 }
375
376 if pending_seen {
379 Poll::Pending
380 } else if self.is_complete() {
381 Poll::Ready(Ok(None))
382 } else if let Some(err) = &self.abort {
383 Poll::Ready(Err(err.clone()))
384 } else {
385 Poll::Pending
386 }
387 }
388
389 fn poll_next_in_range(
399 &self,
400 next_sequence: u64,
401 end_sequence: Option<u64>,
402 ) -> Poll<Result<Option<group::Consumer>>> {
403 if let Some(end) = end_sequence
407 && end < next_sequence
408 {
409 if let Some(err) = &self.abort {
410 return Poll::Ready(Err(err.clone()));
411 }
412 return Poll::Pending;
413 }
414
415 let best = self
418 .lookup
419 .range(next_sequence..)
420 .map(|(_, slot)| &slot.group)
421 .take_while(|group| end_sequence.is_none_or(|end| group.sequence <= end))
422 .find(|group| !group.is_aborted());
423
424 if let Some(group) = best {
425 group.cache_refresh();
427 return Poll::Ready(Ok(Some(group.consume())));
428 }
429
430 if let Some(err) = &self.abort {
432 return Poll::Ready(Err(err.clone()));
433 }
434 if let Some(fin) = self.final_sequence
437 && next_sequence >= fin
438 {
439 return Poll::Ready(Ok(None));
440 }
441 Poll::Pending
442 }
443
444 fn latency_bound(&self) -> Option<Duration> {
447 self.info.as_ref().map(|info| info.latency_max)
448 }
449
450 fn poll_fetch_cached(&self, sequence: u64) -> Poll<Result<group::Consumer>> {
455 if let Some(slot) = self.lookup.get(&sequence)
456 && !slot.group.is_aborted()
457 {
458 slot.group.cache_refresh();
462 return Poll::Ready(Ok(slot.group.consume()));
463 }
464
465 if let Some(err) = &self.abort {
466 return Poll::Ready(Err(err.clone()));
467 }
468
469 if self.final_sequence.is_some_and(|fin| sequence >= fin) {
471 return Poll::Ready(Err(Error::NotFound));
472 }
473
474 Poll::Pending
475 }
476
477 fn evict_expired(&mut self, max_age: Duration) {
486 let now = self.cache.pool().now();
487 let max_ticks = cache::Pool::ticks(max_age);
488
489 let len = self.evict.len();
490 if len > 0 {
491 let start = self.expire_cursor % len;
492 for step in 0..len.min(EVICT_SCAN) {
493 let (sequence, stamp) = self.evict[(start + step) % len];
494 let Some(slot) = self.lookup.get(&sequence) else {
495 continue;
496 };
497 if slot.stamp != stamp {
498 continue;
500 }
501 if slot.group.is_aborted() {
504 self.lookup.remove(&sequence);
505 continue;
506 }
507 if Some(sequence) == self.latest_group || now.saturating_sub(slot.group.cache_accessed()) <= max_ticks {
508 continue;
509 }
510 let slot = self.lookup.remove(&sequence).unwrap();
514 let _ = slot.group.abort(Error::Old);
515 }
516 self.expire_cursor = (start + EVICT_SCAN) % len;
517 }
518
519 while let Some((sequence, stamp)) = self.arrival.front() {
522 if self.lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp) {
523 break;
524 }
525 self.arrival.pop_front();
526 self.offset += 1;
527 }
528
529 while let Some((sequence, stamp)) = self.evict.front() {
531 if self.lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp) {
532 break;
533 }
534 self.evict.pop_front();
535 }
536
537 if self.evict.len() > 2 * self.lookup.len() + EVICT_SLACK {
540 let lookup = &self.lookup;
541 self.evict
542 .retain(|(sequence, stamp)| lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp));
543 }
544 }
545
546 fn clear_cache(&mut self) {
549 self.lookup.clear();
550 self.arrival.clear();
551 self.evict.clear();
552 self.latest_group = None;
553 self.debt = 0;
554 }
555
556 fn install(&mut self, mut info: Info) {
562 info.latency_max = info.latency_max.min(self.broadcast.origin.cache_duration);
563 self.info = Some(info);
564 }
565
566 fn spawn(broadcast: Arc<broadcast::Info>) -> kio::Producer<Self> {
573 let state = kio::Producer::new(Self {
574 broadcast: broadcast.clone(),
575 ..Default::default()
576 });
577 let cache = cache::Track::new(broadcast.origin.pool.clone(), state.downgrade());
578 state.write().ok().expect("a new track is open").cache = cache;
579 state
580 }
581
582 fn claim_sequence(&mut self, sequence: u64) -> Result<()> {
588 if let Some(slot) = self.lookup.get(&sequence) {
589 if !slot.group.is_aborted() {
590 return Err(Error::Duplicate);
591 }
592 self.lookup.remove(&sequence);
593 }
594 Ok(())
595 }
596
597 fn insert_group(&mut self, group: &group::Producer, visible: bool) {
604 let sequence = group.sequence;
605 self.next_stamp = self.next_stamp.wrapping_add(1);
606 let stamp = self.next_stamp;
607
608 if self.latest_group.is_none_or(|latest| sequence >= latest) {
612 if let Some(latest) = self.latest_group
615 && sequence > latest
616 && let Some(prev) = self.lookup.get(&latest)
617 {
618 prev.group.cache_demote();
619 self.evict.push_back((latest, prev.stamp));
620 }
621 self.latest_group = Some(sequence);
622 } else {
623 group.cache_demote();
624 self.evict.push_back((sequence, stamp));
625 }
626
627 self.max_sequence = Some(self.max_sequence.map_or(sequence, |max| max.max(sequence)));
628 self.lookup.insert(
629 sequence,
630 Slot {
631 group: group.clone(),
632 stamp,
633 },
634 );
635 if visible {
636 self.arrival.push_back((sequence, stamp));
637 }
638 }
639
640 fn commit_group(&mut self, group: &group::Producer, visible: bool, latency_max: Duration) {
644 self.charge_debt();
645 self.insert_group(group, visible);
646 self.evict_expired(latency_max);
647 }
648
649 pub(super) fn charge_debt(&mut self) {
662 let written = self.cache.take_written();
663 let pool = self.cache.pool().clone();
664 match pool.accrue(written) {
665 Some(mut accrued) => {
666 if self.oldest_is_stale(&pool) {
667 accrued = accrued.saturating_mul(2);
668 }
669 self.debt = self.debt.saturating_add(accrued).min(pool.used());
672 self.pay_debt(&pool, written.saturating_mul(2));
675 }
676 None => self.debt = 0,
679 }
680 }
681
682 fn oldest_is_stale(&self, pool: &cache::Pool) -> bool {
686 let Some(average) = pool.average() else {
687 return false;
688 };
689 let Some((sequence, stamp)) = self.evict.front() else {
690 return false;
691 };
692 let Some(slot) = self.lookup.get(sequence) else {
693 return false;
694 };
695 slot.stamp == *stamp && !slot.group.is_aborted() && slot.group.cache_accessed() <= average
696 }
697
698 fn pay_debt(&mut self, pool: &cache::Pool, cap: u64) {
711 let average = pool.average().unwrap_or(0);
712 let mut paid = 0u64;
713 let mut scanned = 0usize;
714 for _ in 0..self.evict.len() {
715 if self.debt == 0 || paid >= cap || scanned >= EVICT_SCAN {
716 return;
717 }
718 let Some((sequence, stamp)) = self.evict.pop_front() else {
719 return;
720 };
721 let Some(slot) = self.lookup.get(&sequence) else {
722 continue;
724 };
725 if slot.stamp != stamp {
726 continue;
728 }
729 if slot.group.is_aborted() {
730 self.lookup.remove(&sequence);
732 continue;
733 }
734 if Some(sequence) == self.latest_group {
735 self.evict.push_back((sequence, stamp));
737 continue;
738 }
739
740 scanned += 1;
741 if slot.group.cache_accessed() > average {
745 self.evict.push_back((sequence, stamp));
746 continue;
747 }
748 let size = slot.group.cache_size();
751 if size > self.debt {
752 self.evict.push_front((sequence, stamp));
753 return;
754 }
755
756 self.debt -= size;
757 paid = paid.saturating_add(size);
758 let slot = self.lookup.remove(&sequence).unwrap();
759 let _ = slot.group.abort(Error::Evicted);
760 }
761 }
762
763 fn set_final(&mut self, final_sequence: u64) -> Result<()> {
766 if self.final_sequence.is_some() {
767 return Err(Error::Closed);
768 }
769 if let Some(max) = self.max_sequence
770 && final_sequence <= max
771 {
772 return Err(Error::ProtocolViolation);
773 }
774 self.final_sequence = Some(final_sequence);
775 Ok(())
776 }
777
778 fn is_complete(&self) -> bool {
784 self.final_sequence
785 .is_some_and(|fin| self.max_sequence.map_or(0, |max| max.saturating_add(1)) >= fin)
786 }
787
788 fn poll_finished(&self) -> Poll<Result<u64>> {
789 if let Some(fin) = self.final_sequence {
790 Poll::Ready(Ok(fin))
791 } else if let Some(err) = &self.abort {
792 Poll::Ready(Err(err.clone()))
793 } else {
794 Poll::Pending
795 }
796 }
797
798 fn modify(producer: &kio::Producer<Self>) -> Result<kio::Mut<'_, Self>> {
799 producer.write().map_err(|r| r.abort.clone().unwrap_or(Error::Dropped))
800 }
801
802 fn insert_group_request(&mut self, sequence: u64, info: Option<Info>) -> Result<group::Producer> {
808 if let Some(err) = &self.abort {
809 return Err(err.clone());
810 }
811 if let Some(fin) = self.final_sequence
812 && sequence >= fin
813 {
814 return Err(Error::Closed);
815 }
816
817 if self.info.is_none() {
821 self.install(info.unwrap_or_default());
822 }
823 let info = self.info.clone().unwrap();
824
825 self.claim_sequence(sequence)?;
827
828 let latency_max = info.latency_max;
829 let group = group::Producer::new(group::Info { sequence }, info, self.cache.clone());
830 group.cache_refresh();
835 self.commit_group(&group, false, latency_max);
836 Ok(group)
837 }
838}
839
840#[derive(Clone)]
842pub struct Producer {
843 name: Arc<str>,
844 broadcast: Arc<broadcast::Info>,
847 state: kio::Producer<TrackState>,
848 prev_subscription: Option<Subscription>,
849 alive: Arc<Alive>,
851 stats: stats::Scope,
855}
856
857impl Producer {
858 pub(crate) fn new(
866 broadcast: Arc<broadcast::Info>,
867 name: impl Into<Arc<str>>,
868 info: impl Into<Option<Info>>,
869 ) -> Self {
870 let name = name.into();
871 let state = TrackState::spawn(broadcast.clone());
872 state
873 .write()
874 .ok()
875 .expect("a new track is open")
876 .install(info.into().unwrap_or_default());
877 let alive = Alive::new(name.clone(), state.clone());
878 alive.publish(None);
879 Self {
880 name,
881 state,
882 broadcast,
883 prev_subscription: None,
884 alive,
885 stats: stats::Scope::default(),
886 }
887 }
888
889 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
893 self.alive.publish(Some(&scope));
894 self.stats = scope;
895 self
896 }
897
898 pub fn name(&self) -> &str {
900 &self.name
901 }
902
903 pub fn broadcast(&self) -> &broadcast::Info {
905 &self.broadcast
906 }
907
908 pub fn create_group(&mut self, group: group::Info) -> Result<group::Producer> {
910 let mut state = self.modify()?;
911 if let Some(fin) = state.final_sequence
912 && group.sequence >= fin
913 {
914 return Err(Error::Closed);
915 }
916 let track = state.info.clone().unwrap();
917 let latency_max = track.latency_max;
918
919 state.claim_sequence(group.sequence)?;
921
922 let group = group::Producer::new(group, track, state.cache.clone()).with_meter(self.stats.meter());
923 state.commit_group(&group, true, latency_max);
924
925 Ok(group)
926 }
927
928 pub fn append_group(&mut self) -> Result<group::Producer> {
930 let mut state = self.modify()?;
931 let sequence = match state.max_sequence {
932 Some(s) => s.checked_add(1).ok_or(coding::BoundsExceeded)?,
933 None => 0,
934 };
935 if let Some(fin) = state.final_sequence
936 && sequence >= fin
937 {
938 return Err(Error::Closed);
939 }
940
941 let track = state.info.clone().unwrap();
942 let latency_max = track.latency_max;
943
944 let group =
945 group::Producer::new(group::Info { sequence }, track, state.cache.clone()).with_meter(self.stats.meter());
946 state.commit_group(&group, true, latency_max);
947
948 Ok(group)
949 }
950
951 pub fn append_datagram<B: crate::IntoBytes>(&mut self, timestamp: Timestamp, payload: B) -> Result<u64> {
963 let payload = payload.into_bytes();
964 if payload.len() > super::datagram::MAX_DATAGRAM_PAYLOAD {
965 return Err(Error::FrameTooLarge);
966 }
967 let meter = self.stats.meter();
969 let mut state = self.modify()?;
970 let timescale = state.info.as_ref().unwrap().timescale;
972 let timestamp = timestamp.convert(timescale).map_err(|_| Error::TimestampMismatch)?;
973 let sequence = match state.max_sequence {
974 Some(s) => s.checked_add(1).ok_or(coding::BoundsExceeded)?,
975 None => 0,
976 };
977 if let Some(fin) = state.final_sequence
978 && sequence >= fin
979 {
980 return Err(Error::Closed);
981 }
982 state.max_sequence = Some(sequence);
983 meter.datagram(payload.len() as u64);
984 state.push_datagram(Datagram {
985 sequence,
986 timestamp,
987 payload,
988 });
989 Ok(sequence)
990 }
991
992 pub fn write_datagram(&mut self, mut datagram: Datagram) -> Result<()> {
998 if datagram.payload.len() > super::datagram::MAX_DATAGRAM_PAYLOAD {
999 return Err(Error::FrameTooLarge);
1000 }
1001 let meter = self.stats.meter();
1003 let mut state = self.modify()?;
1004 let timescale = state.info.as_ref().unwrap().timescale;
1006 datagram.timestamp = datagram
1007 .timestamp
1008 .convert(timescale)
1009 .map_err(|_| Error::TimestampMismatch)?;
1010 if let Some(fin) = state.final_sequence
1011 && datagram.sequence >= fin
1012 {
1013 return Err(Error::Closed);
1014 }
1015 state.max_sequence = Some(state.max_sequence.unwrap_or(0).max(datagram.sequence));
1016 meter.datagram(datagram.payload.len() as u64);
1017 state.push_datagram(datagram);
1018 Ok(())
1019 }
1020
1021 pub fn write_frame<B: crate::IntoBytes>(&mut self, timestamp: Timestamp, frame: B) -> Result<()> {
1026 let frame = crate::IntoBytes::into_bytes(frame);
1027 if frame.len() as u64 > group::MAX_CACHE_BYTES {
1028 return Err(Error::FrameTooLarge);
1029 }
1030 let mut group = self.append_group()?;
1031 group.write_frame(timestamp, frame)?;
1032 group.finish()?;
1033 Ok(())
1034 }
1035
1036 pub fn finish(&mut self) -> Result<()> {
1042 let mut state = self.modify()?;
1043 let final_sequence = match state.max_sequence {
1044 Some(max) => max.checked_add(1).ok_or(coding::BoundsExceeded)?,
1045 None => 0,
1046 };
1047 state.set_final(final_sequence)
1048 }
1049
1050 pub fn finish_at(&mut self, final_sequence: u64) -> Result<()> {
1063 self.modify()?.set_final(final_sequence)
1064 }
1065
1066 pub fn final_sequence(&self) -> Option<u64> {
1071 self.state.read().final_sequence
1072 }
1073
1074 pub fn abort(self, err: Error) -> Result<()> {
1085 let mut guard = self.modify()?;
1086 guard.abort = Some(err);
1087 guard.clear_cache();
1088 guard.datagrams.clear();
1089 guard.close();
1090 Ok(())
1091 }
1092
1093 pub async fn unused(&self) -> Result<()> {
1095 self.state.unused().await.map_err(|_| self.abort_reason())
1096 }
1097
1098 pub async fn used(&self) -> Result<()> {
1100 self.state.used().await.map_err(|_| self.abort_reason())
1101 }
1102
1103 pub async fn closed(&self) -> Error {
1105 kio::wait(|waiter| self.poll_closed(waiter)).await
1106 }
1107
1108 pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<Error> {
1110 self.state.poll_closed(waiter).map(|()| self.abort_reason())
1111 }
1112
1113 fn abort_reason(&self) -> Error {
1115 self.state.read().abort.clone().unwrap_or(Error::Dropped)
1116 }
1117
1118 pub fn is_closed(&self) -> bool {
1120 self.state.read().is_closed()
1121 }
1122
1123 pub fn latest(&self) -> Option<u64> {
1125 self.state.read().max_sequence
1126 }
1127
1128 pub fn is_clone(&self, other: &Self) -> bool {
1130 self.state.same_channel(&other.state)
1131 }
1132
1133 pub(crate) fn weak(&self) -> TrackWeak {
1135 TrackWeak {
1136 name: self.name.clone(),
1137 state: self.state.weak(),
1138 }
1139 }
1140
1141 pub fn demand(&self) -> Demand {
1149 Demand {
1150 name: self.name.clone(),
1151 state: self.state.weak(),
1152 }
1153 }
1154
1155 pub fn consume(&self) -> Consumer {
1160 Consumer::plain(self.name.clone(), self.state.consume())
1161 }
1162
1163 pub fn subscribe(&self, subscription: impl Into<Option<Subscription>>) -> Subscriber {
1168 let preferences = subscription.into().unwrap_or_default();
1169
1170 let info = self.state.read().info.clone().expect("producer always has info");
1175 let subscription = kio::Producer::new(preferences);
1176 register_subscription(self.state.read(), &subscription);
1177
1178 Subscriber {
1179 name: self.name.clone(),
1180 info,
1181 inner: SubscriberKind::Plain(PlainSubscriber {
1182 state: self.state.consume(),
1183 subscription,
1184 index: 0,
1185 datagram_index: 0,
1186 min_sequence: 0,
1187 next_sequence: 0,
1188 end_sequence: None,
1189 parked: BTreeMap::new(),
1190 }),
1191 stats: stats::Scope::default(),
1193 _stats_sub: stats::Subscription::default(),
1194 }
1195 }
1196
1197 pub async fn subscription_changed(&mut self) -> Result<Option<Subscription>> {
1203 kio::wait(|waiter| self.poll_subscription_changed(waiter)).await
1204 }
1205
1206 pub fn subscription(&self) -> Option<Subscription> {
1214 let state = self.state.read();
1215 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
1216 drop(state);
1217 snapshot_subscription(&subs, bound)
1218 }
1219
1220 pub fn poll_subscription_changed(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Subscription>>> {
1224 if self.state.poll_closed(waiter).is_ready() {
1227 let abort = self.state.read().abort.clone();
1228 return Poll::Ready(Err(abort.unwrap_or(Error::Dropped)));
1229 }
1230
1231 let state = self.state.read();
1233 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
1234 drop(state);
1235
1236 let prev = &self.prev_subscription;
1237 let mut combined = None;
1238 let mut guard = ready!(subs.poll(waiter, |subs| {
1239 let next = combined_subscription(subs, bound, waiter);
1240 if &next == prev {
1241 Poll::Pending
1242 } else {
1243 combined = next;
1244 Poll::Ready(())
1245 }
1246 }));
1247 guard.retain(|sub| !sub.is_closed());
1249 drop(guard);
1250 self.prev_subscription = combined.clone();
1251 Poll::Ready(Ok(combined))
1252 }
1253
1254 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
1256 self.state.poll_unused(waiter).map(|_| ())
1257 }
1258
1259 pub fn dynamic(&self) -> Dynamic {
1263 Dynamic::new(self.name.clone(), self.state.clone(), self.alive.clone())
1264 }
1265
1266 fn modify(&self) -> Result<kio::Mut<'_, TrackState>> {
1267 TrackState::modify(&self.state)
1268 }
1269}
1270
1271fn poll_requested_group(
1275 state: &kio::Producer<TrackState>,
1276 fetch: &kio::Shared<FetchState>,
1277 waiter: &kio::Waiter,
1278) -> Poll<Result<GroupRequest>> {
1279 if let Poll::Ready(mut guard) = fetch.poll(waiter, |fetch| {
1281 if fetch.has_queued() {
1282 Poll::Ready(())
1283 } else {
1284 Poll::Pending
1285 }
1286 }) {
1287 let sequence = guard.pop().expect("predicate guaranteed a request");
1288 let pending = guard.get(&sequence).expect("popped key must be pending");
1292 let priority = pending.priority;
1293 let result = pending.result.clone();
1294 drop(guard);
1295 return Poll::Ready(Ok(GroupRequest {
1296 state: state.clone(),
1297 fetch: fetch.clone(),
1298 sequence,
1299 priority,
1300 result,
1301 done: false,
1302 }));
1303 }
1304
1305 match state.poll_ref(waiter, |state| match &state.abort {
1307 Some(err) => Poll::Ready(err.clone()),
1308 None => Poll::Pending,
1309 }) {
1310 Poll::Ready(Ok(err)) => Poll::Ready(Err(err)),
1311 Poll::Ready(Err(closed)) => Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped))),
1312 Poll::Pending => Poll::Pending,
1313 }
1314}
1315
1316pub struct Dynamic {
1326 name: Arc<str>,
1327 state: kio::Producer<TrackState>,
1329 fetch: kio::Shared<FetchState>,
1331 alive: Arc<Alive>,
1334}
1335
1336impl Dynamic {
1337 fn new(name: Arc<str>, state: kio::Producer<TrackState>, alive: Arc<Alive>) -> Self {
1338 let fetch = state.read().fetch.clone();
1339 fetch.lock().add_handler();
1340 Self {
1341 name,
1342 state,
1343 fetch,
1344 alive,
1345 }
1346 }
1347
1348 pub fn name(&self) -> &str {
1350 &self.name
1351 }
1352
1353 pub async fn requested_group(&self) -> Result<GroupRequest> {
1359 kio::wait(|waiter| self.poll_requested_group(waiter)).await
1360 }
1361
1362 pub fn poll_requested_group(&self, waiter: &kio::Waiter) -> Poll<Result<GroupRequest>> {
1364 poll_requested_group(&self.state, &self.fetch, waiter)
1365 }
1366
1367 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
1369 self.state.poll_unused(waiter).map(|_| ())
1370 }
1371}
1372
1373impl Clone for Dynamic {
1374 fn clone(&self) -> Self {
1375 self.fetch.lock().add_handler();
1377 Self {
1378 name: self.name.clone(),
1379 state: self.state.clone(),
1380 fetch: self.fetch.clone(),
1381 alive: self.alive.clone(),
1382 }
1383 }
1384}
1385
1386impl Drop for Dynamic {
1387 fn drop(&mut self) {
1388 let mut fetch = self.fetch.lock();
1394 if fetch.remove_handler() {
1395 fetch.drain_queued();
1396 }
1397 }
1398}
1399
1400struct Alive {
1409 name: Arc<str>,
1410 state: kio::Producer<TrackState>,
1411
1412 published: AtomicBool,
1415
1416 stats: OnceLock<stats::Subscription>,
1419}
1420
1421impl Alive {
1422 fn new(name: Arc<str>, state: kio::Producer<TrackState>) -> Arc<Self> {
1423 Arc::new(Self {
1424 name,
1425 state,
1426 published: Default::default(),
1427 stats: Default::default(),
1428 })
1429 }
1430
1431 fn publish(&self, stats: Option<&stats::Scope>) {
1435 self.published.store(true, Ordering::Relaxed);
1436 if let Some(scope) = stats {
1437 let _ = self.stats.set(scope.subscribe());
1440 }
1441 }
1442}
1443
1444impl Drop for Alive {
1445 fn drop(&mut self) {
1446 if !self.published.load(Ordering::Relaxed) {
1448 return;
1449 }
1450 match self.state.write() {
1458 Ok(mut state) => {
1459 if state.final_sequence.is_some() || state.abort.is_some() {
1460 return;
1461 }
1462 tracing::warn!(
1463 track = %self.name,
1464 "track::Producer dropped without finish() or abort()"
1465 );
1466 state.clear_cache();
1467 state.datagrams.clear();
1468 }
1469 Err(state) => {
1470 if state.final_sequence.is_some() || state.abort.is_some() {
1471 return;
1472 }
1473 tracing::warn!(
1474 track = %self.name,
1475 "track::Producer dropped without finish() or abort()"
1476 );
1477 }
1478 }
1479 }
1480}
1481
1482fn combined_subscription(subs: &Subscriptions, bound: Option<Duration>, waiter: &kio::Waiter) -> Option<Subscription> {
1488 let mut combined = None;
1489 for sub in subs.iter() {
1490 if sub.is_closed() {
1495 continue;
1496 }
1497 let _ = sub.poll_closed(waiter);
1502 if let Poll::Ready(Ok(sub)) = sub.poll(waiter, |sub| sub.poll_combined(&combined)) {
1503 combined = Some(sub);
1504 }
1505 }
1506 clamp_combined(combined, bound)
1507}
1508
1509fn snapshot_subscription(subs: &kio::Shared<Subscriptions>, bound: Option<Duration>) -> Option<Subscription> {
1511 let mut combined: Option<Subscription> = None;
1512 for sub in subs.read().iter() {
1513 if sub.is_closed() {
1515 continue;
1516 }
1517 if let Poll::Ready(merged) = sub.read().poll_combined(&combined) {
1518 combined = Some(merged);
1519 }
1520 }
1521 clamp_combined(combined, bound)
1522}
1523
1524fn clamp_combined(combined: Option<Subscription>, bound: Option<Duration>) -> Option<Subscription> {
1532 let mut combined = combined?;
1533 if let Some(bound) = bound {
1534 combined.latency_max = combined.latency_max.min(bound);
1535 }
1536 Some(combined)
1537}
1538
1539fn register_subscription(state: kio::Ref<'_, TrackState>, subscription: &kio::Producer<Subscription>) {
1543 if state.is_closed() {
1544 return;
1545 }
1546 let subs = state.subscriptions.clone();
1547 drop(state);
1548 subs.lock().push(subscription.consume());
1549}
1550
1551#[derive(Clone)]
1553pub(crate) struct TrackWeak {
1554 name: Arc<str>,
1555 state: kio::ProducerWeak<TrackState>,
1556}
1557
1558impl TrackWeak {
1559 pub fn consume(&self) -> Consumer {
1560 Consumer::plain(self.name.clone(), self.state.consume())
1561 }
1562
1563 pub(crate) fn name(&self) -> &Arc<str> {
1566 &self.name
1567 }
1568
1569 pub(crate) fn is_used(&self) -> bool {
1572 !self.state.is_closed() && self.state.is_used()
1573 }
1574
1575 pub(crate) fn poll_used(&self, waiter: &kio::Waiter) {
1578 let _ = self.state.poll_used(waiter);
1579 }
1580
1581 pub(crate) fn poll_unused(&self, waiter: &kio::Waiter) {
1584 let _ = self.state.poll_unused(waiter);
1585 }
1586}
1587
1588impl super::WeakEntry for TrackWeak {
1589 fn is_closed(&self) -> bool {
1590 self.state.is_closed()
1591 }
1592
1593 fn same_channel(&self, other: &Self) -> bool {
1594 self.state.same_channel(&other.state)
1595 }
1596}
1597
1598#[derive(Clone)]
1607pub struct Demand {
1608 name: Arc<str>,
1609 state: kio::ProducerWeak<TrackState>,
1610}
1611
1612impl Demand {
1613 pub fn name(&self) -> &str {
1615 &self.name
1616 }
1617
1618 pub async fn used(&self) -> Result<()> {
1620 self.state.used().await.map_err(|_| self.abort_reason())
1621 }
1622
1623 pub async fn unused(&self) -> Result<()> {
1625 self.state.unused().await.map_err(|_| self.abort_reason())
1626 }
1627
1628 pub async fn closed(&self) -> Error {
1630 self.state.closed().await;
1631 self.abort_reason()
1632 }
1633
1634 fn abort_reason(&self) -> Error {
1636 self.state.read().abort.clone().unwrap_or(Error::Dropped)
1637 }
1638}
1639
1640#[derive(Clone)]
1651pub struct Consumer {
1652 name: Arc<str>,
1653 inner: ConsumerKind,
1654 stats: stats::Scope,
1657}
1658
1659#[derive(Clone)]
1660enum ConsumerKind {
1661 Plain(kio::Consumer<TrackState>),
1662 Spliced(super::resume::Consumer),
1663}
1664
1665impl Consumer {
1666 fn plain(name: Arc<str>, state: kio::Consumer<TrackState>) -> Self {
1667 Self {
1668 name,
1669 inner: ConsumerKind::Plain(state),
1670 stats: stats::Scope::default(),
1671 }
1672 }
1673
1674 pub(crate) fn spliced(name: Arc<str>, resume: super::resume::Consumer) -> Self {
1676 Self {
1677 name,
1678 inner: ConsumerKind::Spliced(resume),
1679 stats: stats::Scope::default(),
1680 }
1681 }
1682
1683 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
1686 self.stats = scope;
1687 self
1688 }
1689
1690 pub fn name(&self) -> &str {
1692 &self.name
1693 }
1694
1695 pub fn subscribe(&self, subscription: impl Into<Option<Subscription>>) -> kio::Pending<Subscribing> {
1701 let subscription = kio::Producer::new(subscription.into().unwrap_or_default());
1702
1703 let inner = match &self.inner {
1704 ConsumerKind::Plain(state) => {
1705 register_subscription(state.read(), &subscription);
1708 SubscribingKind::Plain(state.clone())
1709 }
1710 ConsumerKind::Spliced(resume) => SubscribingKind::Spliced(resume.clone()),
1712 };
1713
1714 kio::Pending::new(Subscribing {
1715 name: self.name.clone(),
1716 inner,
1717 subscription,
1718 stats: self.stats.clone(),
1719 })
1720 }
1721
1722 pub(crate) fn peek_latest(&self) -> Option<group::Consumer> {
1726 match &self.inner {
1727 ConsumerKind::Plain(state) => {
1728 let sequence = state.read().max_sequence?;
1729 self.peek_group(sequence)
1730 }
1731 ConsumerKind::Spliced(resume) => resume.peek_latest(),
1732 }
1733 }
1734
1735 pub(crate) fn peek_before(&self, sequence: u64) -> Option<group::Consumer> {
1739 match &self.inner {
1740 ConsumerKind::Plain(state) => {
1741 let state = state.read();
1742 state
1743 .lookup
1744 .range(..sequence)
1745 .rev()
1746 .map(|(_, slot)| &slot.group)
1747 .find(|group| !group.is_aborted())
1748 .map(|group| group.consume())
1749 }
1750 ConsumerKind::Spliced(resume) => resume.peek_before(sequence),
1751 }
1752 }
1753
1754 pub(crate) fn peek_group(&self, sequence: u64) -> Option<group::Consumer> {
1759 match &self.inner {
1760 ConsumerKind::Plain(state) => {
1761 let state = state.read();
1762 let slot = state.lookup.get(&sequence)?;
1763 if slot.group.is_aborted() {
1764 return None;
1765 }
1766 Some(slot.group.consume())
1767 }
1768 ConsumerKind::Spliced(resume) => resume.peek_group(sequence),
1769 }
1770 }
1771
1772 pub fn fetch_group(&self, sequence: u64, options: impl Into<Option<group::Fetch>>) -> kio::Pending<Fetching> {
1784 let options = options.into().unwrap_or_default();
1785
1786 self.stats.fetch();
1790
1791 let state = match &self.inner {
1792 ConsumerKind::Plain(state) => state,
1793 ConsumerKind::Spliced(resume) => {
1796 return kio::Pending::new(Fetching {
1797 inner: FetchingKind::Spliced(resume.fetch_group(sequence, options)),
1798 stats: self.stats.clone(),
1799 });
1800 }
1801 };
1802
1803 let mut result = None;
1804
1805 let (fetch, unresolved) = {
1809 let state = state.read();
1810 (state.fetch.clone(), state.poll_fetch_cached(sequence).is_pending())
1811 };
1812
1813 if unresolved {
1814 let mut fetch = fetch.lock();
1815 if let Some(pending) = fetch.join(&sequence) {
1816 pending.priority = pending.priority.max(options.priority);
1819 result = Some(pending.result.consume());
1820 } else {
1821 let producer = kio::Producer::<FetchOutcome>::default();
1825 let consumer = producer.consume();
1826 let attempt = PendingFetch {
1827 priority: options.priority,
1828 result: producer,
1829 };
1830 if fetch.insert(sequence, attempt).is_ok() {
1831 result = Some(consumer);
1832 }
1833 }
1834 }
1835
1836 kio::Pending::new(Fetching {
1837 inner: FetchingKind::Plain {
1838 state: state.clone(),
1839 fetch,
1840 sequence,
1841 result,
1842 },
1843 stats: self.stats.clone(),
1844 })
1845 }
1846
1847 pub fn info(&self) -> kio::Pending<Querying> {
1854 kio::Pending::new(Querying {
1855 inner: match &self.inner {
1856 ConsumerKind::Plain(state) => QueryingKind::Plain(state.clone()),
1857 ConsumerKind::Spliced(resume) => QueryingKind::Spliced(resume.clone()),
1858 },
1859 })
1860 }
1861
1862 pub fn latest(&self) -> Option<u64> {
1864 match &self.inner {
1865 ConsumerKind::Plain(state) => state.read().max_sequence,
1866 ConsumerKind::Spliced(resume) => resume.latest(),
1867 }
1868 }
1869
1870 pub(crate) fn poll_complete(&self, waiter: &kio::Waiter) -> Poll<Result<()>> {
1875 let ConsumerKind::Plain(state) = &self.inner else {
1876 return Poll::Pending;
1878 };
1879 match ready!(state.poll(waiter, |state| {
1880 if state.is_complete() {
1881 Poll::Ready(())
1882 } else {
1883 Poll::Pending
1884 }
1885 })) {
1886 Ok(_) => Poll::Ready(Ok(())),
1887 Err(closed) => Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped))),
1890 }
1891 }
1892}
1893
1894pub struct Subscribing {
1897 name: Arc<str>,
1898 inner: SubscribingKind,
1899 subscription: kio::Producer<Subscription>,
1900 stats: stats::Scope,
1901}
1902
1903enum SubscribingKind {
1904 Plain(kio::Consumer<TrackState>),
1905 Spliced(super::resume::Consumer),
1906}
1907
1908impl Subscribing {
1909 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<Subscriber>> {
1912 match &self.inner {
1913 SubscribingKind::Plain(state) => {
1914 let info = ready!(state.poll(waiter, |state| state.poll_info()))
1916 .map_err(|e| e.abort.clone().unwrap_or(Error::Dropped))??;
1917
1918 Poll::Ready(Ok(Subscriber {
1919 name: self.name.clone(),
1920 info,
1921 inner: SubscriberKind::Plain(PlainSubscriber {
1922 state: state.clone(),
1923 subscription: self.subscription.clone(),
1924 index: 0,
1925 datagram_index: 0,
1926 min_sequence: 0,
1927 next_sequence: 0,
1928 end_sequence: None,
1929 parked: BTreeMap::new(),
1930 }),
1931 stats: self.stats.clone(),
1932 _stats_sub: self.stats.subscribe(),
1933 }))
1934 }
1935 SubscribingKind::Spliced(resume) => {
1936 let info = ready!(resume.poll_info(waiter))?;
1939
1940 Poll::Ready(Ok(Subscriber {
1941 name: self.name.clone(),
1942 info,
1943 inner: SubscriberKind::Spliced(Box::new(resume.subscribe_shared(self.subscription.clone()))),
1944 stats: self.stats.clone(),
1945 _stats_sub: self.stats.subscribe(),
1946 }))
1947 }
1948 }
1949 }
1950
1951 pub fn update(&mut self, subscription: Subscription) -> Result<()> {
1956 let mut state = self.subscription.write().map_err(|_| Error::Closed)?;
1957 *state = subscription;
1958 Ok(())
1959 }
1960}
1961
1962impl kio::Pollable for Subscribing {
1963 type Output = Result<Subscriber>;
1964
1965 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
1966 self.poll_ok(waiter)
1967 }
1968}
1969
1970pub struct Querying {
1973 inner: QueryingKind,
1974}
1975
1976enum QueryingKind {
1977 Plain(kio::Consumer<TrackState>),
1978 Spliced(super::resume::Consumer),
1979}
1980
1981impl Querying {
1982 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<Info>> {
1984 match &self.inner {
1985 QueryingKind::Plain(state) => {
1986 let info = ready!(state.poll(waiter, |state| state.poll_info()))
1988 .map_err(|e| e.abort.clone().unwrap_or(Error::Dropped))??;
1989 Poll::Ready(Ok(info))
1990 }
1991 QueryingKind::Spliced(resume) => resume.poll_info(waiter),
1992 }
1993 }
1994}
1995
1996impl kio::Pollable for Querying {
1997 type Output = Result<Info>;
1998
1999 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
2000 self.poll_ok(waiter)
2001 }
2002}
2003
2004pub struct GroupRequest {
2013 state: kio::Producer<TrackState>,
2014 fetch: kio::Shared<FetchState>,
2016 sequence: u64,
2017 priority: u8,
2018 result: kio::Producer<FetchOutcome>,
2020 done: bool,
2021}
2022
2023impl GroupRequest {
2024 pub fn sequence(&self) -> u64 {
2026 self.sequence
2027 }
2028
2029 pub fn priority(&self) -> u8 {
2031 self.priority
2032 }
2033
2034 pub fn accept(mut self, info: impl Into<Option<Info>>) -> Result<group::Producer> {
2042 self.done = true;
2043 let res = TrackState::modify(&self.state)
2047 .and_then(|mut state| state.insert_group_request(self.sequence, info.into()));
2048 self.remove();
2049 res
2050 }
2051
2052 pub fn reject(mut self, err: Error) {
2054 self.done = true;
2055 self.remove();
2058 if let Ok(mut outcome) = self.result.write() {
2059 outcome.rejected = Some(err);
2060 }
2061 }
2062
2063 fn remove(&self) {
2066 self.fetch
2067 .lock()
2068 .remove_if(&self.sequence, |pending| pending.result.same_channel(&self.result));
2069 }
2070}
2071
2072impl Drop for GroupRequest {
2073 fn drop(&mut self) {
2074 if self.done {
2075 return;
2076 }
2077 self.remove();
2078 if let Ok(mut outcome) = self.result.write() {
2079 outcome.rejected = Some(Error::Dropped);
2080 }
2081 }
2082}
2083
2084pub struct Fetching {
2090 inner: FetchingKind,
2091 stats: stats::Scope,
2094}
2095
2096enum FetchingKind {
2097 Plain {
2098 state: kio::Consumer<TrackState>,
2099 fetch: kio::Shared<FetchState>,
2100 sequence: u64,
2101 result: Option<kio::Consumer<FetchOutcome>>,
2103 },
2104 Spliced(kio::Pending<super::resume::Fetching>),
2106}
2107
2108impl kio::Pollable for Fetching {
2109 type Output = Result<group::Consumer>;
2110
2111 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
2112 let (state, fetch, sequence, result) = match &self.inner {
2113 FetchingKind::Plain {
2114 state,
2115 fetch,
2116 sequence,
2117 result,
2118 } => (state, fetch, *sequence, result.as_ref()),
2119 FetchingKind::Spliced(spliced) => {
2120 return kio::Pollable::poll(&**spliced, waiter)
2123 .map(|res| res.map(|group| group.with_meter(self.stats.meter())));
2124 }
2125 };
2126
2127 match state.poll(waiter, |state| state.poll_fetch_cached(sequence)) {
2130 Poll::Ready(Ok(res)) => return Poll::Ready(res.map(|group| group.with_meter(self.stats.meter()))),
2131 Poll::Ready(Err(closed)) => {
2132 return Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped)));
2133 }
2134 Poll::Pending => {}
2135 }
2136
2137 let Some(result) = result else {
2139 return match fetch.poll(waiter, |fetch| match fetch.has_handlers() {
2142 false => Poll::Ready(()),
2143 true => Poll::Pending,
2144 }) {
2145 Poll::Ready(_guard) => Poll::Ready(Err(Error::NotFound)),
2146 Poll::Pending => Poll::Pending,
2147 };
2148 };
2149
2150 match result.poll(waiter, |outcome| match &outcome.rejected {
2153 Some(err) => Poll::Ready(err.clone()),
2154 None => Poll::Pending,
2155 }) {
2156 Poll::Ready(Ok(err)) => Poll::Ready(Err(err)),
2157 Poll::Ready(Err(_closed)) => Poll::Ready(Err(Error::NotFound)),
2158 Poll::Pending => Poll::Pending,
2159 }
2160 }
2161}
2162
2163pub struct Subscriber {
2186 name: Arc<str>,
2187 info: Info,
2188 inner: SubscriberKind,
2189 stats: stats::Scope,
2192 _stats_sub: stats::Subscription,
2195}
2196
2197enum SubscriberKind {
2198 Plain(PlainSubscriber),
2199 Spliced(Box<super::resume::Subscriber>),
2201}
2202
2203struct PlainSubscriber {
2205 state: kio::Consumer<TrackState>,
2206
2207 subscription: kio::Producer<Subscription>,
2208 index: usize,
2210 datagram_index: usize,
2212 min_sequence: u64,
2214 next_sequence: u64,
2217 end_sequence: Option<u64>,
2222 parked: BTreeMap<u64, group::Consumer>,
2227}
2228
2229impl PlainSubscriber {
2230 fn poll<F, R>(&self, waiter: &kio::Waiter, f: F) -> Poll<Result<R>>
2232 where
2233 F: Fn(&kio::Ref<'_, TrackState>) -> Poll<Result<R>>,
2234 {
2235 Poll::Ready(match ready!(self.state.poll(waiter, f)) {
2236 Ok(res) => res,
2237 Err(state) => Err(state.abort.clone().unwrap_or(Error::Dropped)),
2239 })
2240 }
2241
2242 fn poll_recv_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2243 let watch = |group: &group::Consumer| match group.poll_closed(waiter) {
2251 Poll::Pending => true,
2252 Poll::Ready(()) => !group.is_aborted(),
2253 };
2254
2255 let min_sequence = self.min_sequence;
2260 self.parked
2261 .retain(|sequence, group| *sequence >= min_sequence && watch(group));
2262
2263 if let Some(&sequence) = self.parked.keys().next()
2265 && self.end_sequence.is_none_or(|end| sequence <= end)
2266 {
2267 let group = self.parked.remove(&sequence).expect("parked key just observed");
2268 group.cache_refresh();
2270 return Poll::Ready(Ok(Some(group)));
2271 }
2272
2273 loop {
2274 let Some((consumer, found_index)) =
2275 ready!(self.poll(waiter, |state| state.poll_recv_group(self.index, self.min_sequence))?)
2276 else {
2277 if self.parked.is_empty() {
2280 return Poll::Ready(Ok(None));
2281 }
2282 return Poll::Pending;
2283 };
2284 self.index = found_index + 1;
2285
2286 if self.end_sequence.is_some_and(|end| consumer.sequence > end) {
2289 if watch(&consumer) {
2293 self.parked.insert(consumer.sequence, consumer);
2294 }
2295 continue;
2296 }
2297 return Poll::Ready(Ok(Some(consumer)));
2298 }
2299 }
2300
2301 fn poll_recv_datagram(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Datagram>>> {
2302 let Some((datagram, found_index)) =
2303 ready!(self.poll(waiter, |state| state.poll_recv_datagram(self.datagram_index))?)
2304 else {
2305 return Poll::Ready(Ok(None));
2306 };
2307
2308 self.datagram_index = found_index + 1;
2309 self.next_sequence = self.next_sequence.max(datagram.sequence.saturating_add(1));
2310 Poll::Ready(Ok(Some(datagram)))
2311 }
2312
2313 fn poll_next_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2314 let floor = self.next_sequence.max(self.min_sequence);
2315 let Some(group) = ready!(self.poll(waiter, |state| state.poll_next_in_range(floor, self.end_sequence))?) else {
2316 return Poll::Ready(Ok(None));
2317 };
2318 self.next_sequence = group.sequence.saturating_add(1);
2319 Poll::Ready(Ok(Some(group)))
2320 }
2321
2322 fn poll_read_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Frame>>> {
2323 let lower = self.min_sequence.max(self.next_sequence);
2324 let Some((frame, found_index, sequence)) =
2325 ready!(self.poll(waiter, |state| { state.poll_read_frame(self.index, lower, waiter) })?)
2326 else {
2327 return Poll::Ready(Ok(None));
2328 };
2329
2330 self.index = found_index + 1;
2331 self.next_sequence = sequence.saturating_add(1);
2332 Poll::Ready(Ok(Some(frame)))
2333 }
2334}
2335
2336#[derive(Clone)]
2342pub struct SubscriberControl {
2343 subscription: kio::Producer<Subscription>,
2344}
2345
2346impl SubscriberControl {
2347 pub fn subscription(&self) -> Subscription {
2349 self.subscription.read().clone()
2350 }
2351
2352 pub fn update(&self, subscription: Subscription) -> Result<()> {
2357 let mut state = self.subscription.write().map_err(|_| Error::Closed)?;
2358 *state = subscription;
2359 Ok(())
2360 }
2361}
2362
2363impl Subscriber {
2364 pub fn info(&self) -> &Info {
2369 &self.info
2370 }
2371
2372 pub fn name(&self) -> &str {
2374 &self.name
2375 }
2376
2377 pub fn control(&self) -> SubscriberControl {
2379 SubscriberControl {
2380 subscription: match &self.inner {
2381 SubscriberKind::Plain(plain) => plain.subscription.clone(),
2382 SubscriberKind::Spliced(spliced) => spliced.prefs(),
2383 },
2384 }
2385 }
2386
2387 pub fn poll_recv_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2404 let meter = self.stats.meter();
2405 let res = match &mut self.inner {
2406 SubscriberKind::Plain(plain) => plain.poll_recv_group(waiter),
2407 SubscriberKind::Spliced(spliced) => spliced.poll_recv_group(waiter),
2408 };
2409 res.map(|res| res.map(|group| group.map(|group| group.with_meter(meter))))
2410 }
2411
2412 pub async fn recv_group(&mut self) -> Result<Option<group::Consumer>> {
2419 kio::wait(|waiter| self.poll_recv_group(waiter)).await
2420 }
2421
2422 pub fn poll_recv_datagram(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Datagram>>> {
2433 let meter = self.stats.meter();
2434 let res = match &mut self.inner {
2435 SubscriberKind::Plain(plain) => plain.poll_recv_datagram(waiter),
2436 SubscriberKind::Spliced(spliced) => spliced.poll_recv_datagram(waiter),
2437 };
2438 if let Poll::Ready(Ok(Some(datagram))) = &res {
2441 meter.datagram(datagram.payload.len() as u64);
2442 }
2443 res
2444 }
2445
2446 pub async fn recv_datagram(&mut self) -> Result<Option<Datagram>> {
2453 kio::wait(|waiter| self.poll_recv_datagram(waiter)).await
2454 }
2455
2456 pub fn poll_next_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2465 let meter = self.stats.meter();
2466 let res = match &mut self.inner {
2467 SubscriberKind::Plain(plain) => plain.poll_next_group(waiter),
2468 SubscriberKind::Spliced(spliced) => spliced.poll_next_group(waiter),
2469 };
2470 res.map(|res| res.map(|group| group.map(|group| group.with_meter(meter))))
2471 }
2472
2473 pub async fn next_group(&mut self) -> Result<Option<group::Consumer>> {
2479 kio::wait(|waiter| self.poll_next_group(waiter)).await
2480 }
2481
2482 pub fn poll_read_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Frame>>> {
2486 let meter = self.stats.meter();
2487 let res = match &mut self.inner {
2488 SubscriberKind::Plain(plain) => plain.poll_read_frame(waiter),
2489 SubscriberKind::Spliced(spliced) => spliced.poll_read_frame(waiter),
2490 };
2491 if let Poll::Ready(Ok(Some(frame))) = &res {
2494 meter.group();
2495 meter.frames(1);
2496 meter.bytes(frame.payload.len() as u64);
2497 }
2498 res
2499 }
2500
2501 pub async fn read_frame(&mut self) -> Result<Option<frame::Frame>> {
2506 kio::wait(|waiter| self.poll_read_frame(waiter)).await
2507 }
2508
2509 pub fn is_clone(&self, other: &Self) -> bool {
2511 match (&self.inner, &other.inner) {
2512 (SubscriberKind::Plain(a), SubscriberKind::Plain(b)) => a.state.same_channel(&b.state),
2513 (SubscriberKind::Spliced(a), SubscriberKind::Spliced(b)) => a.is_clone(b),
2514 _ => false,
2515 }
2516 }
2517
2518 pub fn poll_finished(&mut self, waiter: &kio::Waiter) -> Poll<Result<u64>> {
2520 match &mut self.inner {
2521 SubscriberKind::Plain(plain) => plain.poll(waiter, |state| state.poll_finished()),
2522 SubscriberKind::Spliced(spliced) => spliced.poll_finished(waiter),
2523 }
2524 }
2525
2526 pub async fn finished(&mut self) -> Result<u64> {
2534 kio::wait(|waiter| self.poll_finished(waiter)).await
2535 }
2536
2537 pub fn start_at(&mut self, sequence: u64) {
2544 match &mut self.inner {
2545 SubscriberKind::Plain(plain) => plain.min_sequence = sequence,
2546 SubscriberKind::Spliced(spliced) => spliced.start_at(sequence),
2547 }
2548 }
2549
2550 pub fn end_at(&mut self, sequence: impl Into<Option<u64>>) {
2563 match &mut self.inner {
2564 SubscriberKind::Plain(plain) => plain.end_sequence = sequence.into(),
2565 SubscriberKind::Spliced(spliced) => spliced.end_at(sequence),
2566 }
2567 }
2568
2569 pub fn subscription(&self) -> Subscription {
2571 self.control().subscription()
2572 }
2573
2574 pub fn update(&mut self, subscription: Subscription) -> Result<()> {
2580 match &mut self.inner {
2581 SubscriberKind::Plain(plain) => {
2582 let mut state = plain.subscription.write().map_err(|_| Error::Closed)?;
2583 *state = subscription;
2584 }
2585 SubscriberKind::Spliced(spliced) => spliced.update(subscription),
2586 }
2587 Ok(())
2588 }
2589
2590 pub fn latest(&self) -> Option<u64> {
2592 match &self.inner {
2593 SubscriberKind::Plain(plain) => plain.state.read().max_sequence,
2594 SubscriberKind::Spliced(spliced) => spliced.latest(),
2595 }
2596 }
2597}
2598
2599pub struct Request {
2611 name: Arc<str>,
2612 broadcast: Arc<broadcast::Info>,
2614 state: kio::Producer<TrackState>,
2615
2616 prev_subscription: Option<Subscription>,
2618
2619 alive: Arc<Alive>,
2622
2623 _dynamic: Dynamic,
2628
2629 stats: stats::Scope,
2632}
2633
2634impl Request {
2635 pub(crate) fn new(broadcast: Arc<broadcast::Info>, name: impl Into<Arc<str>>) -> Self {
2636 let name = name.into();
2637 let state = TrackState::spawn(broadcast.clone());
2638 let alive = Alive::new(name.clone(), state.clone());
2639 let dynamic = Dynamic::new(name.clone(), state.clone(), alive.clone());
2640 Self {
2641 name,
2642 broadcast,
2643 state,
2644 prev_subscription: None,
2645 alive,
2646 _dynamic: dynamic,
2647 stats: stats::Scope::default(),
2648 }
2649 }
2650
2651 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
2654 self.stats = scope;
2655 self
2656 }
2657
2658 pub fn name(&self) -> &str {
2660 &self.name
2661 }
2662
2663 pub fn consume(&self) -> Consumer {
2665 Consumer::plain(self.name.clone(), self.state.consume())
2666 }
2667
2668 pub fn dynamic(&self) -> Dynamic {
2672 Dynamic::new(self.name.clone(), self.state.clone(), self.alive.clone())
2673 }
2674
2675 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
2678 self.state.poll_unused(waiter).map(|_| ())
2679 }
2680
2681 pub fn accept(self, info: impl Into<Option<Info>>) -> Producer {
2688 if let Ok(mut state) = self.state.write() {
2691 state.install(info.into().unwrap_or_default());
2692 }
2693 self.alive.publish(Some(&self.stats));
2696 Producer {
2697 name: self.name,
2698 broadcast: self.broadcast,
2699 state: self.state,
2700 prev_subscription: None,
2701 alive: self.alive,
2702 stats: self.stats,
2703 }
2704 }
2705
2706 pub fn reject(self, err: Error) {
2708 if let Ok(mut state) = self.state.write() {
2709 state.abort = Some(err);
2710 }
2711 }
2712
2713 pub fn subscription(&self) -> Option<Subscription> {
2716 let state = self.state.read();
2717 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
2718 drop(state);
2719 snapshot_subscription(&subs, bound)
2720 }
2721
2722 pub async fn subscription_changed(&mut self) -> Option<Subscription> {
2725 kio::wait(|waiter| self.poll_subscription_changed(waiter)).await
2726 }
2727
2728 pub fn poll_subscription_changed(&mut self, waiter: &kio::Waiter) -> Poll<Option<Subscription>> {
2730 let state = self.state.read();
2731 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
2732 drop(state);
2733
2734 let prev = &self.prev_subscription;
2735 let mut combined = None;
2736 let mut guard = ready!(subs.poll(waiter, |subs| {
2737 let next = combined_subscription(subs, bound, waiter);
2738 if &next == prev {
2739 Poll::Pending
2740 } else {
2741 combined = next;
2742 Poll::Ready(())
2743 }
2744 }));
2745 guard.retain(|sub| !sub.is_closed());
2747 drop(guard);
2748 self.prev_subscription = combined.clone();
2749 Poll::Ready(combined)
2750 }
2751
2752 pub(super) fn weak(&self) -> TrackWeak {
2753 TrackWeak {
2754 name: self.name.clone(),
2755 state: self.state.weak(),
2756 }
2757 }
2758}
2759
2760#[cfg(test)]
2761use futures::FutureExt;
2762
2763#[cfg(test)]
2764#[allow(missing_docs)] impl Subscriber {
2766 pub fn assert_group(&mut self) -> group::Consumer {
2767 self.recv_group()
2768 .now_or_never()
2769 .expect("group would have blocked")
2770 .expect("would have errored")
2771 .expect("track was closed")
2772 }
2773
2774 pub fn assert_no_group(&mut self) {
2775 assert!(
2776 self.recv_group().now_or_never().is_none(),
2777 "recv_group would not have blocked"
2778 );
2779 }
2780
2781 pub fn assert_not_closed(&mut self) {
2782 assert!(self.finished().now_or_never().is_none(), "should not be closed");
2783 }
2784
2785 pub fn assert_closed(&mut self) {
2786 assert!(self.finished().now_or_never().is_some(), "should be closed");
2787 }
2788
2789 pub fn assert_error(&mut self) {
2791 assert!(
2792 self.finished().now_or_never().expect("should not block").is_err(),
2793 "should be error"
2794 );
2795 }
2796
2797 pub fn assert_is_clone(&self, other: &Self) {
2798 assert!(self.is_clone(other), "should be clone");
2799 }
2800
2801 pub fn assert_not_clone(&self, other: &Self) {
2802 assert!(!self.is_clone(other), "should not be clone");
2803 }
2804}
2805
2806#[cfg(test)]
2807mod test {
2808 use super::*;
2809 use crate::model::test_tracing::count_drop_warnings;
2810
2811 fn track_producer(name: impl Into<Arc<str>>, info: impl Into<Option<Info>>) -> Producer {
2814 Producer::new(Arc::new(broadcast::Info::default()), name, info)
2815 }
2816
2817 fn live_groups(state: &TrackState) -> usize {
2819 state.lookup.len()
2820 }
2821
2822 fn first_live_sequence(state: &TrackState) -> u64 {
2824 state
2825 .arrival
2826 .iter()
2827 .find(|(sequence, stamp)| state.lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp))
2828 .map(|(sequence, _)| *sequence)
2829 .unwrap()
2830 }
2831
2832 fn recv_datagram(dg: &mut Subscriber) -> Datagram {
2834 dg.recv_datagram()
2835 .now_or_never()
2836 .expect("datagram would have blocked")
2837 .expect("would have errored")
2838 .expect("track was closed")
2839 }
2840
2841 #[tokio::test]
2842 async fn append_datagram_shares_group_sequence() {
2843 let mut producer = track_producer("test", None);
2844 let ts = Timestamp::from_millis(10).unwrap();
2845
2846 assert_eq!(producer.append_group().unwrap().sequence, 0);
2848 assert_eq!(producer.append_datagram(ts, &b"a"[..]).unwrap(), 1);
2849 assert_eq!(producer.append_group().unwrap().sequence, 2);
2850 assert_eq!(producer.append_datagram(ts, &b"b"[..]).unwrap(), 3);
2851 assert_eq!(producer.latest(), Some(3));
2852 }
2853
2854 #[tokio::test]
2855 async fn append_datagram_roundtrip() {
2856 let mut producer = track_producer("test", None);
2857 let mut dg = producer.subscribe(None);
2858
2859 let ts = Timestamp::from_millis(42).unwrap();
2860 let seq = producer.append_datagram(ts, &b"hello"[..]).unwrap();
2861
2862 let got = recv_datagram(&mut dg);
2863 assert_eq!(got.sequence, seq);
2864 assert_eq!(got.timestamp, ts);
2865 assert_eq!(&got.payload[..], b"hello");
2866 }
2867
2868 #[tokio::test]
2869 async fn write_datagram_preserves_sequence() {
2870 let mut producer = track_producer("test", None);
2871 let mut dg = producer.subscribe(None);
2872
2873 let ts = Timestamp::from_millis(5).unwrap();
2874 producer
2876 .write_datagram(Datagram {
2877 sequence: 100,
2878 timestamp: ts,
2879 payload: bytes::Bytes::from_static(b"x"),
2880 })
2881 .unwrap();
2882
2883 assert_eq!(recv_datagram(&mut dg).sequence, 100);
2884 assert_eq!(producer.append_group().unwrap().sequence, 101);
2886 }
2887
2888 #[tokio::test]
2889 async fn recv_datagram_advances_ordered_group_cursor() {
2890 let mut producer = track_producer("test", None);
2891 let mut subscriber = producer.subscribe(None);
2892 let ts = Timestamp::from_millis(5).unwrap();
2893
2894 producer
2895 .write_datagram(Datagram {
2896 sequence: 5,
2897 timestamp: ts,
2898 payload: bytes::Bytes::from_static(b"x"),
2899 })
2900 .unwrap();
2901 assert_eq!(recv_datagram(&mut subscriber).sequence, 5);
2902
2903 producer.create_group(group::Info { sequence: 3 }).unwrap();
2904 producer.create_group(group::Info { sequence: 6 }).unwrap();
2905
2906 let group = subscriber
2907 .next_group()
2908 .now_or_never()
2909 .expect("group would have blocked")
2910 .expect("would have errored")
2911 .expect("track was closed");
2912 assert_eq!(group.sequence, 6);
2913 }
2914
2915 #[tokio::test]
2916 async fn datagram_normalized_to_track_timescale() {
2917 let info = Info::default().with_timescale(Timescale::MICRO);
2918 let mut producer = track_producer("test", info);
2919 let mut dg = producer.subscribe(None);
2920
2921 producer
2923 .append_datagram(Timestamp::from_millis(2).unwrap(), &b"z"[..])
2924 .unwrap();
2925 let got = recv_datagram(&mut dg);
2926 assert_eq!(got.timestamp.scale(), Timescale::MICRO);
2927 assert_eq!(got.timestamp.value(), 2_000);
2928 }
2929
2930 #[tokio::test]
2931 async fn datagram_rejects_oversized() {
2932 let mut producer = track_producer("test", None);
2933 let big = bytes::Bytes::from(vec![0u8; crate::model::datagram::MAX_DATAGRAM_PAYLOAD + 1]);
2934 let ts = Timestamp::from_millis(0).unwrap();
2935 assert!(matches!(
2936 producer.append_datagram(ts, big.clone()),
2937 Err(Error::FrameTooLarge)
2938 ));
2939 assert!(matches!(
2940 producer.write_datagram(Datagram {
2941 sequence: 0,
2942 timestamp: ts,
2943 payload: big,
2944 }),
2945 Err(Error::FrameTooLarge)
2946 ));
2947 }
2948
2949 #[tokio::test]
2950 async fn datagram_fanout_to_subscribers() {
2951 let mut producer = track_producer("test", None);
2952 let mut a = producer.subscribe(None);
2954 let mut b = producer.subscribe(None);
2955 let ts = Timestamp::from_millis(1).unwrap();
2956
2957 producer.append_datagram(ts, &b"first"[..]).unwrap();
2958 producer.append_datagram(ts, &b"second"[..]).unwrap();
2959
2960 assert_eq!(&recv_datagram(&mut a).payload[..], b"first");
2962 assert_eq!(&recv_datagram(&mut a).payload[..], b"second");
2963 assert_eq!(&recv_datagram(&mut b).payload[..], b"first");
2964 assert_eq!(&recv_datagram(&mut b).payload[..], b"second");
2965 }
2966
2967 #[tokio::test]
2968 async fn datagram_evicts_stale() {
2969 tokio::time::pause();
2970
2971 let mut producer = track_producer("test", None);
2972 let mut dg = producer.subscribe(None);
2973 let ts = Timestamp::from_millis(0).unwrap();
2974
2975 producer.append_datagram(ts, &b"old"[..]).unwrap(); tokio::time::advance(MAX_DATAGRAM_AGE + Duration::from_millis(10)).await;
2979 producer.append_datagram(ts, &b"new"[..]).unwrap(); let got = recv_datagram(&mut dg);
2983 assert_eq!(got.sequence, 1);
2984 assert_eq!(&got.payload[..], b"new");
2985 }
2986
2987 #[tokio::test]
2988 async fn datagram_recv_pends_until_written() {
2989 let mut producer = track_producer("test", None);
2990 let mut dg = producer.subscribe(None);
2991
2992 assert!(
2993 dg.recv_datagram().now_or_never().is_none(),
2994 "should block with no datagrams"
2995 );
2996
2997 producer
2998 .append_datagram(Timestamp::from_millis(0).unwrap(), &b"go"[..])
2999 .unwrap();
3000 assert_eq!(&recv_datagram(&mut dg).payload[..], b"go");
3001 }
3002
3003 #[tokio::test]
3007 async fn datagram_wire_roundtrip_between_tracks() {
3008 use crate::coding::{Decode, Encode};
3009 use crate::lite;
3010
3011 let version = lite::Version::Lite05;
3012
3013 let mut origin = track_producer("test", None);
3015 let mut origin_dg = origin.subscribe(None);
3016 let ts = Timestamp::from_millis(7).unwrap();
3017 let seq = origin.append_datagram(ts, &b"payload"[..]).unwrap();
3018
3019 let d = recv_datagram(&mut origin_dg);
3020 let body = lite::Datagram {
3021 subscribe: 5,
3022 sequence: d.sequence,
3023 timestamp: d.timestamp.value(),
3024 payload: d.payload.clone(),
3025 }
3026 .encode_bytes(version)
3027 .unwrap();
3028
3029 let mut slice = &body[..];
3031 let wire = lite::Datagram::decode(&mut slice, version).unwrap();
3032 let mut downstream = track_producer("test", None);
3033 let mut downstream_dg = downstream.subscribe(None);
3034 downstream
3035 .write_datagram(Datagram {
3036 sequence: wire.sequence,
3037 timestamp: Timestamp::new(wire.timestamp, Timescale::MILLI).unwrap(),
3038 payload: wire.payload,
3039 })
3040 .unwrap();
3041
3042 let got = recv_datagram(&mut downstream_dg);
3043 assert_eq!(got.sequence, seq);
3044 assert_eq!(got.timestamp, ts);
3045 assert_eq!(&got.payload[..], b"payload");
3046 }
3047
3048 #[tokio::test]
3049 async fn evict_expired_groups() {
3050 tokio::time::pause();
3051
3052 let mut producer = track_producer("test", None);
3053
3054 producer.append_group().unwrap(); producer.append_group().unwrap(); producer.append_group().unwrap(); {
3060 let state = producer.state.read();
3061 assert_eq!(live_groups(&state), 3);
3062 assert_eq!(state.offset, 0);
3063 }
3064
3065 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3067
3068 producer.append_group().unwrap(); {
3075 let state = producer.state.read();
3076 assert_eq!(live_groups(&state), 1);
3077 assert_eq!(first_live_sequence(&state), 3);
3078 assert_eq!(state.offset, 3);
3079 assert!(!state.lookup.contains_key(&0));
3080 assert!(!state.lookup.contains_key(&1));
3081 assert!(!state.lookup.contains_key(&2));
3082 assert!(state.lookup.contains_key(&3));
3083 }
3084 }
3085
3086 #[tokio::test]
3090 async fn aging_out_a_finished_group_keeps_the_clean_end() {
3091 tokio::time::pause();
3092
3093 let mut producer = track_producer("test", None);
3094 let mut group = producer.create_group(group::Info { sequence: 0 }).unwrap();
3095 let mut consumer = group.consume();
3096
3097 group
3098 .write_frame(Timestamp::from_millis(0).unwrap(), b"hello".as_slice())
3099 .unwrap();
3100 assert_eq!(consumer.next_frame().await.unwrap().unwrap().size, 5);
3101
3102 tokio::time::advance(DEFAULT_LATENCY_MAX * 12).await;
3104 group.finish().unwrap();
3105 let _next = producer.create_group(group::Info { sequence: 1 }).unwrap();
3106
3107 assert!(consumer.next_frame().await.unwrap().is_none());
3108 }
3109
3110 #[tokio::test]
3114 async fn active_reader_survives_expiry() {
3115 tokio::time::pause();
3116
3117 let mut producer = track_producer("test", None);
3118 let mut subscriber = producer.subscribe(None);
3119
3120 let mut group = producer.create_group(0u64.into()).unwrap();
3122 for _ in 0..10 {
3123 group.write_frame(Timestamp::ZERO, b"x".as_slice()).unwrap();
3124 }
3125 group.finish().unwrap();
3126 let mut reading = subscriber.assert_group();
3127
3128 producer.create_group(1u64.into()).unwrap().finish().unwrap();
3130
3131 for seq in 2..12u64 {
3134 tokio::time::advance(DEFAULT_LATENCY_MAX / 2).await;
3135 let frame = reading.next_frame().await;
3136 assert!(
3137 matches!(frame, Ok(Some(_))),
3138 "an actively-read group must not expire mid-read (step {seq})"
3139 );
3140 producer.create_group(seq.into()).unwrap().finish().unwrap();
3141 }
3142
3143 let state = producer.state.read();
3144 assert!(state.lookup.contains_key(&0), "the read group survived");
3145 assert!(!state.lookup.contains_key(&1), "the unread group still expired");
3146 }
3147
3148 #[tokio::test]
3152 async fn slow_frame_reader_survives_expiry() {
3153 tokio::time::pause();
3154
3155 let mut producer = track_producer("test", None);
3156 let mut subscriber = producer.subscribe(None);
3157
3158 let mut group = producer.create_group(0u64.into()).unwrap();
3159 for _ in 0..20 {
3160 group.write_frame(Timestamp::ZERO, b"x".as_slice()).unwrap();
3161 }
3162 group.finish().unwrap();
3163 let mut reading = subscriber.assert_group();
3164
3165 for seq in 1..20u64 {
3167 tokio::time::advance(DEFAULT_LATENCY_MAX / 2).await;
3168 let frame = reading.read_frame().await;
3169 assert!(
3170 matches!(frame, Ok(Some(_))),
3171 "a slow reader must not expire mid-read (step {seq})"
3172 );
3173 producer.create_group(seq.into()).unwrap().finish().unwrap();
3174 }
3175 }
3176
3177 #[tokio::test]
3182 async fn slow_batch_reader_survives_expiry_with_keep_alive() {
3183 tokio::time::pause();
3184
3185 let mut producer = track_producer("test", None);
3186 let mut subscriber = producer.subscribe(None);
3187
3188 let mut group = producer.create_group(0u64.into()).unwrap();
3189 for _ in 0..20 {
3190 group.write_frame(Timestamp::ZERO, b"x".as_slice()).unwrap();
3191 }
3192 group.finish().unwrap();
3193 let mut reading = subscriber.assert_group();
3194
3195 let mut buf = crate::frame::Buffer::<8>::new();
3198 let count = reading.read_frames(&mut buf).await.unwrap().len();
3199 assert_eq!(count, 8, "the batch is bounded by the buffer");
3200
3201 for step in 0..8u64 {
3202 tokio::time::advance(DEFAULT_LATENCY_MAX / 2).await;
3203 reading.keep_alive();
3204 producer.create_group((step + 1).into()).unwrap().finish().unwrap();
3206 }
3207
3208 let rest = reading
3210 .read_frames(&mut buf)
3211 .await
3212 .expect("a batch reader that kept the group alive must not be expired");
3213 assert_eq!(rest.len(), 8, "the next batch picks up where the last one stopped");
3214 }
3215
3216 #[tokio::test]
3220 async fn delivery_restarts_the_expiry_clock() {
3221 tokio::time::pause();
3222
3223 let mut producer = track_producer("test", None);
3224 let mut subscriber = producer.subscribe(None);
3225
3226 let mut group = producer.create_group(0u64.into()).unwrap();
3227 group.write_frame(Timestamp::ZERO, b"x".as_slice()).unwrap();
3228 group.finish().unwrap();
3229 producer.create_group(1u64.into()).unwrap().finish().unwrap();
3231
3232 tokio::time::advance(DEFAULT_LATENCY_MAX - Duration::from_secs(1)).await;
3234 let mut reading = subscriber.assert_group();
3235
3236 tokio::time::advance(DEFAULT_LATENCY_MAX - Duration::from_secs(1)).await;
3239 producer.create_group(2u64.into()).unwrap().finish().unwrap();
3240
3241 let frame = reading.read_frame().await.unwrap();
3242 assert!(frame.is_some(), "a just-delivered group must not expire unread");
3243 }
3244
3245 #[tokio::test]
3249 async fn streaming_frame_writes_keep_the_group_alive() {
3250 tokio::time::pause();
3251
3252 let mut producer = track_producer("test", None);
3253 let mut straggler = producer.create_group(0u64.into()).unwrap();
3254 producer.create_group(1u64.into()).unwrap().finish().unwrap();
3256
3257 let mut frame = straggler
3258 .create_frame(frame::Info {
3259 size: 10,
3260 timestamp: Timestamp::ZERO,
3261 })
3262 .unwrap();
3263 for seq in 2..12u64 {
3266 tokio::time::advance(DEFAULT_LATENCY_MAX / 2).await;
3267 frame.write(bytes::Bytes::from_static(b"x")).unwrap();
3268 producer.create_group(seq.into()).unwrap().finish().unwrap();
3269 }
3270 frame.finish().unwrap();
3271 straggler.finish().unwrap();
3272
3273 let state = producer.state.read();
3274 assert!(
3275 state.lookup.contains_key(&0),
3276 "a group streaming a frame survives expiry"
3277 );
3278 }
3279
3280 #[tokio::test]
3283 async fn parked_reoffer_restarts_the_expiry_clock() {
3284 tokio::time::pause();
3285
3286 let mut producer = track_producer("test", None);
3287 let mut subscriber = producer.subscribe(None);
3288 subscriber.end_at(0);
3289
3290 for seq in 0..2u64 {
3291 let mut group = producer.create_group(seq.into()).unwrap();
3292 group.write_frame(Timestamp::ZERO, b"x".as_slice()).unwrap();
3293 group.finish().unwrap();
3294 }
3295
3296 assert_eq!(subscriber.assert_group().sequence, 0);
3298 subscriber.assert_no_group();
3299
3300 tokio::time::advance(DEFAULT_LATENCY_MAX - Duration::from_secs(1)).await;
3302 subscriber.end_at(1);
3303 let mut reading = subscriber.assert_group();
3304 assert_eq!(reading.sequence, 1);
3305
3306 tokio::time::advance(DEFAULT_LATENCY_MAX - Duration::from_secs(1)).await;
3309 producer.create_group(2u64.into()).unwrap().finish().unwrap();
3310
3311 let frame = reading.read_frame().await.unwrap();
3312 assert!(frame.is_some(), "a just-re-offered group must not expire unread");
3313 }
3314
3315 #[tokio::test]
3316 async fn evict_keeps_max_sequence() {
3317 tokio::time::pause();
3318
3319 let mut producer = track_producer("test", None);
3320 producer.append_group().unwrap(); tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3324
3325 producer.append_group().unwrap(); {
3329 let state = producer.state.read();
3330 assert_eq!(live_groups(&state), 1);
3331 assert_eq!(first_live_sequence(&state), 1);
3332 assert_eq!(state.offset, 1);
3333 }
3334 }
3335
3336 #[tokio::test]
3337 async fn no_eviction_when_fresh() {
3338 tokio::time::pause();
3339
3340 let mut producer = track_producer("test", None);
3341 producer.append_group().unwrap(); producer.append_group().unwrap(); producer.append_group().unwrap(); {
3346 let state = producer.state.read();
3347 assert_eq!(live_groups(&state), 3);
3348 assert_eq!(state.offset, 0);
3349 }
3350 }
3351
3352 #[tokio::test]
3353 async fn consumer_skips_evicted_groups() {
3354 tokio::time::pause();
3355
3356 let mut producer = track_producer("test", None);
3357 producer.append_group().unwrap(); let mut consumer = producer.subscribe(None);
3360
3361 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3362 producer.append_group().unwrap(); let group = consumer.assert_group();
3366 assert_eq!(group.sequence, 1);
3367 }
3368
3369 #[tokio::test]
3370 async fn cache_age_controls_eviction() {
3371 tokio::time::pause();
3372
3373 let mut producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(1)));
3375 producer.append_group().unwrap(); tokio::time::advance(Duration::from_secs(2)).await;
3379 producer.append_group().unwrap(); let state = producer.state.read();
3383 assert_eq!(live_groups(&state), 1);
3384 assert_eq!(first_live_sequence(&state), 1);
3385 }
3386
3387 #[test]
3388 fn latency_max_clamped_to_cache() {
3389 let producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(2)));
3390
3391 let mut subscriber = producer.subscribe(Subscription::default().with_latency_max(Duration::from_secs(10)));
3395 assert_eq!(subscriber.subscription().latency_max, Duration::from_secs(10));
3396 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3397
3398 subscriber
3400 .update(Subscription::default().with_latency_max(Duration::from_millis(500)))
3401 .unwrap();
3402 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_millis(500));
3403
3404 subscriber
3405 .update(Subscription::default().with_latency_max(Duration::ZERO))
3406 .unwrap();
3407 assert_eq!(producer.subscription().unwrap().latency_max, Duration::ZERO);
3408 }
3409
3410 fn track_producer_capped(name: impl Into<Arc<str>>, info: Info, cap: Duration) -> Producer {
3413 let origin = crate::origin::Info::default().with_cache_duration(cap);
3414 Producer::new(Arc::new(broadcast::Info { origin }), name, info)
3415 }
3416
3417 #[test]
3418 fn origin_cache_duration_clamps_latency_max() {
3419 let capped = track_producer_capped(
3422 "test",
3423 Info::default().with_latency_max(Duration::from_secs(60)),
3424 Duration::from_secs(1),
3425 );
3426 assert_eq!(capped.state.read().latency_bound(), Some(Duration::from_secs(1)));
3427
3428 let under = track_producer_capped(
3429 "test",
3430 Info::default().with_latency_max(Duration::from_millis(500)),
3431 Duration::from_secs(1),
3432 );
3433 assert_eq!(under.state.read().latency_bound(), Some(Duration::from_millis(500)));
3434 }
3435
3436 #[tokio::test]
3437 async fn origin_cache_duration_caps_eviction() {
3438 tokio::time::pause();
3439
3440 let mut producer = track_producer_capped(
3442 "test",
3443 Info::default().with_latency_max(Duration::from_secs(60)),
3444 Duration::from_secs(1),
3445 );
3446 producer.append_group().unwrap(); tokio::time::advance(Duration::from_secs(2)).await;
3450 producer.append_group().unwrap(); let state = producer.state.read();
3454 assert_eq!(live_groups(&state), 1);
3455 assert_eq!(first_live_sequence(&state), 1);
3456 }
3457
3458 #[test]
3459 fn latency_max_clamped_via_every_update_path() {
3460 let producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(2)));
3461 let over = Subscription::default().with_latency_max(Duration::from_secs(10));
3462
3463 let mut subscriber = producer.subscribe(over.clone());
3466 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3467
3468 subscriber.control().update(over.clone()).unwrap();
3469 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3470
3471 subscriber.update(over).unwrap();
3472 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3473 }
3474
3475 #[test]
3476 fn latency_max_aggregate_clamps_the_max_across_subscribers() {
3477 let producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(2)));
3478
3479 let _a = producer.subscribe(Subscription::default().with_latency_max(Duration::from_millis(500)));
3482 let _b = producer.subscribe(Subscription::default().with_latency_max(Duration::from_secs(10)));
3483
3484 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3485 }
3486
3487 #[test]
3488 fn subscriber_control_updates_while_read_future_is_pending() {
3489 let producer = track_producer("test", None);
3490 let mut subscriber = producer.subscribe(None);
3491 let control = subscriber.control();
3492
3493 let mut recv = Box::pin(subscriber.recv_group());
3494 assert!(recv.as_mut().now_or_never().is_none());
3495
3496 control
3497 .update(Subscription::default().with_priority(7).with_ordered(false))
3498 .unwrap();
3499
3500 let aggregate = producer.subscription().expect("expected an active subscription");
3501 assert_eq!(aggregate.priority, 7);
3502 assert!(!aggregate.ordered);
3503 }
3504
3505 #[test]
3506 fn dropped_subscriber_leaves_no_ghost_in_aggregate() {
3507 let mut producer = track_producer("test", None);
3512 let a = producer.subscribe(Subscription::default().with_priority(5));
3513
3514 let waiter = kio::Waiter::noop();
3516 assert!(
3517 matches!(producer.poll_subscription_changed(&waiter), Poll::Ready(Ok(Some(_)))),
3518 "one live subscriber should aggregate to Some",
3519 );
3520
3521 drop(a);
3523
3524 assert!(
3526 matches!(producer.poll_subscription_changed(&waiter), Poll::Ready(Ok(None))),
3527 "a dropped subscriber must not linger in the aggregate",
3528 );
3529
3530 assert!(
3532 producer.subscription().is_none(),
3533 "snapshot must exclude a dropped subscriber",
3534 );
3535 }
3536
3537 #[test]
3538 fn dropped_subscriber_wakes_the_aggregate() {
3539 use std::sync::atomic::{AtomicBool, Ordering};
3546
3547 let mut producer = track_producer("test", None);
3548 let a = producer.subscribe(Subscription::default().with_priority(5));
3549
3550 let woken = Arc::new(AtomicBool::new(false));
3551 let waiter = kio::Waiter::new(futures::task::waker(Arc::new(FlagWake(woken.clone()))));
3552
3553 assert!(matches!(
3555 producer.poll_subscription_changed(&waiter),
3556 Poll::Ready(Ok(Some(_)))
3557 ));
3558 assert!(
3559 producer.poll_subscription_changed(&waiter).is_pending(),
3560 "the aggregate is unchanged, so this poll must park",
3561 );
3562 assert!(!woken.load(Ordering::SeqCst), "nothing happened yet");
3563
3564 drop(a);
3565 assert!(
3566 woken.load(Ordering::SeqCst),
3567 "the last subscriber leaving must wake the aggregate watcher",
3568 );
3569 }
3570
3571 struct FlagWake(Arc<std::sync::atomic::AtomicBool>);
3573
3574 impl futures::task::ArcWake for FlagWake {
3575 fn wake_by_ref(arc_self: &Arc<Self>) {
3576 arc_self.0.store(true, std::sync::atomic::Ordering::SeqCst);
3577 }
3578 }
3579
3580 #[tokio::test]
3581 async fn out_of_order_max_sequence_at_front() {
3582 tokio::time::pause();
3583
3584 let mut producer = track_producer("test", None);
3585
3586 producer.create_group(group::Info { sequence: 5 }).unwrap();
3588 producer.create_group(group::Info { sequence: 3 }).unwrap();
3589 producer.create_group(group::Info { sequence: 4 }).unwrap();
3590
3591 {
3593 let state = producer.state.read();
3594 assert_eq!(state.max_sequence, Some(5));
3595 }
3596
3597 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3599
3600 producer.append_group().unwrap(); {
3606 let state = producer.state.read();
3607 assert_eq!(live_groups(&state), 1);
3608 assert_eq!(first_live_sequence(&state), 6);
3609 assert!(!state.lookup.contains_key(&3));
3610 assert!(!state.lookup.contains_key(&4));
3611 assert!(!state.lookup.contains_key(&5));
3612 assert!(state.lookup.contains_key(&6));
3613 }
3614 }
3615
3616 #[tokio::test]
3617 async fn max_sequence_at_front_blocks_trim() {
3618 tokio::time::pause();
3619
3620 let mut producer = track_producer("test", None);
3621
3622 producer.create_group(group::Info { sequence: 5 }).unwrap();
3624
3625 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3626
3627 producer.create_group(group::Info { sequence: 3 }).unwrap();
3629
3630 {
3633 let state = producer.state.read();
3634 assert_eq!(live_groups(&state), 2);
3635 assert_eq!(state.offset, 0);
3636 }
3637
3638 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3640
3641 producer.create_group(group::Info { sequence: 2 }).unwrap();
3643
3644 {
3649 let state = producer.state.read();
3650 assert_eq!(live_groups(&state), 2);
3651 assert_eq!(state.offset, 0);
3652 assert!(state.lookup.contains_key(&5));
3653 assert!(!state.lookup.contains_key(&3));
3654 assert!(state.lookup.contains_key(&2));
3655 }
3656
3657 let mut consumer = producer.subscribe(None);
3659 let group = consumer.assert_group();
3660 assert_eq!(group.sequence, 5);
3662 }
3663
3664 #[tokio::test]
3665 async fn abort_clears_cached_groups() {
3666 let mut producer = track_producer("test", None);
3667 producer.append_group().unwrap();
3668 producer.append_group().unwrap();
3669
3670 let mut consumer = producer.subscribe(None);
3672 assert_eq!(live_groups(&producer.state.read()), 2);
3673
3674 producer.clone().abort(Error::Cancel).unwrap();
3675
3676 {
3677 let state = producer.state.read();
3678 assert!(state.lookup.is_empty(), "cached groups should be dropped on abort");
3679 assert!(state.arrival.is_empty());
3680 assert!(state.evict.is_empty());
3681 }
3682
3683 let result = consumer.recv_group().now_or_never().expect("should not block");
3685 assert!(matches!(result, Err(Error::Cancel)));
3686 }
3687
3688 #[tokio::test]
3689 async fn drop_unfinished_clears_cached_groups() {
3690 let producer = track_producer("test", None);
3691 let mut writer = producer.clone();
3692 writer.append_group().unwrap();
3693
3694 let mut consumer = producer.subscribe(None);
3696 assert_eq!(live_groups(&producer.state.read()), 1);
3697
3698 drop(writer);
3700 drop(producer);
3701
3702 let result = consumer.recv_group().now_or_never().expect("should not block");
3703 assert!(matches!(result, Err(Error::Dropped)));
3704 }
3705
3706 #[tokio::test]
3707 async fn drop_after_abort_does_not_warn() {
3708 let warns = count_drop_warnings("track::Producer dropped without finish", || {
3711 let producer = track_producer("test", None);
3712 let keep = producer.clone();
3713 let mut writer = producer.clone();
3714 let mut group = writer.append_group().unwrap();
3715 group.finish().unwrap();
3716 let _consumer = producer.subscribe(None);
3717 writer.abort(Error::Cancel).unwrap();
3718 drop(keep);
3719 });
3720 assert_eq!(warns, 0, "abort-then-drop must not emit unfinished-producer WARN");
3721 }
3722
3723 #[tokio::test]
3724 async fn drop_unfinished_warns() {
3725 let warns = count_drop_warnings("track::Producer dropped without finish", || {
3726 let producer = track_producer("test", None);
3727 let mut writer = producer.clone();
3728 writer.append_group().unwrap();
3729 let _consumer = producer.subscribe(None);
3730 drop(writer);
3731 drop(producer);
3732 });
3733 assert!(warns >= 1, "unfinished drop must emit unfinished-producer WARN");
3734 }
3735
3736 #[tokio::test]
3737 async fn drop_finished_keeps_cached_groups() {
3738 let mut producer = track_producer("test", None);
3739 producer.append_group().unwrap();
3740 producer.finish().unwrap();
3741
3742 let mut consumer = producer.subscribe(None);
3743 drop(producer);
3744
3745 assert_eq!(consumer.assert_group().sequence, 0);
3747 let done = consumer.recv_group().now_or_never().expect("should not block").unwrap();
3748 assert!(done.is_none(), "consumer should drain then see clean finish");
3749 }
3750
3751 #[test]
3752 fn append_finish_cannot_be_rewritten() {
3753 let mut producer = track_producer("test", None);
3754
3755 assert!(producer.finish().is_ok());
3757 assert!(producer.finish().is_err());
3758 assert!(producer.append_group().is_err());
3759 }
3760
3761 #[test]
3762 fn finish_after_groups() {
3763 let mut producer = track_producer("test", None);
3764
3765 producer.append_group().unwrap();
3766 assert!(producer.finish().is_ok());
3767 assert!(producer.finish().is_err());
3768 assert!(producer.append_group().is_err());
3769 }
3770
3771 #[test]
3772 fn finish_at_rejects_a_boundary_at_or_below_the_live_edge() {
3773 let mut producer = track_producer("test", None);
3774 producer.create_group(group::Info { sequence: 5 }).unwrap();
3775
3776 assert!(producer.finish_at(4).is_err());
3779 assert!(producer.finish_at(5).is_err());
3780 assert!(producer.finish_at(6).is_ok());
3781
3782 {
3783 let state = producer.state.read();
3784 assert_eq!(state.final_sequence, Some(6));
3785 }
3786
3787 assert!(producer.finish_at(6).is_err());
3789 assert!(producer.create_group(group::Info { sequence: 4 }).is_ok());
3790 assert!(producer.create_group(group::Info { sequence: 6 }).is_err());
3791 }
3792
3793 #[test]
3794 fn final_sequence_reports_the_declared_boundary() {
3795 let mut producer = track_producer("test", None);
3796 assert_eq!(producer.final_sequence(), None);
3797
3798 producer.create_group(group::Info { sequence: 5 }).unwrap();
3799 assert_eq!(producer.final_sequence(), None, "a group does not declare a boundary");
3800
3801 producer.finish_at(9).unwrap();
3802 assert_eq!(producer.final_sequence(), Some(9));
3803
3804 assert!(producer.finish().is_err());
3806 }
3807
3808 #[test]
3809 fn final_sequence_reports_the_live_edge_after_finish() {
3810 let mut producer = track_producer("test", None);
3811 producer.create_group(group::Info { sequence: 5 }).unwrap();
3812 producer.finish().unwrap();
3813 assert_eq!(producer.final_sequence(), Some(6));
3814 }
3815
3816 #[tokio::test]
3817 async fn finish_at_declares_a_future_boundary() {
3818 let mut producer = track_producer("test", None);
3819 producer.create_group(group::Info { sequence: 5 }).unwrap();
3820
3821 producer.finish_at(7).unwrap();
3823
3824 let mut consumer = producer.subscribe(None);
3825 assert_eq!(consumer.assert_group().sequence, 5);
3826
3827 let boundary = consumer
3830 .finished()
3831 .now_or_never()
3832 .expect("boundary is known immediately")
3833 .expect("would have errored");
3834 assert_eq!(boundary, 7);
3835 assert!(
3836 consumer.recv_group().now_or_never().is_none(),
3837 "should wait for the outstanding group"
3838 );
3839
3840 producer.create_group(group::Info { sequence: 6 }).unwrap();
3842 assert_eq!(consumer.assert_group().sequence, 6);
3843 let done = consumer
3844 .recv_group()
3845 .now_or_never()
3846 .expect("should not block")
3847 .expect("would have errored");
3848 assert!(done.is_none(), "track completes once the boundary is reached");
3849 }
3850
3851 #[tokio::test]
3852 async fn recv_group_finishes_without_waiting_for_gaps() {
3853 let mut producer = track_producer("test", None);
3854 producer.create_group(group::Info { sequence: 1 }).unwrap();
3855 producer.finish().unwrap();
3856
3857 let mut consumer = producer.subscribe(None);
3858 assert_eq!(consumer.assert_group().sequence, 1);
3859
3860 let done = consumer
3861 .recv_group()
3862 .now_or_never()
3863 .expect("should not block")
3864 .expect("would have errored");
3865 assert!(done.is_none(), "track should finish without waiting for gaps");
3866 }
3867
3868 #[tokio::test]
3869 async fn next_group_skips_late_arrivals() {
3870 let mut producer = track_producer("test", None);
3871 let mut consumer = producer.subscribe(None);
3872
3873 producer.create_group(group::Info { sequence: 5 }).unwrap();
3875 let group = consumer
3876 .next_group()
3877 .now_or_never()
3878 .expect("should not block")
3879 .expect("would have errored")
3880 .expect("track should not be closed");
3881 assert_eq!(group.sequence, 5);
3882
3883 producer.create_group(group::Info { sequence: 3 }).unwrap();
3885 producer.create_group(group::Info { sequence: 4 }).unwrap();
3887 producer.create_group(group::Info { sequence: 7 }).unwrap();
3889
3890 let group = consumer
3891 .next_group()
3892 .now_or_never()
3893 .expect("should not block")
3894 .expect("would have errored")
3895 .expect("track should not be closed");
3896 assert_eq!(group.sequence, 7);
3897
3898 assert!(
3900 consumer.next_group().now_or_never().is_none(),
3901 "should block waiting for a higher sequence"
3902 );
3903 }
3904
3905 #[tokio::test]
3906 async fn next_group_returns_arrivals_in_order() {
3907 let mut producer = track_producer("test", None);
3908 let mut consumer = producer.subscribe(None);
3909
3910 producer.create_group(group::Info { sequence: 3 }).unwrap();
3912 producer.create_group(group::Info { sequence: 5 }).unwrap();
3913
3914 let group = consumer
3915 .next_group()
3916 .now_or_never()
3917 .expect("should not block")
3918 .expect("would have errored")
3919 .expect("track should not be closed");
3920 assert_eq!(group.sequence, 3);
3921
3922 let group = consumer
3923 .next_group()
3924 .now_or_never()
3925 .expect("should not block")
3926 .expect("would have errored")
3927 .expect("track should not be closed");
3928 assert_eq!(group.sequence, 5);
3929 }
3930
3931 #[tokio::test]
3932 async fn next_group_and_recv_group_use_independent_cursors() {
3933 let mut producer = track_producer("test", None);
3934 let mut consumer = producer.subscribe(None);
3935
3936 producer.create_group(group::Info { sequence: 5 }).unwrap();
3938 producer.create_group(group::Info { sequence: 3 }).unwrap();
3939
3940 let group = consumer
3943 .next_group()
3944 .now_or_never()
3945 .expect("should not block")
3946 .expect("would have errored")
3947 .expect("track should not be closed");
3948 assert_eq!(group.sequence, 3);
3949
3950 assert_eq!(consumer.assert_group().sequence, 5);
3953 }
3954
3955 #[tokio::test]
3956 async fn end_at_caps_next_group() {
3957 let mut producer = track_producer("test", None);
3958 let mut consumer = producer.subscribe(None);
3959
3960 for s in 0..6 {
3961 producer.create_group(group::Info { sequence: s }).unwrap();
3962 }
3963
3964 consumer.end_at(2);
3965
3966 assert_eq!(
3968 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3969 0
3970 );
3971 assert_eq!(
3972 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3973 1
3974 );
3975 assert_eq!(
3976 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3977 2
3978 );
3979
3980 assert!(
3982 consumer.next_group().now_or_never().is_none(),
3983 "capped consumer must block instead of returning out-of-range groups"
3984 );
3985 }
3986
3987 #[tokio::test]
3988 async fn end_at_release_drains_cached_groups() {
3989 let mut producer = track_producer("test", None);
3990 let mut consumer = producer.subscribe(None);
3991
3992 for s in 0..6 {
3993 producer.create_group(group::Info { sequence: s }).unwrap();
3994 }
3995
3996 consumer.end_at(1);
3997 assert_eq!(
3998 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3999 0
4000 );
4001 assert_eq!(
4002 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4003 1
4004 );
4005 assert!(consumer.next_group().now_or_never().is_none(), "capped at 1");
4006
4007 consumer.end_at(4);
4009 assert_eq!(
4010 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4011 2
4012 );
4013 assert_eq!(
4014 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4015 3
4016 );
4017 assert_eq!(
4018 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4019 4
4020 );
4021 assert!(consumer.next_group().now_or_never().is_none(), "capped at 4");
4022
4023 consumer.end_at(None);
4025 assert_eq!(
4026 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4027 5
4028 );
4029 assert!(consumer.next_group().now_or_never().is_none(), "no more groups");
4030 }
4031
4032 #[tokio::test]
4033 async fn end_at_lower_than_cursor_parks_consumer() {
4034 let mut producer = track_producer("test", None);
4035 let mut consumer = producer.subscribe(None);
4036
4037 for s in 0..3 {
4038 producer.create_group(group::Info { sequence: s }).unwrap();
4039 }
4040
4041 assert_eq!(
4043 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4044 0
4045 );
4046 assert_eq!(
4047 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4048 1
4049 );
4050 assert_eq!(
4051 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4052 2
4053 );
4054
4055 consumer.end_at(1);
4057 producer.create_group(group::Info { sequence: 3 }).unwrap();
4058 producer.create_group(group::Info { sequence: 4 }).unwrap();
4059 assert!(
4060 consumer.next_group().now_or_never().is_none(),
4061 "cap is below cursor; nothing returnable until cap rises"
4062 );
4063
4064 consumer.end_at(None);
4066 assert_eq!(
4067 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4068 3
4069 );
4070 assert_eq!(
4071 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4072 4
4073 );
4074 }
4075
4076 #[tokio::test]
4077 async fn end_at_toggling_around_late_arrivals() {
4078 let mut producer = track_producer("test", None);
4079 let mut consumer = producer.subscribe(None);
4080
4081 consumer.end_at(5);
4082
4083 producer.create_group(group::Info { sequence: 2 }).unwrap();
4085 producer.create_group(group::Info { sequence: 5 }).unwrap();
4086 producer.create_group(group::Info { sequence: 3 }).unwrap();
4087 producer.create_group(group::Info { sequence: 8 }).unwrap();
4089 producer.create_group(group::Info { sequence: 4 }).unwrap();
4090
4091 assert_eq!(
4093 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4094 2
4095 );
4096 assert_eq!(
4097 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4098 3
4099 );
4100 assert_eq!(
4101 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4102 4
4103 );
4104 assert_eq!(
4105 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4106 5
4107 );
4108 assert!(consumer.next_group().now_or_never().is_none());
4110
4111 consumer.end_at(10);
4113 assert_eq!(
4114 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4115 8
4116 );
4117 }
4118
4119 #[tokio::test]
4123 async fn end_at_parks_recv_group() {
4124 let mut producer = track_producer("test", None);
4125 let mut consumer = producer.subscribe(None);
4126
4127 for s in 0..3 {
4128 producer.create_group(group::Info { sequence: s }).unwrap();
4129 }
4130
4131 consumer.end_at(1);
4132 assert_eq!(
4133 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4134 0
4135 );
4136 assert_eq!(
4137 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4138 1
4139 );
4140 assert!(consumer.recv_group().now_or_never().is_none(), "capped at 1");
4141
4142 producer.finish().unwrap();
4144 assert!(
4145 consumer.recv_group().now_or_never().is_none(),
4146 "still parked after finish"
4147 );
4148
4149 consumer.end_at(None);
4150 assert_eq!(
4151 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4152 2
4153 );
4154 assert!(
4155 matches!(consumer.recv_group().now_or_never(), Some(Ok(None))),
4156 "finished once the parked group drains"
4157 );
4158 }
4159
4160 #[tokio::test]
4163 async fn recv_group_serves_arrivals_behind_the_cap() {
4164 let mut producer = track_producer("test", None);
4165 let mut consumer = producer.subscribe(None);
4166
4167 consumer.end_at(1);
4168
4169 producer.create_group(group::Info { sequence: 2 }).unwrap();
4171 producer.create_group(group::Info { sequence: 0 }).unwrap();
4172 producer.create_group(group::Info { sequence: 1 }).unwrap();
4173
4174 assert_eq!(
4175 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4176 0
4177 );
4178 assert_eq!(
4179 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4180 1
4181 );
4182 assert!(consumer.recv_group().now_or_never().is_none(), "capped at 1");
4183
4184 consumer.end_at(2);
4185 assert_eq!(
4186 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4187 2
4188 );
4189 }
4190
4191 #[tokio::test]
4194 async fn start_at_drops_parked_recv_groups() {
4195 let mut producer = track_producer("test", None);
4196 let mut consumer = producer.subscribe(None);
4197
4198 consumer.end_at(0);
4199 producer.create_group(group::Info { sequence: 1 }).unwrap();
4200 assert!(
4201 consumer.recv_group().now_or_never().is_none(),
4202 "group 1 parked at the cap"
4203 );
4204
4205 consumer.start_at(2);
4206 consumer.end_at(None);
4207 producer.create_group(group::Info { sequence: 2 }).unwrap();
4208 assert_eq!(
4209 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4210 2,
4211 "the overtaken parked group is dropped, not re-offered"
4212 );
4213 }
4214
4215 #[tokio::test]
4219 async fn evicted_parked_recv_groups_are_dropped() {
4220 let mut producer = track_producer("test", None);
4221 let mut consumer = producer.subscribe(None);
4222
4223 producer.create_group(group::Info { sequence: 0 }).unwrap();
4224 assert_eq!(
4225 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4226 0
4227 );
4228
4229 consumer.end_at(0);
4230 let straggler = producer.create_group(group::Info { sequence: 1 }).unwrap();
4231 assert!(
4232 consumer.recv_group().now_or_never().is_none(),
4233 "group 1 parked at the cap"
4234 );
4235
4236 straggler.abort(Error::Old).unwrap();
4238 producer.finish().unwrap();
4239
4240 consumer.end_at(None);
4241 assert!(
4242 matches!(consumer.recv_group().now_or_never(), Some(Ok(None))),
4243 "a dead parked group must not be delivered or hold the stream open"
4244 );
4245 }
4246
4247 #[tokio::test]
4251 async fn evicted_parked_group_wakes_the_clean_end() {
4252 use std::sync::atomic::{AtomicUsize, Ordering};
4253 use std::task::{Context, Wake};
4254
4255 struct CountWaker(AtomicUsize);
4258 impl Wake for CountWaker {
4259 fn wake(self: std::sync::Arc<Self>) {
4260 self.0.fetch_add(1, Ordering::SeqCst);
4261 }
4262 }
4263
4264 let mut producer = track_producer("test", None);
4265 let mut consumer = producer.subscribe(None);
4266
4267 producer.create_group(group::Info { sequence: 0 }).unwrap();
4268 assert_eq!(
4269 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4270 0
4271 );
4272
4273 consumer.end_at(0);
4274 let straggler = producer.create_group(group::Info { sequence: 1 }).unwrap();
4275 assert!(consumer.recv_group().now_or_never().is_none(), "parked at the cap");
4276 producer.finish().unwrap();
4277
4278 let counter = std::sync::Arc::new(CountWaker(AtomicUsize::new(0)));
4279 let waker = std::task::Waker::from(counter.clone());
4280 let mut cx = Context::from_waker(&waker);
4281 let mut fut = std::pin::pin!(consumer.recv_group());
4282 assert!(
4283 fut.as_mut().poll(&mut cx).is_pending(),
4284 "the parked group holds it open"
4285 );
4286
4287 straggler.abort(Error::Old).unwrap();
4288 assert!(counter.0.load(Ordering::SeqCst) > 0, "the eviction wakeup was lost");
4289 assert!(matches!(fut.as_mut().poll(&mut cx), Poll::Ready(Ok(None))));
4290 }
4291
4292 #[tokio::test]
4293 async fn read_frame_returns_single_frame_per_group() {
4294 let mut producer = track_producer("test", None);
4295 let mut consumer = producer.subscribe(None);
4296
4297 producer.write_frame(Timestamp::ZERO, b"hello".as_slice()).unwrap();
4298 producer.write_frame(Timestamp::ZERO, b"world".as_slice()).unwrap();
4299
4300 let frame = consumer
4301 .read_frame()
4302 .now_or_never()
4303 .expect("should not block")
4304 .expect("would have errored")
4305 .expect("track should not be closed");
4306 assert_eq!(&frame.payload[..], b"hello");
4307
4308 let frame = consumer
4309 .read_frame()
4310 .now_or_never()
4311 .expect("should not block")
4312 .expect("would have errored")
4313 .expect("track should not be closed");
4314 assert_eq!(&frame.payload[..], b"world");
4315 }
4316
4317 #[test]
4318 fn write_frame_rejects_an_oversized_frame_before_appending_its_group() {
4319 let mut producer = track_producer("test", None);
4320 let frame = bytes::Bytes::from(vec![0; group::MAX_CACHE_BYTES as usize + 1]);
4321
4322 assert!(matches!(
4323 producer.write_frame(Timestamp::ZERO, frame),
4324 Err(Error::FrameTooLarge)
4325 ));
4326 assert_eq!(producer.latest(), None, "the rejected frame did not publish a group");
4327 }
4328
4329 #[tokio::test]
4330 async fn read_frame_preserves_timestamp() {
4331 let mut producer = track_producer("test", None);
4332 let mut consumer = producer.subscribe(None);
4333
4334 producer
4335 .write_frame(Timestamp::from_micros(20_000).unwrap(), b"hello".as_slice())
4336 .unwrap();
4337
4338 let frame = consumer
4339 .read_frame()
4340 .now_or_never()
4341 .expect("should not block")
4342 .expect("would have errored")
4343 .expect("track should not be closed");
4344 assert_eq!(frame.timestamp.as_micros(), 20_000);
4345 assert_eq!(&frame.payload[..], b"hello");
4346 }
4347
4348 #[tokio::test]
4349 async fn read_frame_skips_stalled_group_for_newer_ready_frame() {
4350 let mut producer = track_producer("test", None);
4351 let mut consumer = producer.subscribe(None);
4352
4353 let _stalled = producer.create_group(group::Info { sequence: 3 }).unwrap();
4355 let mut g5 = producer.create_group(group::Info { sequence: 5 }).unwrap();
4357 g5.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"later"))
4358 .unwrap();
4359 g5.finish().unwrap();
4360
4361 let frame = consumer
4363 .read_frame()
4364 .now_or_never()
4365 .expect("should not block on stalled earlier group")
4366 .expect("would have errored")
4367 .expect("track should not be closed");
4368 assert_eq!(&frame.payload[..], b"later");
4369 }
4370
4371 #[tokio::test]
4372 async fn read_frame_discards_rest_of_multi_frame_group() {
4373 let mut producer = track_producer("test", None);
4374 let mut consumer = producer.subscribe(None);
4375
4376 let mut g0 = producer.create_group(group::Info { sequence: 0 }).unwrap();
4378 g0.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"one"))
4379 .unwrap();
4380 g0.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"two"))
4381 .unwrap();
4382 g0.finish().unwrap();
4383
4384 producer.write_frame(Timestamp::ZERO, b"next".as_slice()).unwrap();
4386
4387 let frame = consumer
4388 .read_frame()
4389 .now_or_never()
4390 .expect("should not block")
4391 .expect("would have errored")
4392 .expect("track should not be closed");
4393 assert_eq!(&frame.payload[..], b"one");
4394
4395 let frame = consumer
4397 .read_frame()
4398 .now_or_never()
4399 .expect("should not block")
4400 .expect("would have errored")
4401 .expect("track should not be closed");
4402 assert_eq!(&frame.payload[..], b"next");
4403 }
4404
4405 #[tokio::test]
4406 async fn read_frame_waits_for_pending_group_after_finish() {
4407 let mut producer = track_producer("test", None);
4410 let mut consumer = producer.subscribe(None);
4411
4412 let mut g0 = producer.create_group(group::Info { sequence: 0 }).unwrap();
4413 producer.finish().unwrap();
4414
4415 assert!(
4417 consumer.read_frame().now_or_never().is_none(),
4418 "read_frame must block on a pending group even after finish()"
4419 );
4420
4421 g0.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"late"))
4423 .unwrap();
4424 let frame = consumer
4425 .read_frame()
4426 .now_or_never()
4427 .expect("should not block once a frame is written")
4428 .expect("would have errored")
4429 .expect("track should not be closed");
4430 assert_eq!(&frame.payload[..], b"late");
4431 }
4432
4433 #[tokio::test]
4434 async fn read_frame_respects_start_at() {
4435 let mut producer = track_producer("test", None);
4438 let mut consumer = producer.subscribe(None);
4439 consumer.start_at(5);
4440
4441 let mut g3 = producer.create_group(group::Info { sequence: 3 }).unwrap();
4443 g3.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"skip-me"))
4444 .unwrap();
4445 g3.finish().unwrap();
4446
4447 let mut g5 = producer.create_group(group::Info { sequence: 5 }).unwrap();
4448 g5.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"keep"))
4449 .unwrap();
4450 g5.finish().unwrap();
4451
4452 let frame = consumer
4453 .read_frame()
4454 .now_or_never()
4455 .expect("should not block")
4456 .expect("would have errored")
4457 .expect("track should not be closed");
4458 assert_eq!(&frame.payload[..], b"keep");
4459 }
4460
4461 #[tokio::test]
4462 async fn read_frame_returns_none_when_finished() {
4463 let mut producer = track_producer("test", None);
4464 let mut consumer = producer.subscribe(None);
4465
4466 producer.write_frame(Timestamp::ZERO, b"only".as_slice()).unwrap();
4467 producer.finish().unwrap();
4468
4469 let frame = consumer
4470 .read_frame()
4471 .now_or_never()
4472 .expect("should not block")
4473 .expect("would have errored")
4474 .expect("track should not be closed");
4475 assert_eq!(&frame.payload[..], b"only");
4476
4477 let done = consumer
4478 .read_frame()
4479 .now_or_never()
4480 .expect("should not block")
4481 .expect("would have errored");
4482 assert!(done.is_none());
4483 }
4484
4485 #[test]
4486 fn append_group_returns_bounds_exceeded_on_sequence_overflow() {
4487 let mut producer = track_producer("test", None);
4488 {
4489 let mut state = producer.state.write().ok().unwrap();
4490 state.max_sequence = Some(u64::MAX);
4491 }
4492
4493 assert!(matches!(producer.append_group(), Err(Error::BoundsExceeded(_))));
4494 }
4495
4496 #[tokio::test]
4497 async fn fetch_cache_hit() {
4498 let mut producer = track_producer("test", None);
4499
4500 let mut group = producer.append_group().unwrap(); group
4503 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hello"))
4504 .unwrap();
4505 group.finish().unwrap();
4506
4507 let dynamic = producer.dynamic();
4510 let consumer = producer.consume();
4511 assert!(consumer.peek_group(0).is_some());
4512 let mut g = consumer.fetch_group(0, None).await.unwrap();
4513 assert_eq!(g.sequence, 0);
4514 assert_eq!(&g.read_frame().await.unwrap().unwrap().payload[..], b"hello");
4515
4516 assert!(dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending());
4518 }
4519
4520 #[tokio::test]
4521 async fn fetch_miss_signals_dynamic() {
4522 let producer = track_producer("test", None);
4523 let dynamic = producer.dynamic();
4524 let consumer = producer.consume();
4525
4526 assert!(consumer.peek_group(5).is_none());
4530 let pending = consumer.fetch_group(5, group::Fetch::default().with_priority(7));
4531 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
4532
4533 let req = dynamic
4534 .requested_group()
4535 .now_or_never()
4536 .expect("should not block")
4537 .unwrap();
4538 assert_eq!(req.sequence(), 5);
4539 assert_eq!(req.priority(), 7);
4540
4541 let mut group = req.accept(None).unwrap();
4543 group
4544 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hi"))
4545 .unwrap();
4546 group.finish().unwrap();
4547
4548 let mut g = pending.await.unwrap();
4549 assert_eq!(g.sequence, 5);
4550 assert_eq!(&g.read_frame().await.unwrap().unwrap().payload[..], b"hi");
4551 }
4552
4553 #[tokio::test]
4554 async fn fetch_miss_rejects() {
4555 let producer = track_producer("test", None);
4556 let dynamic = producer.dynamic();
4557 let consumer = producer.consume();
4558
4559 let pending = consumer.fetch_group(5, None);
4560 let req = dynamic
4561 .requested_group()
4562 .now_or_never()
4563 .expect("should not block")
4564 .unwrap();
4565
4566 req.reject(Error::Cancel);
4567 assert!(matches!(pending.await, Err(Error::Cancel)));
4568 let fetch = producer.state.read().fetch.clone();
4569 assert!(fetch.read().is_empty());
4570 }
4571
4572 #[tokio::test]
4573 async fn fetch_miss_drop_rejects() {
4574 let producer = track_producer("test", None);
4575 let dynamic = producer.dynamic();
4576 let consumer = producer.consume();
4577
4578 let pending = consumer.fetch_group(5, None);
4579 let req = dynamic
4580 .requested_group()
4581 .now_or_never()
4582 .expect("should not block")
4583 .unwrap();
4584
4585 drop(req);
4586 assert!(matches!(pending.await, Err(Error::Dropped)));
4587 }
4588
4589 #[tokio::test]
4590 async fn fetch_reject_does_not_poison_retry() {
4591 let producer = track_producer("test", None);
4592 let dynamic = producer.dynamic();
4593 let consumer = producer.consume();
4594
4595 let pending = consumer.fetch_group(5, None);
4596 let req = dynamic
4597 .requested_group()
4598 .now_or_never()
4599 .expect("should not block")
4600 .unwrap();
4601 req.reject(Error::Cancel);
4602 assert!(matches!(pending.await, Err(Error::Cancel)));
4603
4604 let retry = consumer.fetch_group(5, None);
4605 let req = dynamic
4606 .requested_group()
4607 .now_or_never()
4608 .expect("should not block")
4609 .unwrap();
4610 let mut group = req.accept(None).unwrap();
4611 group
4612 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"retry"))
4613 .unwrap();
4614 group.finish().unwrap();
4615
4616 let mut group = retry.await.unwrap();
4617 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"retry");
4618 }
4619
4620 #[tokio::test]
4621 async fn fetch_coalesces_concurrent() {
4622 let producer = track_producer("test", None);
4623 let dynamic = producer.dynamic();
4624 let consumer = producer.consume();
4625
4626 let first = consumer.fetch_group(5, group::Fetch::default().with_priority(1));
4629 let second = consumer.fetch_group(5, group::Fetch::default().with_priority(7));
4630 assert!(kio::Pollable::poll(&*first, &kio::Waiter::noop()).is_pending());
4631
4632 let req = dynamic
4633 .requested_group()
4634 .now_or_never()
4635 .expect("should not block")
4636 .unwrap();
4637 assert_eq!(req.sequence(), 5);
4638 assert_eq!(req.priority(), 7);
4639 assert!(
4640 dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending(),
4641 "the second fetch queued a duplicate request"
4642 );
4643
4644 let third = consumer.fetch_group(5, None);
4646
4647 let mut group = req.accept(None).unwrap();
4649 group
4650 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hi"))
4651 .unwrap();
4652 group.finish().unwrap();
4653
4654 assert_eq!(first.await.unwrap().sequence, 5);
4655 assert_eq!(second.await.unwrap().sequence, 5);
4656 assert_eq!(third.await.unwrap().sequence, 5);
4657 }
4658
4659 #[tokio::test]
4660 async fn fetch_coalesced_reject_fails_all() {
4661 let producer = track_producer("test", None);
4662 let dynamic = producer.dynamic();
4663 let consumer = producer.consume();
4664
4665 let first = consumer.fetch_group(5, None);
4666 let second = consumer.fetch_group(5, None);
4667 let req = dynamic
4668 .requested_group()
4669 .now_or_never()
4670 .expect("should not block")
4671 .unwrap();
4672 req.reject(Error::Cancel);
4673
4674 assert!(matches!(first.await, Err(Error::Cancel)));
4675 assert!(matches!(second.await, Err(Error::Cancel)));
4676
4677 let retry = consumer.fetch_group(5, None);
4679 assert!(kio::Pollable::poll(&*retry, &kio::Waiter::noop()).is_pending());
4680 let req = dynamic
4681 .requested_group()
4682 .now_or_never()
4683 .expect("should not block")
4684 .unwrap();
4685 assert_eq!(req.sequence(), 5);
4686 }
4687
4688 #[tokio::test]
4689 async fn fetch_queued_fails_when_handlers_leave() {
4690 let producer = track_producer("test", None);
4691 let dynamic = producer.dynamic();
4692 let consumer = producer.consume();
4693
4694 let pending = consumer.fetch_group(5, None);
4696 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
4697 drop(dynamic);
4698 assert!(matches!(pending.await, Err(Error::NotFound)));
4699
4700 let fetch = producer.state.read().fetch.clone();
4702 assert!(fetch.read().is_empty());
4703 }
4704
4705 #[tokio::test]
4706 async fn fetch_miss_no_dynamic_not_found() {
4707 let mut producer = track_producer("test", None);
4710 producer.append_group().unwrap(); let consumer = producer.consume();
4712 assert!(matches!(consumer.fetch_group(5, None).await, Err(Error::NotFound)));
4713 }
4714
4715 #[tokio::test]
4716 async fn fetch_past_final_not_found() {
4717 let mut producer = track_producer("test", None);
4718 producer.append_group().unwrap(); producer.finish().unwrap(); let dynamic = producer.dynamic();
4724 let consumer = producer.consume();
4725 assert!(matches!(consumer.fetch_group(5, None).await, Err(Error::NotFound)));
4726
4727 assert!(dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending());
4729 }
4730
4731 fn pooled_producer(capacity: u64) -> (Producer, cache::Pool) {
4733 let pool = cache::Pool::new(capacity);
4734 let broadcast = broadcast::Info {
4735 origin: crate::origin::Info::default().with_pool(pool.clone()),
4736 ..Default::default()
4737 };
4738 let producer = Producer::new(Arc::new(broadcast), "test", None);
4739 (producer, pool)
4740 }
4741
4742 fn finished_group(producer: &mut Producer, size: usize) -> u64 {
4743 let mut group = producer.append_group().unwrap();
4744 group
4745 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; size]))
4746 .unwrap();
4747 group.finish().unwrap();
4748 group.sequence
4749 }
4750
4751 #[tokio::test]
4754 async fn debt_evicts_oldest_group() {
4755 tokio::time::pause();
4756
4757 let (mut producer, pool) = pooled_producer(10_000);
4759
4760 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let consumer = producer.consume();
4765 assert!(consumer.peek_group(0).is_none(), "oldest group is evicted");
4766 assert!(consumer.peek_group(2).is_some(), "latest group survives");
4767 assert!(pool.used() <= 21_000, "usage hovers near capacity: {}", pool.used());
4770
4771 let mut subscriber = producer.subscribe(None);
4773 assert!(subscriber.assert_group().sequence > 0, "evicted group is not delivered");
4774 }
4775
4776 #[tokio::test]
4778 async fn latest_group_never_evicted() {
4779 tokio::time::pause();
4780
4781 let (mut producer, pool) = pooled_producer(100);
4783 finished_group(&mut producer, 1000); assert!(pool.used() > 100, "the latest may exceed the budget");
4785
4786 finished_group(&mut producer, 1000); finished_group(&mut producer, 1000); let consumer = producer.consume();
4791 assert!(consumer.peek_group(0).is_none());
4792 let mut group = consumer.peek_group(2).expect("latest survives");
4793 assert_eq!(group.read_frame().await.unwrap().unwrap().payload.len(), 1000);
4794 }
4795
4796 #[tokio::test]
4800 async fn fetch_refresh_survives_eviction() {
4801 tokio::time::pause();
4802
4803 let (mut producer, _pool) = pooled_producer(10_000);
4804 let consumer = producer.consume();
4805
4806 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_secs(1)).await;
4808 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_secs(1)).await;
4810 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_millis(500)).await;
4812
4813 let mut fetched = consumer.fetch_group(0, None).await.unwrap();
4815 assert_eq!(fetched.read_frame().await.unwrap().unwrap().payload.len(), 3_000);
4816 tokio::time::advance(Duration::from_millis(500)).await;
4817
4818 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_secs(1)).await;
4822 finished_group(&mut producer, 3_000); assert!(consumer.peek_group(0).is_some(), "refreshed group survives");
4825 assert!(consumer.peek_group(1).is_none(), "unread group is evicted instead");
4826 }
4827
4828 #[tokio::test]
4831 async fn eviction_aborts_readers() {
4832 tokio::time::pause();
4833
4834 let (mut producer, _pool) = pooled_producer(10_000);
4835 let mut subscriber = producer.subscribe(None);
4836
4837 finished_group(&mut producer, 10_000); let mut group0 = subscriber.assert_group();
4839
4840 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let read = group0.read_frame().await;
4844 assert!(matches!(read, Err(Error::Evicted)), "expected Evicted, got {read:?}");
4845 }
4846
4847 #[tokio::test]
4851 async fn small_writes_carry_debt() {
4852 tokio::time::pause();
4853
4854 let (mut producer, pool) = pooled_producer(22_000);
4855 let consumer = producer.consume();
4856
4857 finished_group(&mut producer, 20_000); for _ in 0..3 {
4862 finished_group(&mut producer, 1_000);
4863 }
4864 assert!(consumer.peek_group(0).is_some(), "debt smaller than the victim carries");
4865
4866 for _ in 0..20 {
4868 finished_group(&mut producer, 1_000);
4869 }
4870 assert!(
4871 consumer.peek_group(0).is_none(),
4872 "accumulated debt evicts the large group"
4873 );
4874 assert!(pool.used() <= 24_000, "usage hovers near capacity: {}", pool.used());
4877 }
4878
4879 #[tokio::test]
4883 async fn payment_capped_per_write() {
4884 tokio::time::pause();
4885
4886 let (mut producer, pool) = pooled_producer(1 << 40);
4887 for _ in 0..10 {
4888 finished_group(&mut producer, 1_000);
4889 }
4890
4891 pool.resize(100);
4893 let before = pool.used();
4894
4895 finished_group(&mut producer, 1_000);
4897
4898 let consumer = producer.consume();
4899 assert!(consumer.peek_group(0).is_none(), "the oldest groups are evicted");
4900 assert!(consumer.peek_group(1).is_none());
4901 assert!(consumer.peek_group(2).is_some(), "the backlog drains gradually");
4902 assert!(pool.used() > before - 4_000, "one write must not dump the backlog");
4903 }
4904
4905 #[tokio::test]
4909 async fn accept_preserves_write_accounting() {
4910 tokio::time::pause();
4911
4912 let pool = cache::Pool::new(12_000);
4913 let broadcast = broadcast::Info {
4914 origin: crate::origin::Info::default().with_pool(pool.clone()),
4915 ..Default::default()
4916 };
4917 let request = Request::new(Arc::new(broadcast), "test");
4918 let dynamic = request.dynamic();
4919 let consumer = request.consume();
4920
4921 let pending = consumer.fetch_group(0, None);
4923 let req = dynamic
4924 .requested_group()
4925 .now_or_never()
4926 .expect("should not block")
4927 .unwrap();
4928 let mut backfill = req.accept(None).unwrap();
4929 pending.await.unwrap();
4930 backfill
4931 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 30_000]))
4932 .unwrap();
4933
4934 let mut producer = request.accept(None);
4937 producer.append_group().unwrap().finish().unwrap();
4938 producer.append_group().unwrap().finish().unwrap();
4939
4940 assert!(
4941 producer.consume().peek_group(0).is_none(),
4942 "pre-accept backfill growth is reclaimed after accept"
4943 );
4944 assert!(pool.used() <= 13_000, "usage converges: {}", pool.used());
4945 }
4946
4947 #[tokio::test]
4950 async fn recreated_sequence_bounds_eviction_hints() {
4951 let (mut producer, _pool) = pooled_producer(1 << 40);
4952 producer.create_group(5u64.into()).unwrap().finish().unwrap();
4953
4954 for _ in 0..200 {
4955 let group = producer.create_group(1u64.into()).unwrap();
4956 group.abort(Error::Cancel).unwrap();
4957 }
4958
4959 let state = producer.state.read();
4960 assert!(
4961 state.evict.len() <= 2 * state.lookup.len() + EVICT_SLACK,
4962 "stale hints are compacted: {} entries for {} slots",
4963 state.evict.len(),
4964 state.lookup.len()
4965 );
4966 }
4967
4968 #[tokio::test]
4971 async fn same_tick_write_outranks_inserted() {
4972 tokio::time::pause();
4973
4974 let (mut producer, _pool) = pooled_producer(10_000);
4976
4977 producer.append_group().unwrap().finish().unwrap(); finished_group(&mut producer, 3_000); finished_group(&mut producer, 3_000); finished_group(&mut producer, 3_000); finished_group(&mut producer, 3_000); let consumer = producer.consume();
4984 assert!(consumer.peek_group(0).is_none(), "insert-only content pays first");
4985 assert!(consumer.peek_group(1).is_some(), "same-tick written content survives");
4986 }
4987
4988 #[tokio::test]
4991 async fn frame_only_writer_pays() {
4992 tokio::time::pause();
4993
4994 let (mut producer, pool) = pooled_producer(2_000);
4995 let mut demoted = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); demoted
5001 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 300_000]))
5002 .unwrap();
5003
5004 assert!(
5005 pool.used() <= 5_000,
5006 "the frame write settled the debt: {}",
5007 pool.used()
5008 );
5009 assert!(matches!(demoted.finish(), Err(Error::Evicted)));
5010 }
5011
5012 #[tokio::test]
5015 async fn each_track_owns_its_account() {
5016 let broadcast = Arc::new(broadcast::Info::default());
5017 let info = Info::default();
5018 let a = Producer::new(broadcast.clone(), "a", info.clone());
5019 let b = Producer::new(broadcast, "b", info);
5020
5021 let a = a.state.read().cache.clone();
5022 let b = b.state.read().cache.clone();
5023 assert!(!Arc::ptr_eq(&a, &b), "each track owns its account");
5024 }
5025
5026 #[tokio::test]
5029 async fn a_dynamic_defers_teardown() {
5030 let (mut producer, pool) = pooled_producer(1 << 40);
5031 let dynamic = producer.dynamic();
5032 finished_group(&mut producer, 100);
5033
5034 drop(producer);
5035 assert!(pool.used() > 0, "the handler still serves the cache");
5036
5037 drop(dynamic);
5038 assert_eq!(pool.used(), 0, "the last handle tears it down");
5039 }
5040
5041 #[tokio::test]
5047 async fn finished_track_frees_its_cache() {
5048 let (mut producer, pool) = pooled_producer(1 << 40);
5049 finished_group(&mut producer, 100);
5050 producer.finish().unwrap();
5051
5052 let state = producer.state.downgrade();
5053 drop(producer);
5054
5055 assert!(state.upgrade().is_none(), "the track state is freed");
5056 assert_eq!(pool.used(), 0, "so are its cached bytes");
5057 }
5058
5059 #[tokio::test]
5063 async fn teardown_ignores_a_settling_group() {
5064 let (mut producer, pool) = pooled_producer(1 << 40);
5065 finished_group(&mut producer, 100);
5066
5067 let settling = producer.state.downgrade().upgrade().expect("open");
5069 drop(producer);
5070
5071 assert_eq!(pool.used(), 0, "the abrupt teardown still released the cache");
5072 drop(settling);
5073 }
5074
5075 #[tokio::test]
5078 async fn cached_group_outlives_its_track() {
5079 let (mut producer, pool) = pooled_producer(1 << 40);
5080 let sequence = finished_group(&mut producer, 100);
5081 let group = producer.consume().peek_group(sequence).expect("cached");
5082 producer.finish().unwrap();
5083
5084 let state = producer.state.downgrade();
5085 drop(producer);
5086 assert!(state.upgrade().is_none(), "the track state is freed");
5087 assert!(pool.used() > 0, "the retained group keeps its own bytes");
5088
5089 drop(group);
5090 assert_eq!(pool.used(), 0, "which it releases when dropped");
5091 }
5092
5093 #[tokio::test]
5097 async fn pre_accept_backfill_settles_late_writes() {
5098 tokio::time::pause();
5099
5100 let pool = cache::Pool::new(2_000);
5101 let broadcast = broadcast::Info {
5102 origin: crate::origin::Info::default().with_pool(pool.clone()),
5103 ..Default::default()
5104 };
5105 let request = Request::new(Arc::new(broadcast), "test");
5106 let dynamic = request.dynamic();
5107 let consumer = request.consume();
5108
5109 let pending = consumer.fetch_group(0, None);
5111 let req = dynamic
5112 .requested_group()
5113 .now_or_never()
5114 .expect("should not block")
5115 .unwrap();
5116 let mut backfill = req.accept(None).unwrap();
5117 pending.await.unwrap();
5118
5119 let mut producer = request.accept(None);
5121 producer.append_group().unwrap().finish().unwrap();
5122
5123 backfill
5126 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 300_000]))
5127 .unwrap();
5128
5129 assert!(
5130 pool.used() <= 5_000,
5131 "the frame write settled the debt: {}",
5132 pool.used()
5133 );
5134 }
5135
5136 #[tokio::test]
5140 async fn write_restarts_retention_clock() {
5141 tokio::time::pause();
5142
5143 let (mut producer, _pool) = pooled_producer(1 << 40);
5144 let mut straggler = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
5149 straggler
5150 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 100]))
5151 .unwrap();
5152 producer.append_group().unwrap().finish().unwrap(); let consumer = producer.consume();
5155 assert!(consumer.peek_group(0).is_some(), "the write restarted the clock");
5156
5157 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
5159 producer.append_group().unwrap().finish().unwrap(); assert!(consumer.peek_group(0).is_none(), "idle content still expires");
5161 }
5162
5163 #[tokio::test]
5166 async fn refreshed_front_does_not_starve_expiry() {
5167 tokio::time::pause();
5168
5169 let (mut producer, _pool) = pooled_producer(1 << 40);
5170 let dynamic = producer.dynamic();
5171 let consumer = producer.consume();
5172
5173 producer.create_group(10u64.into()).unwrap().finish().unwrap();
5174 for sequence in 1..=5u64 {
5175 let pending = consumer.fetch_group(sequence, None);
5176 let req = dynamic
5177 .requested_group()
5178 .now_or_never()
5179 .expect("should not block")
5180 .unwrap();
5181 let mut group = req.accept(None).unwrap();
5182 group
5183 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 100]))
5184 .unwrap();
5185 group.finish().unwrap();
5186 pending.await.unwrap();
5187 }
5188
5189 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
5192 for sequence in 1..=4u64 {
5193 consumer.fetch_group(sequence, None).await.unwrap();
5194 }
5195
5196 for _ in 0..3 {
5198 producer.append_group().unwrap().finish().unwrap();
5199 }
5200 assert!(consumer.peek_group(5).is_none(), "expired backfill is reclaimed");
5201 assert!(consumer.peek_group(1).is_some(), "refreshed backfill survives");
5202 }
5203
5204 #[tokio::test]
5207 async fn recreated_sequence_delivered_once() {
5208 let (mut producer, _pool) = pooled_producer(1 << 40);
5209
5210 producer.create_group(0u64.into()).unwrap().finish().unwrap();
5211 let aborted = producer.create_group(1u64.into()).unwrap();
5212 aborted.abort(Error::Cancel).unwrap();
5213 producer.create_group(2u64.into()).unwrap().finish().unwrap();
5214 producer.create_group(1u64.into()).unwrap().finish().unwrap();
5215
5216 let mut subscriber = producer.subscribe(None);
5217 assert_eq!(subscriber.assert_group().sequence, 0);
5218 assert_eq!(subscriber.assert_group().sequence, 2);
5219 assert_eq!(
5220 subscriber.assert_group().sequence,
5221 1,
5222 "replacement arrives at its own position"
5223 );
5224 subscriber.assert_no_group();
5225 }
5226
5227 #[tokio::test]
5231 async fn datagrams_do_not_block_eviction() {
5232 tokio::time::pause();
5233
5234 let (mut producer, pool) = pooled_producer(1_000);
5235 for _ in 0..10 {
5236 finished_group(&mut producer, 1_000);
5237 producer.append_datagram(Timestamp::ZERO, &b"beat"[..]).unwrap();
5238 }
5239
5240 let consumer = producer.consume();
5241 assert!(consumer.peek_group(0).is_none(), "old groups still evict");
5242 assert!(
5243 pool.used() < 4 * 1_256,
5244 "interleaved datagrams must not bypass the budget: {}",
5245 pool.used()
5246 );
5247 }
5248
5249 #[tokio::test]
5253 async fn aborted_group_leaves_no_ghost_sample() {
5254 tokio::time::pause();
5255
5256 let (mut producer, pool) = pooled_producer(1 << 40);
5257 let group0 = producer.append_group().unwrap();
5258 producer.append_group().unwrap(); assert!(pool.average().is_some(), "demoted group is sampled");
5261 group0.abort(Error::Cancel).unwrap();
5262 assert_eq!(pool.average(), None, "the abort must remove the sample");
5263 }
5264
5265 #[tokio::test]
5268 async fn empty_groups_repay_overhead() {
5269 tokio::time::pause();
5270
5271 let (mut producer, pool) = pooled_producer(1_000);
5272 for _ in 0..100 {
5273 let mut group = producer.append_group().unwrap();
5274 group.finish().unwrap();
5275 }
5276
5277 assert!(
5278 pool.used() <= 3_000,
5279 "empty-group overhead must stay near the budget: {}",
5280 pool.used()
5281 );
5282 }
5283
5284 #[tokio::test]
5287 async fn growth_on_demoted_group_is_billed() {
5288 tokio::time::pause();
5289
5290 let (mut producer, pool) = pooled_producer(2_000);
5291 let mut straggler = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); straggler
5296 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 10_000]))
5297 .unwrap();
5298
5299 producer.append_group().unwrap().finish().unwrap(); let consumer = producer.consume();
5303 assert!(consumer.peek_group(0).is_none(), "the ballooned group is evicted");
5304 assert!(pool.used() <= 3_000, "growth is reclaimed: {}", pool.used());
5305 }
5306
5307 #[tokio::test]
5310 async fn refilled_sequence_stays_out_of_subscriptions() {
5311 let (mut producer, _pool) = pooled_producer(1 << 40);
5312 let dynamic = producer.dynamic();
5313 let consumer = producer.consume();
5314
5315 producer.create_group(0u64.into()).unwrap().finish().unwrap();
5316 let aborted = producer.create_group(1u64.into()).unwrap();
5317 aborted.abort(Error::Cancel).unwrap();
5318 producer.create_group(2u64.into()).unwrap().finish().unwrap();
5319
5320 let pending = consumer.fetch_group(1, None);
5323 let req = dynamic
5324 .requested_group()
5325 .now_or_never()
5326 .expect("should not block")
5327 .unwrap();
5328 let mut group = req.accept(None).unwrap();
5329 group
5330 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"backfill"))
5331 .unwrap();
5332 group.finish().unwrap();
5333 pending.await.unwrap();
5334
5335 assert!(consumer.peek_group(1).is_some());
5337 let mut subscriber = producer.subscribe(None);
5338 assert_eq!(subscriber.assert_group().sequence, 0);
5339 assert_eq!(subscriber.assert_group().sequence, 2);
5340 subscriber.assert_no_group();
5341 }
5342
5343 #[tokio::test]
5346 async fn expired_backfill_behind_refreshed_reclaimed() {
5347 tokio::time::pause();
5348
5349 let (mut producer, _pool) = pooled_producer(1 << 40);
5350 let dynamic = producer.dynamic();
5351 let consumer = producer.consume();
5352
5353 producer.create_group(5u64.into()).unwrap().finish().unwrap();
5354 for sequence in [2u64, 3u64] {
5355 let pending = consumer.fetch_group(sequence, None);
5356 let req = dynamic
5357 .requested_group()
5358 .now_or_never()
5359 .expect("should not block")
5360 .unwrap();
5361 let mut group = req.accept(None).unwrap();
5362 group
5363 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 100]))
5364 .unwrap();
5365 group.finish().unwrap();
5366 pending.await.unwrap();
5367 }
5368
5369 tokio::time::advance(Duration::from_secs(4)).await;
5371 consumer.fetch_group(2, None).await.unwrap();
5372 tokio::time::advance(DEFAULT_LATENCY_MAX - Duration::from_secs(2)).await;
5373 producer.create_group(6u64.into()).unwrap().finish().unwrap();
5374
5375 let consumer = producer.consume();
5376 assert!(consumer.peek_group(2).is_some(), "refreshed backfill survives");
5377 assert!(consumer.peek_group(3).is_none(), "expired backfill is reclaimed");
5378 }
5379
5380 #[tokio::test]
5383 async fn same_tick_fetch_protects() {
5384 tokio::time::pause();
5385
5386 let (mut producer, _pool) = pooled_producer(10_000);
5388 let consumer = producer.consume();
5389
5390 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();
5395
5396 finished_group(&mut producer, 3_000); finished_group(&mut producer, 3_000); assert!(consumer.peek_group(0).is_some(), "same-tick refresh protects");
5400 assert!(consumer.peek_group(1).is_none(), "the unread group dies instead");
5401 }
5402
5403 #[tokio::test]
5407 async fn refetched_latest_stays_protected() {
5408 tokio::time::pause();
5409
5410 let (mut producer, _pool) = pooled_producer(10_000);
5411 let dynamic = producer.dynamic();
5412 let consumer = producer.consume();
5413
5414 let straggler = producer.append_group().unwrap(); let latest = producer.append_group().unwrap(); latest.abort(Error::Cancel).unwrap();
5419
5420 let pending = consumer.fetch_group(1, None);
5422 let req = dynamic
5423 .requested_group()
5424 .now_or_never()
5425 .expect("should not block")
5426 .unwrap();
5427 let mut group = req.accept(None).unwrap();
5428 group
5429 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 1000]))
5430 .unwrap();
5431 group.finish().unwrap();
5432 pending.await.unwrap();
5433
5434 {
5437 let state = producer.state.read();
5438 assert!(state.lookup.contains_key(&1), "refetched group is cached");
5439 assert!(
5440 state.evict.iter().all(|(sequence, _)| *sequence != 1),
5441 "the live edge must not be an eviction candidate"
5442 );
5443 }
5444 drop(straggler);
5445 }
5446
5447 #[tokio::test]
5450 async fn eviction_allows_refetch() {
5451 tokio::time::pause();
5452
5453 let (mut producer, _pool) = pooled_producer(10_000);
5454 let dynamic = producer.dynamic();
5455
5456 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let consumer = producer.consume();
5461 assert!(consumer.peek_group(0).is_none());
5462 let pending = consumer.fetch_group(0, None);
5463
5464 let req = dynamic
5465 .requested_group()
5466 .now_or_never()
5467 .expect("should not block")
5468 .unwrap();
5469 assert_eq!(req.sequence(), 0);
5470
5471 let mut group = req.accept(None).unwrap();
5472 group
5473 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"refetched"))
5474 .unwrap();
5475 group.finish().unwrap();
5476
5477 let mut group = pending.await.unwrap();
5478 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"refetched");
5479 }
5480
5481 #[tokio::test]
5484 async fn fetched_backfill_not_subscribed() {
5485 let (mut producer, _pool) = pooled_producer(1 << 40);
5486 let dynamic = producer.dynamic();
5487 let consumer = producer.consume();
5488
5489 producer.create_group(5u64.into()).unwrap().finish().unwrap();
5491 producer.create_group(6u64.into()).unwrap().finish().unwrap();
5492
5493 let pending = consumer.fetch_group(2, None);
5495 let req = dynamic
5496 .requested_group()
5497 .now_or_never()
5498 .expect("should not block")
5499 .unwrap();
5500 let mut group = req.accept(None).unwrap();
5501 group
5502 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"backfill"))
5503 .unwrap();
5504 group.finish().unwrap();
5505 let mut fetched = pending.await.unwrap();
5506 assert_eq!(&fetched.read_frame().await.unwrap().unwrap().payload[..], b"backfill");
5507 assert!(consumer.peek_group(2).is_some(), "backfill is cached for later fetches");
5508
5509 let mut subscriber = producer.subscribe(None);
5511 assert_eq!(subscriber.assert_group().sequence, 5);
5512 assert_eq!(subscriber.assert_group().sequence, 6);
5513 subscriber.assert_no_group();
5514 }
5515
5516 #[tokio::test]
5519 async fn expired_backfill_reclaimed() {
5520 tokio::time::pause();
5521
5522 let (mut producer, pool) = pooled_producer(1 << 40);
5523 let dynamic = producer.dynamic();
5524 let consumer = producer.consume();
5525
5526 producer.create_group(5u64.into()).unwrap().finish().unwrap();
5527
5528 let pending = consumer.fetch_group(2, None);
5530 let req = dynamic
5531 .requested_group()
5532 .now_or_never()
5533 .expect("should not block")
5534 .unwrap();
5535 let mut group = req.accept(None).unwrap();
5536 group
5537 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 1000]))
5538 .unwrap();
5539 group.finish().unwrap();
5540 pending.await.unwrap();
5541 let used = pool.used();
5542
5543 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
5545 producer.create_group(6u64.into()).unwrap().finish().unwrap();
5546
5547 assert!(consumer.peek_group(2).is_none(), "expired backfill is reclaimed");
5548 assert!(pool.used() < used, "its bytes are released");
5549 }
5550
5551 #[tokio::test]
5552 async fn fetch_aborts_with_track() {
5553 let producer = track_producer("test", None);
5554 let dynamic = producer.dynamic();
5555 let consumer = producer.consume();
5556
5557 let pending = consumer.fetch_group(3, None);
5558 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
5559
5560 producer.abort(Error::Cancel).unwrap();
5561 assert!(pending.await.is_err());
5562 drop(dynamic);
5563 }
5564}