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 published: bool,
139
140 broadcast: Arc<broadcast::Info>,
143
144 cache: Arc<cache::Track>,
148
149 lookup: BTreeMap<u64, Slot>,
157
158 arrival: VecDeque<(u64, u32)>,
163
164 evict: VecDeque<(u64, u32)>,
172
173 debt: u64,
178
179 datagrams: VecDeque<(Datagram, web_async::time::Instant)>,
183
184 datagram_offset: usize,
187
188 offset: usize,
191
192 max_sequence: Option<u64>,
195
196 latest_group: Option<u64>,
201
202 next_stamp: u32,
204
205 expire_cursor: usize,
208
209 final_sequence: Option<u64>,
211
212 abort: Option<Error>,
214
215 subscriptions: kio::Shared<Subscriptions>,
219
220 fetch: kio::Shared<FetchState>,
223}
224
225struct Slot {
231 group: group::Producer,
232
233 stamp: u32,
238}
239
240pub(crate) const CACHE_OVERHEAD: u64 = 2 * (size_of::<u64>() + size_of::<Slot>() + 2 * size_of::<(u64, u32)>()) as u64;
249
250type Subscriptions = Vec<kio::Consumer<Subscription>>;
252
253type FetchState = Requests<u64, PendingFetch>;
258
259struct PendingFetch {
261 priority: u8,
263
264 result: kio::Producer<FetchOutcome>,
269}
270
271#[derive(Default)]
274struct FetchOutcome {
275 rejected: Option<Error>,
276}
277
278impl TrackState {
279 fn normalize_info(broadcast: &broadcast::Info, mut info: Info) -> Info {
280 info.latency_max = info.latency_max.min(broadcast.origin.cache_duration);
281 info
282 }
283
284 fn accept(&mut self, info: Info) {
285 self.published = true;
286 self.install(info);
287 }
288
289 fn poll_info(&self) -> Poll<Result<Info>> {
290 if let Some(info) = &self.info {
291 Poll::Ready(Ok(info.clone()))
292 } else if let Some(err) = &self.abort {
293 Poll::Ready(Err(err.clone()))
296 } else {
297 Poll::Pending
298 }
299 }
300
301 fn poll_recv_group(&self, index: usize, min_sequence: u64) -> Poll<Result<Option<(group::Consumer, usize)>>> {
305 let start = index.saturating_sub(self.offset);
306 for (i, (sequence, stamp)) in self.arrival.iter().enumerate().skip(start) {
307 if *sequence >= min_sequence
308 && let Some(slot) = self.lookup.get(sequence)
309 && slot.stamp == *stamp
310 && !slot.group.is_aborted()
311 {
312 slot.group.cache_refresh();
315 return Poll::Ready(Ok(Some((slot.group.consume(), self.offset + i))));
316 }
317 }
318
319 if self.is_complete() {
321 Poll::Ready(Ok(None))
322 } else if let Some(err) = &self.abort {
323 Poll::Ready(Err(err.clone()))
324 } else {
325 Poll::Pending
326 }
327 }
328
329 fn poll_recv_datagram(&self, index: usize) -> Poll<Result<Option<(Datagram, usize)>>> {
335 let start = index.saturating_sub(self.datagram_offset);
336 if let Some((datagram, _)) = self.datagrams.get(start) {
337 return Poll::Ready(Ok(Some((datagram.clone(), self.datagram_offset + start))));
338 }
339
340 if self.is_complete() {
342 Poll::Ready(Ok(None))
343 } else if let Some(err) = &self.abort {
344 Poll::Ready(Err(err.clone()))
345 } else {
346 Poll::Pending
347 }
348 }
349
350 fn push_datagram(&mut self, datagram: Datagram) {
352 let now = web_async::time::Instant::now();
353 self.datagrams.push_back((datagram, now));
354 while let Some((_, at)) = self.datagrams.front() {
355 if now.duration_since(*at) <= MAX_DATAGRAM_AGE {
356 break;
357 }
358 self.datagrams.pop_front();
359 self.datagram_offset += 1;
360 }
361 }
362
363 fn poll_read_frame(
367 &self,
368 index: usize,
369 next_sequence: u64,
370 waiter: &kio::Waiter,
371 ) -> Poll<Result<Option<(frame::Frame, usize, u64)>>> {
372 let start = index.saturating_sub(self.offset);
373 let mut pending_seen = false;
374 for (i, (sequence, stamp)) in self.arrival.iter().enumerate().skip(start) {
375 if *sequence < next_sequence {
376 continue;
377 }
378 let Some(slot) = self.lookup.get(sequence) else {
379 continue;
380 };
381 if slot.stamp != *stamp {
382 continue;
385 }
386
387 let mut consumer = slot.group.consume();
388 match consumer.poll_read_frame(waiter) {
389 Poll::Ready(Ok(Some(frame))) => {
390 return Poll::Ready(Ok(Some((frame, self.offset + i, *sequence))));
391 }
392 Poll::Ready(Ok(None)) => continue,
393 Poll::Ready(Err(_)) => continue,
396 Poll::Pending => {
397 pending_seen = true;
398 continue;
399 }
400 }
401 }
402
403 if pending_seen {
406 Poll::Pending
407 } else if self.is_complete() {
408 Poll::Ready(Ok(None))
409 } else if let Some(err) = &self.abort {
410 Poll::Ready(Err(err.clone()))
411 } else {
412 Poll::Pending
413 }
414 }
415
416 fn poll_next_in_range(
426 &self,
427 next_sequence: u64,
428 end_sequence: Option<u64>,
429 ) -> Poll<Result<Option<group::Consumer>>> {
430 if let Some(end) = end_sequence
434 && end < next_sequence
435 {
436 if let Some(err) = &self.abort {
437 return Poll::Ready(Err(err.clone()));
438 }
439 return Poll::Pending;
440 }
441
442 let best = self
445 .lookup
446 .range(next_sequence..)
447 .map(|(_, slot)| &slot.group)
448 .take_while(|group| end_sequence.is_none_or(|end| group.sequence <= end))
449 .find(|group| !group.is_aborted());
450
451 if let Some(group) = best {
452 group.cache_refresh();
454 return Poll::Ready(Ok(Some(group.consume())));
455 }
456
457 if let Some(err) = &self.abort {
459 return Poll::Ready(Err(err.clone()));
460 }
461 if let Some(fin) = self.final_sequence
464 && next_sequence >= fin
465 {
466 return Poll::Ready(Ok(None));
467 }
468 Poll::Pending
469 }
470
471 fn latency_bound(&self) -> Option<Duration> {
474 self.info.as_ref().map(|info| info.latency_max)
475 }
476
477 fn poll_fetch_cached(&self, sequence: u64) -> Poll<Result<group::Consumer>> {
482 if let Some(slot) = self.lookup.get(&sequence)
483 && !slot.group.is_aborted()
484 {
485 slot.group.cache_refresh();
489 return Poll::Ready(Ok(slot.group.consume()));
490 }
491
492 if let Some(err) = &self.abort {
493 return Poll::Ready(Err(err.clone()));
494 }
495
496 if self.final_sequence.is_some_and(|fin| sequence >= fin) {
498 return Poll::Ready(Err(Error::NotFound));
499 }
500
501 Poll::Pending
502 }
503
504 fn evict_expired(&mut self, max_age: Duration) {
513 let now = self.cache.pool().now();
514 let max_ticks = cache::Pool::ticks(max_age);
515
516 let len = self.evict.len();
517 if len > 0 {
518 let start = self.expire_cursor % len;
519 for step in 0..len.min(EVICT_SCAN) {
520 let (sequence, stamp) = self.evict[(start + step) % len];
521 let Some(slot) = self.lookup.get(&sequence) else {
522 continue;
523 };
524 if slot.stamp != stamp {
525 continue;
527 }
528 if slot.group.is_aborted() {
531 self.lookup.remove(&sequence);
532 continue;
533 }
534 if Some(sequence) == self.latest_group || now.saturating_sub(slot.group.cache_accessed()) <= max_ticks {
535 continue;
536 }
537 let slot = self.lookup.remove(&sequence).unwrap();
541 let _ = slot.group.abort(Error::Old);
542 }
543 self.expire_cursor = (start + EVICT_SCAN) % len;
544 }
545
546 while let Some((sequence, stamp)) = self.arrival.front() {
549 if self.lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp) {
550 break;
551 }
552 self.arrival.pop_front();
553 self.offset += 1;
554 }
555
556 while let Some((sequence, stamp)) = self.evict.front() {
558 if self.lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp) {
559 break;
560 }
561 self.evict.pop_front();
562 }
563
564 if self.evict.len() > 2 * self.lookup.len() + EVICT_SLACK {
567 let lookup = &self.lookup;
568 self.evict
569 .retain(|(sequence, stamp)| lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp));
570 }
571 }
572
573 fn clear_cache(&mut self) {
576 self.lookup.clear();
577 self.arrival.clear();
578 self.evict.clear();
579 self.latest_group = None;
580 self.debt = 0;
581 }
582
583 fn install(&mut self, info: Info) {
589 let info = Self::normalize_info(&self.broadcast, info);
590 self.info = Some(info);
591 }
592
593 fn spawn(broadcast: Arc<broadcast::Info>) -> kio::Producer<Self> {
600 let state = kio::Producer::new(Self {
601 broadcast: broadcast.clone(),
602 ..Default::default()
603 });
604 let cache = cache::Track::new(broadcast.origin.pool.clone(), state.downgrade());
605 state.write().ok().expect("a new track is open").cache = cache;
606 state
607 }
608
609 fn claim_sequence(&mut self, sequence: u64) -> Result<()> {
615 if let Some(slot) = self.lookup.get(&sequence) {
616 if !slot.group.is_aborted() {
617 return Err(Error::Duplicate);
618 }
619 self.lookup.remove(&sequence);
620 }
621 Ok(())
622 }
623
624 fn insert_group(&mut self, group: &group::Producer, visible: bool) {
631 let sequence = group.sequence;
632 self.next_stamp = self.next_stamp.wrapping_add(1);
633 let stamp = self.next_stamp;
634
635 if self.latest_group.is_none_or(|latest| sequence >= latest) {
639 if let Some(latest) = self.latest_group
642 && sequence > latest
643 && let Some(prev) = self.lookup.get(&latest)
644 {
645 prev.group.cache_demote();
646 self.evict.push_back((latest, prev.stamp));
647 }
648 self.latest_group = Some(sequence);
649 } else {
650 group.cache_demote();
651 self.evict.push_back((sequence, stamp));
652 }
653
654 self.max_sequence = Some(self.max_sequence.map_or(sequence, |max| max.max(sequence)));
655 self.lookup.insert(
656 sequence,
657 Slot {
658 group: group.clone(),
659 stamp,
660 },
661 );
662 if visible {
663 self.arrival.push_back((sequence, stamp));
664 }
665 }
666
667 fn commit_group(&mut self, group: &group::Producer, visible: bool, latency_max: Duration) {
671 self.charge_debt();
672 self.insert_group(group, visible);
673 self.evict_expired(latency_max);
674 }
675
676 pub(super) fn charge_debt(&mut self) {
689 let written = self.cache.take_written();
690 let pool = self.cache.pool().clone();
691 match pool.accrue(written) {
692 Some(mut accrued) => {
693 if self.oldest_is_stale(&pool) {
694 accrued = accrued.saturating_mul(2);
695 }
696 self.debt = self.debt.saturating_add(accrued).min(pool.used());
699 self.pay_debt(&pool, written.saturating_mul(2));
702 }
703 None => self.debt = 0,
706 }
707 }
708
709 fn oldest_is_stale(&self, pool: &cache::Pool) -> bool {
713 let Some(average) = pool.average() else {
714 return false;
715 };
716 let Some((sequence, stamp)) = self.evict.front() else {
717 return false;
718 };
719 let Some(slot) = self.lookup.get(sequence) else {
720 return false;
721 };
722 slot.stamp == *stamp && !slot.group.is_aborted() && slot.group.cache_accessed() <= average
723 }
724
725 fn pay_debt(&mut self, pool: &cache::Pool, cap: u64) {
738 let average = pool.average().unwrap_or(0);
739 let mut paid = 0u64;
740 let mut scanned = 0usize;
741 for _ in 0..self.evict.len() {
742 if self.debt == 0 || paid >= cap || scanned >= EVICT_SCAN {
743 return;
744 }
745 let Some((sequence, stamp)) = self.evict.pop_front() else {
746 return;
747 };
748 let Some(slot) = self.lookup.get(&sequence) else {
749 continue;
751 };
752 if slot.stamp != stamp {
753 continue;
755 }
756 if slot.group.is_aborted() {
757 self.lookup.remove(&sequence);
759 continue;
760 }
761 if Some(sequence) == self.latest_group {
762 self.evict.push_back((sequence, stamp));
764 continue;
765 }
766
767 scanned += 1;
768 if slot.group.cache_accessed() > average {
772 self.evict.push_back((sequence, stamp));
773 continue;
774 }
775 let size = slot.group.cache_size();
778 if size > self.debt {
779 self.evict.push_front((sequence, stamp));
780 return;
781 }
782
783 self.debt -= size;
784 paid = paid.saturating_add(size);
785 let slot = self.lookup.remove(&sequence).unwrap();
786 let _ = slot.group.abort(Error::Evicted);
787 }
788 }
789
790 fn set_final(&mut self, final_sequence: u64) -> Result<()> {
793 if self.final_sequence.is_some() {
794 return Err(Error::Closed);
795 }
796 if let Some(max) = self.max_sequence
797 && final_sequence <= max
798 {
799 return Err(Error::ProtocolViolation);
800 }
801 self.final_sequence = Some(final_sequence);
802 Ok(())
803 }
804
805 fn is_complete(&self) -> bool {
811 self.final_sequence
812 .is_some_and(|fin| self.max_sequence.map_or(0, |max| max.saturating_add(1)) >= fin)
813 }
814
815 fn poll_finished(&self) -> Poll<Result<u64>> {
816 if let Some(fin) = self.final_sequence {
817 Poll::Ready(Ok(fin))
818 } else if let Some(err) = &self.abort {
819 Poll::Ready(Err(err.clone()))
820 } else {
821 Poll::Pending
822 }
823 }
824
825 fn modify(producer: &kio::Producer<Self>) -> Result<kio::Mut<'_, Self>> {
826 producer.write().map_err(|r| r.abort.clone().unwrap_or(Error::Dropped))
827 }
828
829 fn insert_group_request(&mut self, sequence: u64, info: Option<Info>) -> Result<group::Producer> {
835 if let Some(err) = &self.abort {
836 return Err(err.clone());
837 }
838 if let Some(fin) = self.final_sequence
839 && sequence >= fin
840 {
841 return Err(Error::Closed);
842 }
843
844 if self.info.is_none() {
848 self.install(info.unwrap_or_default());
849 }
850 let info = self.info.clone().unwrap();
851
852 self.claim_sequence(sequence)?;
854
855 let latency_max = info.latency_max;
856 let group = group::Producer::new(group::Info { sequence }, info, self.cache.clone());
857 group.cache_refresh();
862 self.commit_group(&group, false, latency_max);
863 Ok(group)
864 }
865}
866
867#[derive(Clone)]
869pub struct Producer {
870 name: Arc<str>,
871 info: Info,
872 broadcast: Arc<broadcast::Info>,
875 state: kio::Producer<TrackState>,
876 prev_subscription: Option<Subscription>,
877 alive: Arc<Alive>,
879 stats: stats::Scope,
883}
884
885impl Producer {
886 pub(crate) fn new(
894 broadcast: Arc<broadcast::Info>,
895 name: impl Into<Arc<str>>,
896 info: impl Into<Option<Info>>,
897 ) -> Self {
898 let name = name.into();
899 let info = TrackState::normalize_info(&broadcast, info.into().unwrap_or_default());
900 let state = TrackState::spawn(broadcast.clone());
901 state.write().ok().expect("a new track is open").accept(info.clone());
902 let alive = Alive::new(name.clone(), state.clone());
903 alive.publish(None);
904 Self {
905 name,
906 info,
907 state,
908 broadcast,
909 prev_subscription: None,
910 alive,
911 stats: stats::Scope::default(),
912 }
913 }
914
915 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
919 self.alive.publish(Some(&scope));
920 self.stats = scope;
921 self
922 }
923
924 pub fn name(&self) -> &str {
926 &self.name
927 }
928
929 pub fn broadcast(&self) -> &broadcast::Info {
931 &self.broadcast
932 }
933
934 pub fn create_group(&mut self, group: group::Info) -> Result<group::Producer> {
936 let mut state = self.modify()?;
937 if let Some(fin) = state.final_sequence
938 && group.sequence >= fin
939 {
940 return Err(Error::Closed);
941 }
942 let track = state.info.clone().unwrap();
943 let latency_max = track.latency_max;
944
945 state.claim_sequence(group.sequence)?;
947
948 let group = group::Producer::new(group, track, state.cache.clone()).with_meter(self.stats.meter());
949 state.commit_group(&group, true, latency_max);
950
951 Ok(group)
952 }
953
954 pub fn append_group(&mut self) -> Result<group::Producer> {
956 let mut state = self.modify()?;
957 let sequence = match state.max_sequence {
958 Some(s) => s.checked_add(1).ok_or(coding::BoundsExceeded)?,
959 None => 0,
960 };
961 if let Some(fin) = state.final_sequence
962 && sequence >= fin
963 {
964 return Err(Error::Closed);
965 }
966
967 let track = state.info.clone().unwrap();
968 let latency_max = track.latency_max;
969
970 let group =
971 group::Producer::new(group::Info { sequence }, track, state.cache.clone()).with_meter(self.stats.meter());
972 state.commit_group(&group, true, latency_max);
973
974 Ok(group)
975 }
976
977 pub fn append_datagram<B: crate::IntoBytes>(&mut self, timestamp: Timestamp, payload: B) -> Result<u64> {
989 let payload = payload.into_bytes();
990 if payload.len() > super::datagram::MAX_DATAGRAM_PAYLOAD {
991 return Err(Error::FrameTooLarge);
992 }
993 let meter = self.stats.meter();
995 let mut state = self.modify()?;
996 let timescale = state.info.as_ref().unwrap().timescale;
998 let timestamp = timestamp.convert(timescale).map_err(|_| Error::TimestampMismatch)?;
999 let sequence = match state.max_sequence {
1000 Some(s) => s.checked_add(1).ok_or(coding::BoundsExceeded)?,
1001 None => 0,
1002 };
1003 if let Some(fin) = state.final_sequence
1004 && sequence >= fin
1005 {
1006 return Err(Error::Closed);
1007 }
1008 state.max_sequence = Some(sequence);
1009 meter.datagram(payload.len() as u64);
1010 state.push_datagram(Datagram {
1011 sequence,
1012 timestamp,
1013 payload,
1014 });
1015 Ok(sequence)
1016 }
1017
1018 pub fn write_datagram(&mut self, mut datagram: Datagram) -> Result<()> {
1024 if datagram.payload.len() > super::datagram::MAX_DATAGRAM_PAYLOAD {
1025 return Err(Error::FrameTooLarge);
1026 }
1027 let meter = self.stats.meter();
1029 let mut state = self.modify()?;
1030 let timescale = state.info.as_ref().unwrap().timescale;
1032 datagram.timestamp = datagram
1033 .timestamp
1034 .convert(timescale)
1035 .map_err(|_| Error::TimestampMismatch)?;
1036 if let Some(fin) = state.final_sequence
1037 && datagram.sequence >= fin
1038 {
1039 return Err(Error::Closed);
1040 }
1041 state.max_sequence = Some(state.max_sequence.unwrap_or(0).max(datagram.sequence));
1042 meter.datagram(datagram.payload.len() as u64);
1043 state.push_datagram(datagram);
1044 Ok(())
1045 }
1046
1047 pub fn write_frame<B: crate::IntoBytes>(&mut self, timestamp: Timestamp, frame: B) -> Result<()> {
1052 let frame = crate::IntoBytes::into_bytes(frame);
1053 if frame.len() as u64 > group::MAX_CACHE_BYTES {
1054 return Err(Error::FrameTooLarge);
1055 }
1056 let mut group = self.append_group()?;
1057 group.write_frame(timestamp, frame)?;
1058 group.finish()?;
1059 Ok(())
1060 }
1061
1062 pub fn finish(&mut self) -> Result<()> {
1068 let mut state = self.modify()?;
1069 let final_sequence = match state.max_sequence {
1070 Some(max) => max.checked_add(1).ok_or(coding::BoundsExceeded)?,
1071 None => 0,
1072 };
1073 state.set_final(final_sequence)
1074 }
1075
1076 pub fn finish_at(&mut self, final_sequence: u64) -> Result<()> {
1089 self.modify()?.set_final(final_sequence)
1090 }
1091
1092 pub fn final_sequence(&self) -> Option<u64> {
1097 self.state.read().final_sequence
1098 }
1099
1100 pub fn abort(self, err: Error) -> Result<()> {
1111 let mut guard = self.modify()?;
1112 guard.abort = Some(err);
1113 guard.clear_cache();
1114 guard.datagrams.clear();
1115 guard.close();
1116 Ok(())
1117 }
1118
1119 pub async fn unused(&self) -> Result<()> {
1121 self.state.unused().await.map_err(|_| self.abort_reason())
1122 }
1123
1124 pub async fn used(&self) -> Result<()> {
1126 self.state.used().await.map_err(|_| self.abort_reason())
1127 }
1128
1129 pub async fn closed(&self) -> Error {
1131 kio::wait(|waiter| self.poll_closed(waiter)).await
1132 }
1133
1134 pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<Error> {
1136 self.state.poll_closed(waiter).map(|()| self.abort_reason())
1137 }
1138
1139 fn abort_reason(&self) -> Error {
1141 self.state.read().abort.clone().unwrap_or(Error::Dropped)
1142 }
1143
1144 pub fn is_closed(&self) -> bool {
1146 self.state.read().is_closed()
1147 }
1148
1149 pub fn latest(&self) -> Option<u64> {
1151 self.state.read().max_sequence
1152 }
1153
1154 pub fn is_clone(&self, other: &Self) -> bool {
1156 self.state.same_channel(&other.state)
1157 }
1158
1159 pub(crate) fn weak(&self) -> TrackWeak {
1161 TrackWeak {
1162 name: self.name.clone(),
1163 state: self.state.weak(),
1164 }
1165 }
1166
1167 pub fn demand(&self) -> Demand {
1175 Demand {
1176 name: self.name.clone(),
1177 state: self.state.weak(),
1178 }
1179 }
1180
1181 pub fn consume(&self) -> Consumer {
1186 Consumer::plain(self.name.clone(), self.state.consume())
1187 }
1188
1189 pub fn subscribe(&self, subscription: impl Into<Option<Subscription>>) -> Subscriber {
1194 let preferences = subscription.into().unwrap_or_default();
1195
1196 let info = self.info.clone();
1201 let subscription = kio::Producer::new(preferences);
1202 register_subscription(self.state.read(), &subscription);
1203
1204 Subscriber {
1205 name: self.name.clone(),
1206 info,
1207 inner: SubscriberKind::Plain(PlainSubscriber {
1208 state: self.state.consume(),
1209 subscription,
1210 index: 0,
1211 datagram_index: 0,
1212 min_sequence: 0,
1213 next_sequence: 0,
1214 end_sequence: None,
1215 parked: BTreeMap::new(),
1216 }),
1217 stats: stats::Scope::default(),
1219 _stats_sub: stats::Subscription::default(),
1220 }
1221 }
1222
1223 pub async fn subscription_changed(&mut self) -> Result<Option<Subscription>> {
1229 kio::wait(|waiter| self.poll_subscription_changed(waiter)).await
1230 }
1231
1232 pub fn subscription(&self) -> Option<Subscription> {
1240 let state = self.state.read();
1241 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
1242 drop(state);
1243 snapshot_subscription(&subs, bound)
1244 }
1245
1246 pub fn poll_subscription_changed(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Subscription>>> {
1250 if self.state.poll_closed(waiter).is_ready() {
1253 let abort = self.state.read().abort.clone();
1254 return Poll::Ready(Err(abort.unwrap_or(Error::Dropped)));
1255 }
1256
1257 let state = self.state.read();
1259 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
1260 drop(state);
1261
1262 let prev = &self.prev_subscription;
1263 let mut combined = None;
1264 let mut guard = ready!(subs.poll(waiter, |subs| {
1265 let next = combined_subscription(subs, bound, waiter);
1266 if &next == prev {
1267 Poll::Pending
1268 } else {
1269 combined = next;
1270 Poll::Ready(())
1271 }
1272 }));
1273 guard.retain(|sub| !sub.is_closed());
1275 drop(guard);
1276 self.prev_subscription = combined.clone();
1277 Poll::Ready(Ok(combined))
1278 }
1279
1280 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
1282 self.state.poll_unused(waiter).map(|_| ())
1283 }
1284
1285 pub fn dynamic(&self) -> Dynamic {
1289 Dynamic::new(self.name.clone(), self.state.clone(), self.alive.clone())
1290 }
1291
1292 fn modify(&self) -> Result<kio::Mut<'_, TrackState>> {
1293 TrackState::modify(&self.state)
1294 }
1295}
1296
1297fn poll_requested_group(
1301 state: &kio::Producer<TrackState>,
1302 fetch: &kio::Shared<FetchState>,
1303 waiter: &kio::Waiter,
1304) -> Poll<Result<GroupRequest>> {
1305 if let Poll::Ready(mut guard) = fetch.poll(waiter, |fetch| {
1307 if fetch.has_queued() {
1308 Poll::Ready(())
1309 } else {
1310 Poll::Pending
1311 }
1312 }) {
1313 let sequence = guard.pop().expect("predicate guaranteed a request");
1314 let pending = guard.get(&sequence).expect("popped key must be pending");
1318 let priority = pending.priority;
1319 let result = pending.result.clone();
1320 drop(guard);
1321 return Poll::Ready(Ok(GroupRequest {
1322 state: state.clone(),
1323 fetch: fetch.clone(),
1324 sequence,
1325 priority,
1326 result,
1327 done: false,
1328 }));
1329 }
1330
1331 match state.poll_ref(waiter, |state| match &state.abort {
1333 Some(err) => Poll::Ready(err.clone()),
1334 None => Poll::Pending,
1335 }) {
1336 Poll::Ready(Ok(err)) => Poll::Ready(Err(err)),
1337 Poll::Ready(Err(closed)) => Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped))),
1338 Poll::Pending => Poll::Pending,
1339 }
1340}
1341
1342pub struct Dynamic {
1352 name: Arc<str>,
1353 state: kio::Producer<TrackState>,
1355 fetch: kio::Shared<FetchState>,
1357 alive: Arc<Alive>,
1360}
1361
1362impl Dynamic {
1363 fn new(name: Arc<str>, state: kio::Producer<TrackState>, alive: Arc<Alive>) -> Self {
1364 let fetch = state.read().fetch.clone();
1365 fetch.lock().add_handler();
1366 Self {
1367 name,
1368 state,
1369 fetch,
1370 alive,
1371 }
1372 }
1373
1374 pub fn name(&self) -> &str {
1376 &self.name
1377 }
1378
1379 pub async fn requested_group(&self) -> Result<GroupRequest> {
1385 kio::wait(|waiter| self.poll_requested_group(waiter)).await
1386 }
1387
1388 pub fn poll_requested_group(&self, waiter: &kio::Waiter) -> Poll<Result<GroupRequest>> {
1390 poll_requested_group(&self.state, &self.fetch, waiter)
1391 }
1392
1393 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
1395 self.state.poll_unused(waiter).map(|_| ())
1396 }
1397}
1398
1399impl Clone for Dynamic {
1400 fn clone(&self) -> Self {
1401 self.fetch.lock().add_handler();
1403 Self {
1404 name: self.name.clone(),
1405 state: self.state.clone(),
1406 fetch: self.fetch.clone(),
1407 alive: self.alive.clone(),
1408 }
1409 }
1410}
1411
1412impl Drop for Dynamic {
1413 fn drop(&mut self) {
1414 let mut fetch = self.fetch.lock();
1420 if fetch.remove_handler() {
1421 fetch.drain_queued();
1422 }
1423 }
1424}
1425
1426struct Alive {
1435 name: Arc<str>,
1436 state: kio::Producer<TrackState>,
1437
1438 published: AtomicBool,
1441
1442 stats: OnceLock<stats::Subscription>,
1445}
1446
1447impl Alive {
1448 fn new(name: Arc<str>, state: kio::Producer<TrackState>) -> Arc<Self> {
1449 Arc::new(Self {
1450 name,
1451 state,
1452 published: Default::default(),
1453 stats: Default::default(),
1454 })
1455 }
1456
1457 fn publish(&self, stats: Option<&stats::Scope>) {
1461 self.published.store(true, Ordering::Relaxed);
1462 if let Some(scope) = stats {
1463 let _ = self.stats.set(scope.subscribe());
1466 }
1467 }
1468}
1469
1470impl Drop for Alive {
1471 fn drop(&mut self) {
1472 if !self.published.load(Ordering::Relaxed) {
1474 return;
1475 }
1476 match self.state.write() {
1484 Ok(mut state) => {
1485 if state.final_sequence.is_some() || state.abort.is_some() {
1486 return;
1487 }
1488 tracing::warn!(
1489 track = %self.name,
1490 "track::Producer dropped without finish() or abort()"
1491 );
1492 state.clear_cache();
1493 state.datagrams.clear();
1494 }
1495 Err(state) => {
1496 if state.final_sequence.is_some() || state.abort.is_some() {
1497 return;
1498 }
1499 tracing::warn!(
1500 track = %self.name,
1501 "track::Producer dropped without finish() or abort()"
1502 );
1503 }
1504 }
1505 }
1506}
1507
1508fn combined_subscription(subs: &Subscriptions, bound: Option<Duration>, waiter: &kio::Waiter) -> Option<Subscription> {
1514 let mut combined = None;
1515 for sub in subs.iter() {
1516 if sub.is_closed() {
1521 continue;
1522 }
1523 let _ = sub.poll_closed(waiter);
1528 if let Poll::Ready(Ok(sub)) = sub.poll(waiter, |sub| sub.poll_combined(&combined)) {
1529 combined = Some(sub);
1530 }
1531 }
1532 clamp_combined(combined, bound)
1533}
1534
1535fn snapshot_subscription(subs: &kio::Shared<Subscriptions>, bound: Option<Duration>) -> Option<Subscription> {
1537 let mut combined: Option<Subscription> = None;
1538 for sub in subs.read().iter() {
1539 if sub.is_closed() {
1541 continue;
1542 }
1543 if let Poll::Ready(merged) = sub.read().poll_combined(&combined) {
1544 combined = Some(merged);
1545 }
1546 }
1547 clamp_combined(combined, bound)
1548}
1549
1550fn clamp_combined(combined: Option<Subscription>, bound: Option<Duration>) -> Option<Subscription> {
1558 let mut combined = combined?;
1559 if let Some(bound) = bound {
1560 combined.latency_max = combined.latency_max.min(bound);
1561 }
1562 Some(combined)
1563}
1564
1565fn register_subscription(state: kio::Ref<'_, TrackState>, subscription: &kio::Producer<Subscription>) {
1569 if state.is_closed() {
1570 return;
1571 }
1572 let subs = state.subscriptions.clone();
1573 drop(state);
1574 subs.lock().push(subscription.consume());
1575}
1576
1577#[derive(Clone)]
1579pub(crate) struct TrackWeak {
1580 name: Arc<str>,
1581 state: kio::ProducerWeak<TrackState>,
1582}
1583
1584impl TrackWeak {
1585 pub fn consume(&self) -> Consumer {
1586 Consumer::plain(self.name.clone(), self.state.consume())
1587 }
1588
1589 pub(crate) fn name(&self) -> &Arc<str> {
1592 &self.name
1593 }
1594
1595 pub(crate) fn reject(&self, err: Error) -> bool {
1604 let Some(producer) = self.state.produce() else {
1605 return false;
1606 };
1607 let Ok(mut state) = producer.write() else {
1608 return false;
1609 };
1610 if state.published || state.abort.is_some() {
1611 return false;
1612 }
1613 state.abort = Some(err);
1614 state.close();
1615 true
1616 }
1617
1618 pub(crate) fn is_used(&self) -> bool {
1621 !self.state.is_closed() && self.state.is_used()
1622 }
1623
1624 pub(crate) fn poll_used(&self, waiter: &kio::Waiter) {
1627 let _ = self.state.poll_used(waiter);
1628 }
1629
1630 pub(crate) fn poll_unused(&self, waiter: &kio::Waiter) {
1633 let _ = self.state.poll_unused(waiter);
1634 }
1635}
1636
1637impl super::WeakEntry for TrackWeak {
1638 fn is_closed(&self) -> bool {
1639 self.state.is_closed()
1640 }
1641
1642 fn same_channel(&self, other: &Self) -> bool {
1643 self.state.same_channel(&other.state)
1644 }
1645}
1646
1647#[derive(Clone)]
1656pub struct Demand {
1657 name: Arc<str>,
1658 state: kio::ProducerWeak<TrackState>,
1659}
1660
1661impl Demand {
1662 pub fn name(&self) -> &str {
1664 &self.name
1665 }
1666
1667 pub async fn used(&self) -> Result<()> {
1669 self.state.used().await.map_err(|_| self.abort_reason())
1670 }
1671
1672 pub async fn unused(&self) -> Result<()> {
1674 self.state.unused().await.map_err(|_| self.abort_reason())
1675 }
1676
1677 pub async fn closed(&self) -> Error {
1679 self.state.closed().await;
1680 self.abort_reason()
1681 }
1682
1683 fn abort_reason(&self) -> Error {
1685 self.state.read().abort.clone().unwrap_or(Error::Dropped)
1686 }
1687}
1688
1689#[derive(Clone)]
1700pub struct Consumer {
1701 name: Arc<str>,
1702 inner: ConsumerKind,
1703 stats: stats::Scope,
1706}
1707
1708#[derive(Clone)]
1709enum ConsumerKind {
1710 Plain(kio::Consumer<TrackState>),
1711 Spliced(super::resume::Consumer),
1712}
1713
1714impl Consumer {
1715 fn plain(name: Arc<str>, state: kio::Consumer<TrackState>) -> Self {
1716 Self {
1717 name,
1718 inner: ConsumerKind::Plain(state),
1719 stats: stats::Scope::default(),
1720 }
1721 }
1722
1723 pub(crate) fn spliced(name: Arc<str>, resume: super::resume::Consumer) -> Self {
1725 Self {
1726 name,
1727 inner: ConsumerKind::Spliced(resume),
1728 stats: stats::Scope::default(),
1729 }
1730 }
1731
1732 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
1735 self.stats = scope;
1736 self
1737 }
1738
1739 pub fn name(&self) -> &str {
1741 &self.name
1742 }
1743
1744 pub fn subscribe(&self, subscription: impl Into<Option<Subscription>>) -> kio::Pending<Subscribing> {
1750 let subscription = kio::Producer::new(subscription.into().unwrap_or_default());
1751
1752 let inner = match &self.inner {
1753 ConsumerKind::Plain(state) => {
1754 register_subscription(state.read(), &subscription);
1757 SubscribingKind::Plain(state.clone())
1758 }
1759 ConsumerKind::Spliced(resume) => SubscribingKind::Spliced(resume.clone()),
1761 };
1762
1763 kio::Pending::new(Subscribing {
1764 name: self.name.clone(),
1765 inner,
1766 subscription,
1767 stats: self.stats.clone(),
1768 })
1769 }
1770
1771 pub(crate) fn peek_latest(&self) -> Option<group::Consumer> {
1775 match &self.inner {
1776 ConsumerKind::Plain(state) => {
1777 let sequence = state.read().max_sequence?;
1778 self.peek_group(sequence)
1779 }
1780 ConsumerKind::Spliced(resume) => resume.peek_latest(),
1781 }
1782 }
1783
1784 pub(crate) fn peek_before(&self, sequence: u64) -> Option<group::Consumer> {
1788 match &self.inner {
1789 ConsumerKind::Plain(state) => {
1790 let state = state.read();
1791 state
1792 .lookup
1793 .range(..sequence)
1794 .rev()
1795 .map(|(_, slot)| &slot.group)
1796 .find(|group| !group.is_aborted())
1797 .map(|group| group.consume())
1798 }
1799 ConsumerKind::Spliced(resume) => resume.peek_before(sequence),
1800 }
1801 }
1802
1803 pub(crate) fn peek_group(&self, sequence: u64) -> Option<group::Consumer> {
1808 match &self.inner {
1809 ConsumerKind::Plain(state) => {
1810 let state = state.read();
1811 let slot = state.lookup.get(&sequence)?;
1812 if slot.group.is_aborted() {
1813 return None;
1814 }
1815 Some(slot.group.consume())
1816 }
1817 ConsumerKind::Spliced(resume) => resume.peek_group(sequence),
1818 }
1819 }
1820
1821 pub fn fetch_group(&self, sequence: u64, options: impl Into<Option<group::Fetch>>) -> kio::Pending<Fetching> {
1833 let options = options.into().unwrap_or_default();
1834
1835 self.stats.fetch();
1839
1840 let state = match &self.inner {
1841 ConsumerKind::Plain(state) => state,
1842 ConsumerKind::Spliced(resume) => {
1845 return kio::Pending::new(Fetching {
1846 inner: FetchingKind::Spliced(resume.fetch_group(sequence, options)),
1847 stats: self.stats.clone(),
1848 });
1849 }
1850 };
1851
1852 let mut result = None;
1853
1854 let (fetch, unresolved) = {
1858 let state = state.read();
1859 (state.fetch.clone(), state.poll_fetch_cached(sequence).is_pending())
1860 };
1861
1862 if unresolved {
1863 let mut fetch = fetch.lock();
1864 if let Some(pending) = fetch.join(&sequence) {
1865 pending.priority = pending.priority.max(options.priority);
1868 result = Some(pending.result.consume());
1869 } else {
1870 let producer = kio::Producer::<FetchOutcome>::default();
1874 let consumer = producer.consume();
1875 let attempt = PendingFetch {
1876 priority: options.priority,
1877 result: producer,
1878 };
1879 if fetch.insert(sequence, attempt).is_ok() {
1880 result = Some(consumer);
1881 }
1882 }
1883 }
1884
1885 kio::Pending::new(Fetching {
1886 inner: FetchingKind::Plain {
1887 state: state.clone(),
1888 fetch,
1889 sequence,
1890 result,
1891 },
1892 stats: self.stats.clone(),
1893 })
1894 }
1895
1896 pub fn info(&self) -> kio::Pending<Querying> {
1903 kio::Pending::new(Querying {
1904 inner: match &self.inner {
1905 ConsumerKind::Plain(state) => QueryingKind::Plain(state.clone()),
1906 ConsumerKind::Spliced(resume) => QueryingKind::Spliced(resume.clone()),
1907 },
1908 })
1909 }
1910
1911 pub fn latest(&self) -> Option<u64> {
1913 match &self.inner {
1914 ConsumerKind::Plain(state) => state.read().max_sequence,
1915 ConsumerKind::Spliced(resume) => resume.latest(),
1916 }
1917 }
1918
1919 pub(crate) fn poll_complete(&self, waiter: &kio::Waiter) -> Poll<Result<()>> {
1924 let ConsumerKind::Plain(state) = &self.inner else {
1925 return Poll::Pending;
1927 };
1928 match ready!(state.poll(waiter, |state| {
1929 if state.is_complete() {
1930 Poll::Ready(())
1931 } else {
1932 Poll::Pending
1933 }
1934 })) {
1935 Ok(_) => Poll::Ready(Ok(())),
1936 Err(closed) => Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped))),
1939 }
1940 }
1941}
1942
1943pub struct Subscribing {
1946 name: Arc<str>,
1947 inner: SubscribingKind,
1948 subscription: kio::Producer<Subscription>,
1949 stats: stats::Scope,
1950}
1951
1952enum SubscribingKind {
1953 Plain(kio::Consumer<TrackState>),
1954 Spliced(super::resume::Consumer),
1955}
1956
1957impl Subscribing {
1958 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<Subscriber>> {
1961 match &self.inner {
1962 SubscribingKind::Plain(state) => {
1963 let info = ready!(state.poll(waiter, |state| state.poll_info()))
1965 .map_err(|e| e.abort.clone().unwrap_or(Error::Dropped))??;
1966
1967 Poll::Ready(Ok(Subscriber {
1968 name: self.name.clone(),
1969 info,
1970 inner: SubscriberKind::Plain(PlainSubscriber {
1971 state: state.clone(),
1972 subscription: self.subscription.clone(),
1973 index: 0,
1974 datagram_index: 0,
1975 min_sequence: 0,
1976 next_sequence: 0,
1977 end_sequence: None,
1978 parked: BTreeMap::new(),
1979 }),
1980 stats: self.stats.clone(),
1981 _stats_sub: self.stats.subscribe(),
1982 }))
1983 }
1984 SubscribingKind::Spliced(resume) => {
1985 let info = ready!(resume.poll_info(waiter))?;
1988
1989 Poll::Ready(Ok(Subscriber {
1990 name: self.name.clone(),
1991 info,
1992 inner: SubscriberKind::Spliced(Box::new(resume.subscribe_shared(self.subscription.clone()))),
1993 stats: self.stats.clone(),
1994 _stats_sub: self.stats.subscribe(),
1995 }))
1996 }
1997 }
1998 }
1999
2000 pub fn update(&mut self, subscription: Subscription) -> Result<()> {
2005 let mut state = self.subscription.write().map_err(|_| Error::Closed)?;
2006 *state = subscription;
2007 Ok(())
2008 }
2009}
2010
2011impl kio::Pollable for Subscribing {
2012 type Output = Result<Subscriber>;
2013
2014 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
2015 self.poll_ok(waiter)
2016 }
2017}
2018
2019pub struct Querying {
2022 inner: QueryingKind,
2023}
2024
2025enum QueryingKind {
2026 Plain(kio::Consumer<TrackState>),
2027 Spliced(super::resume::Consumer),
2028}
2029
2030impl Querying {
2031 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<Info>> {
2033 match &self.inner {
2034 QueryingKind::Plain(state) => {
2035 let info = ready!(state.poll(waiter, |state| state.poll_info()))
2037 .map_err(|e| e.abort.clone().unwrap_or(Error::Dropped))??;
2038 Poll::Ready(Ok(info))
2039 }
2040 QueryingKind::Spliced(resume) => resume.poll_info(waiter),
2041 }
2042 }
2043}
2044
2045impl kio::Pollable for Querying {
2046 type Output = Result<Info>;
2047
2048 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
2049 self.poll_ok(waiter)
2050 }
2051}
2052
2053pub struct GroupRequest {
2062 state: kio::Producer<TrackState>,
2063 fetch: kio::Shared<FetchState>,
2065 sequence: u64,
2066 priority: u8,
2067 result: kio::Producer<FetchOutcome>,
2069 done: bool,
2070}
2071
2072impl GroupRequest {
2073 pub fn sequence(&self) -> u64 {
2075 self.sequence
2076 }
2077
2078 pub fn priority(&self) -> u8 {
2080 self.priority
2081 }
2082
2083 pub fn accept(mut self, info: impl Into<Option<Info>>) -> Result<group::Producer> {
2091 self.done = true;
2092 let res = TrackState::modify(&self.state)
2096 .and_then(|mut state| state.insert_group_request(self.sequence, info.into()));
2097 self.remove();
2098 res
2099 }
2100
2101 pub fn reject(mut self, err: Error) {
2103 self.done = true;
2104 self.remove();
2107 if let Ok(mut outcome) = self.result.write() {
2108 outcome.rejected = Some(err);
2109 }
2110 }
2111
2112 fn remove(&self) {
2115 self.fetch
2116 .lock()
2117 .remove_if(&self.sequence, |pending| pending.result.same_channel(&self.result));
2118 }
2119}
2120
2121impl Drop for GroupRequest {
2122 fn drop(&mut self) {
2123 if self.done {
2124 return;
2125 }
2126 self.remove();
2127 if let Ok(mut outcome) = self.result.write() {
2128 outcome.rejected = Some(Error::Dropped);
2129 }
2130 }
2131}
2132
2133pub struct Fetching {
2139 inner: FetchingKind,
2140 stats: stats::Scope,
2143}
2144
2145enum FetchingKind {
2146 Plain {
2147 state: kio::Consumer<TrackState>,
2148 fetch: kio::Shared<FetchState>,
2149 sequence: u64,
2150 result: Option<kio::Consumer<FetchOutcome>>,
2152 },
2153 Spliced(kio::Pending<super::resume::Fetching>),
2155}
2156
2157impl kio::Pollable for Fetching {
2158 type Output = Result<group::Consumer>;
2159
2160 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
2161 let (state, fetch, sequence, result) = match &self.inner {
2162 FetchingKind::Plain {
2163 state,
2164 fetch,
2165 sequence,
2166 result,
2167 } => (state, fetch, *sequence, result.as_ref()),
2168 FetchingKind::Spliced(spliced) => {
2169 return kio::Pollable::poll(&**spliced, waiter)
2172 .map(|res| res.map(|group| group.with_meter(self.stats.meter())));
2173 }
2174 };
2175
2176 match state.poll(waiter, |state| state.poll_fetch_cached(sequence)) {
2179 Poll::Ready(Ok(res)) => return Poll::Ready(res.map(|group| group.with_meter(self.stats.meter()))),
2180 Poll::Ready(Err(closed)) => {
2181 return Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped)));
2182 }
2183 Poll::Pending => {}
2184 }
2185
2186 let Some(result) = result else {
2188 return match fetch.poll(waiter, |fetch| match fetch.has_handlers() {
2191 false => Poll::Ready(()),
2192 true => Poll::Pending,
2193 }) {
2194 Poll::Ready(_guard) => Poll::Ready(Err(Error::NotFound)),
2195 Poll::Pending => Poll::Pending,
2196 };
2197 };
2198
2199 match result.poll(waiter, |outcome| match &outcome.rejected {
2202 Some(err) => Poll::Ready(err.clone()),
2203 None => Poll::Pending,
2204 }) {
2205 Poll::Ready(Ok(err)) => Poll::Ready(Err(err)),
2206 Poll::Ready(Err(_closed)) => Poll::Ready(Err(Error::NotFound)),
2207 Poll::Pending => Poll::Pending,
2208 }
2209 }
2210}
2211
2212pub struct Subscriber {
2235 name: Arc<str>,
2236 info: Info,
2237 inner: SubscriberKind,
2238 stats: stats::Scope,
2241 _stats_sub: stats::Subscription,
2244}
2245
2246enum SubscriberKind {
2247 Plain(PlainSubscriber),
2248 Spliced(Box<super::resume::Subscriber>),
2250}
2251
2252struct PlainSubscriber {
2254 state: kio::Consumer<TrackState>,
2255
2256 subscription: kio::Producer<Subscription>,
2257 index: usize,
2259 datagram_index: usize,
2261 min_sequence: u64,
2263 next_sequence: u64,
2266 end_sequence: Option<u64>,
2271 parked: BTreeMap<u64, group::Consumer>,
2276}
2277
2278impl PlainSubscriber {
2279 fn poll<F, R>(&self, waiter: &kio::Waiter, f: F) -> Poll<Result<R>>
2281 where
2282 F: Fn(&kio::Ref<'_, TrackState>) -> Poll<Result<R>>,
2283 {
2284 Poll::Ready(match ready!(self.state.poll(waiter, f)) {
2285 Ok(res) => res,
2286 Err(state) => Err(state.abort.clone().unwrap_or(Error::Dropped)),
2288 })
2289 }
2290
2291 fn poll_recv_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2292 let watch = |group: &group::Consumer| match group.poll_closed(waiter) {
2300 Poll::Pending => true,
2301 Poll::Ready(()) => !group.is_aborted(),
2302 };
2303
2304 let min_sequence = self.min_sequence;
2309 self.parked
2310 .retain(|sequence, group| *sequence >= min_sequence && watch(group));
2311
2312 if let Some(&sequence) = self.parked.keys().next()
2314 && self.end_sequence.is_none_or(|end| sequence <= end)
2315 {
2316 let group = self.parked.remove(&sequence).expect("parked key just observed");
2317 group.cache_refresh();
2319 return Poll::Ready(Ok(Some(group)));
2320 }
2321
2322 loop {
2323 let Some((consumer, found_index)) =
2324 ready!(self.poll(waiter, |state| state.poll_recv_group(self.index, self.min_sequence))?)
2325 else {
2326 if self.parked.is_empty() {
2329 return Poll::Ready(Ok(None));
2330 }
2331 return Poll::Pending;
2332 };
2333 self.index = found_index + 1;
2334
2335 if self.end_sequence.is_some_and(|end| consumer.sequence > end) {
2338 if watch(&consumer) {
2342 self.parked.insert(consumer.sequence, consumer);
2343 }
2344 continue;
2345 }
2346 return Poll::Ready(Ok(Some(consumer)));
2347 }
2348 }
2349
2350 fn poll_recv_datagram(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Datagram>>> {
2351 let Some((datagram, found_index)) =
2352 ready!(self.poll(waiter, |state| state.poll_recv_datagram(self.datagram_index))?)
2353 else {
2354 return Poll::Ready(Ok(None));
2355 };
2356
2357 self.datagram_index = found_index + 1;
2358 self.next_sequence = self.next_sequence.max(datagram.sequence.saturating_add(1));
2359 Poll::Ready(Ok(Some(datagram)))
2360 }
2361
2362 fn poll_next_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2363 let floor = self.next_sequence.max(self.min_sequence);
2364 let Some(group) = ready!(self.poll(waiter, |state| state.poll_next_in_range(floor, self.end_sequence))?) else {
2365 return Poll::Ready(Ok(None));
2366 };
2367 self.next_sequence = group.sequence.saturating_add(1);
2368 Poll::Ready(Ok(Some(group)))
2369 }
2370
2371 fn poll_read_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Frame>>> {
2372 let lower = self.min_sequence.max(self.next_sequence);
2373 let Some((frame, found_index, sequence)) =
2374 ready!(self.poll(waiter, |state| { state.poll_read_frame(self.index, lower, waiter) })?)
2375 else {
2376 return Poll::Ready(Ok(None));
2377 };
2378
2379 self.index = found_index + 1;
2380 self.next_sequence = sequence.saturating_add(1);
2381 Poll::Ready(Ok(Some(frame)))
2382 }
2383}
2384
2385#[derive(Clone)]
2391pub struct SubscriberControl {
2392 subscription: kio::Producer<Subscription>,
2393}
2394
2395impl SubscriberControl {
2396 pub fn subscription(&self) -> Subscription {
2398 self.subscription.read().clone()
2399 }
2400
2401 pub fn update(&self, subscription: Subscription) -> Result<()> {
2406 let mut state = self.subscription.write().map_err(|_| Error::Closed)?;
2407 *state = subscription;
2408 Ok(())
2409 }
2410}
2411
2412impl Subscriber {
2413 pub fn info(&self) -> &Info {
2418 &self.info
2419 }
2420
2421 pub fn name(&self) -> &str {
2423 &self.name
2424 }
2425
2426 pub fn control(&self) -> SubscriberControl {
2428 SubscriberControl {
2429 subscription: match &self.inner {
2430 SubscriberKind::Plain(plain) => plain.subscription.clone(),
2431 SubscriberKind::Spliced(spliced) => spliced.prefs(),
2432 },
2433 }
2434 }
2435
2436 pub fn poll_recv_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2453 let meter = self.stats.meter();
2454 let res = match &mut self.inner {
2455 SubscriberKind::Plain(plain) => plain.poll_recv_group(waiter),
2456 SubscriberKind::Spliced(spliced) => spliced.poll_recv_group(waiter),
2457 };
2458 res.map(|res| res.map(|group| group.map(|group| group.with_meter(meter))))
2459 }
2460
2461 pub async fn recv_group(&mut self) -> Result<Option<group::Consumer>> {
2468 kio::wait(|waiter| self.poll_recv_group(waiter)).await
2469 }
2470
2471 pub fn poll_recv_datagram(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Datagram>>> {
2482 let meter = self.stats.meter();
2483 let res = match &mut self.inner {
2484 SubscriberKind::Plain(plain) => plain.poll_recv_datagram(waiter),
2485 SubscriberKind::Spliced(spliced) => spliced.poll_recv_datagram(waiter),
2486 };
2487 if let Poll::Ready(Ok(Some(datagram))) = &res {
2490 meter.datagram(datagram.payload.len() as u64);
2491 }
2492 res
2493 }
2494
2495 pub async fn recv_datagram(&mut self) -> Result<Option<Datagram>> {
2502 kio::wait(|waiter| self.poll_recv_datagram(waiter)).await
2503 }
2504
2505 pub fn poll_next_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2514 let meter = self.stats.meter();
2515 let res = match &mut self.inner {
2516 SubscriberKind::Plain(plain) => plain.poll_next_group(waiter),
2517 SubscriberKind::Spliced(spliced) => spliced.poll_next_group(waiter),
2518 };
2519 res.map(|res| res.map(|group| group.map(|group| group.with_meter(meter))))
2520 }
2521
2522 pub async fn next_group(&mut self) -> Result<Option<group::Consumer>> {
2528 kio::wait(|waiter| self.poll_next_group(waiter)).await
2529 }
2530
2531 pub fn poll_read_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Frame>>> {
2535 let meter = self.stats.meter();
2536 let res = match &mut self.inner {
2537 SubscriberKind::Plain(plain) => plain.poll_read_frame(waiter),
2538 SubscriberKind::Spliced(spliced) => spliced.poll_read_frame(waiter),
2539 };
2540 if let Poll::Ready(Ok(Some(frame))) = &res {
2543 meter.group();
2544 meter.frames(1);
2545 meter.bytes(frame.payload.len() as u64);
2546 }
2547 res
2548 }
2549
2550 pub async fn read_frame(&mut self) -> Result<Option<frame::Frame>> {
2555 kio::wait(|waiter| self.poll_read_frame(waiter)).await
2556 }
2557
2558 pub fn is_clone(&self, other: &Self) -> bool {
2560 match (&self.inner, &other.inner) {
2561 (SubscriberKind::Plain(a), SubscriberKind::Plain(b)) => a.state.same_channel(&b.state),
2562 (SubscriberKind::Spliced(a), SubscriberKind::Spliced(b)) => a.is_clone(b),
2563 _ => false,
2564 }
2565 }
2566
2567 pub fn poll_finished(&mut self, waiter: &kio::Waiter) -> Poll<Result<u64>> {
2569 match &mut self.inner {
2570 SubscriberKind::Plain(plain) => plain.poll(waiter, |state| state.poll_finished()),
2571 SubscriberKind::Spliced(spliced) => spliced.poll_finished(waiter),
2572 }
2573 }
2574
2575 pub async fn finished(&mut self) -> Result<u64> {
2583 kio::wait(|waiter| self.poll_finished(waiter)).await
2584 }
2585
2586 pub fn start_at(&mut self, sequence: u64) {
2593 match &mut self.inner {
2594 SubscriberKind::Plain(plain) => plain.min_sequence = sequence,
2595 SubscriberKind::Spliced(spliced) => spliced.start_at(sequence),
2596 }
2597 }
2598
2599 pub fn end_at(&mut self, sequence: impl Into<Option<u64>>) {
2612 match &mut self.inner {
2613 SubscriberKind::Plain(plain) => plain.end_sequence = sequence.into(),
2614 SubscriberKind::Spliced(spliced) => spliced.end_at(sequence),
2615 }
2616 }
2617
2618 pub fn subscription(&self) -> Subscription {
2620 self.control().subscription()
2621 }
2622
2623 pub fn update(&mut self, subscription: Subscription) -> Result<()> {
2629 match &mut self.inner {
2630 SubscriberKind::Plain(plain) => {
2631 let mut state = plain.subscription.write().map_err(|_| Error::Closed)?;
2632 *state = subscription;
2633 }
2634 SubscriberKind::Spliced(spliced) => spliced.update(subscription),
2635 }
2636 Ok(())
2637 }
2638
2639 pub fn latest(&self) -> Option<u64> {
2641 match &self.inner {
2642 SubscriberKind::Plain(plain) => plain.state.read().max_sequence,
2643 SubscriberKind::Spliced(spliced) => spliced.latest(),
2644 }
2645 }
2646}
2647
2648pub struct Request {
2660 name: Arc<str>,
2661 broadcast: Arc<broadcast::Info>,
2663 state: kio::Producer<TrackState>,
2664
2665 prev_subscription: Option<Subscription>,
2667
2668 alive: Arc<Alive>,
2671
2672 _dynamic: Dynamic,
2677
2678 stats: stats::Scope,
2681}
2682
2683impl Request {
2684 pub(crate) fn new(broadcast: Arc<broadcast::Info>, name: impl Into<Arc<str>>) -> Self {
2685 let name = name.into();
2686 let state = TrackState::spawn(broadcast.clone());
2687 let alive = Alive::new(name.clone(), state.clone());
2688 let dynamic = Dynamic::new(name.clone(), state.clone(), alive.clone());
2689 Self {
2690 name,
2691 broadcast,
2692 state,
2693 prev_subscription: None,
2694 alive,
2695 _dynamic: dynamic,
2696 stats: stats::Scope::default(),
2697 }
2698 }
2699
2700 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
2703 self.stats = scope;
2704 self
2705 }
2706
2707 pub fn name(&self) -> &str {
2709 &self.name
2710 }
2711
2712 pub fn consume(&self) -> Consumer {
2714 Consumer::plain(self.name.clone(), self.state.consume())
2715 }
2716
2717 pub fn dynamic(&self) -> Dynamic {
2721 Dynamic::new(self.name.clone(), self.state.clone(), self.alive.clone())
2722 }
2723
2724 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
2727 self.state.poll_unused(waiter).map(|_| ())
2728 }
2729
2730 pub fn accept(self, info: impl Into<Option<Info>>) -> Producer {
2737 let info = TrackState::normalize_info(&self.broadcast, info.into().unwrap_or_default());
2738 if let Ok(mut state) = self.state.write() {
2741 state.accept(info.clone());
2742 }
2743 self.alive.publish(Some(&self.stats));
2746 Producer {
2747 name: self.name,
2748 info,
2749 broadcast: self.broadcast,
2750 state: self.state,
2751 prev_subscription: None,
2752 alive: self.alive,
2753 stats: self.stats,
2754 }
2755 }
2756
2757 pub fn reject(self, err: Error) {
2759 if let Ok(mut state) = self.state.write() {
2760 state.abort = Some(err);
2761 }
2762 }
2763
2764 pub fn subscription(&self) -> Option<Subscription> {
2767 let state = self.state.read();
2768 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
2769 drop(state);
2770 snapshot_subscription(&subs, bound)
2771 }
2772
2773 pub async fn subscription_changed(&mut self) -> Option<Subscription> {
2776 kio::wait(|waiter| self.poll_subscription_changed(waiter)).await
2777 }
2778
2779 pub fn poll_subscription_changed(&mut self, waiter: &kio::Waiter) -> Poll<Option<Subscription>> {
2781 let state = self.state.read();
2782 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
2783 drop(state);
2784
2785 let prev = &self.prev_subscription;
2786 let mut combined = None;
2787 let mut guard = ready!(subs.poll(waiter, |subs| {
2788 let next = combined_subscription(subs, bound, waiter);
2789 if &next == prev {
2790 Poll::Pending
2791 } else {
2792 combined = next;
2793 Poll::Ready(())
2794 }
2795 }));
2796 guard.retain(|sub| !sub.is_closed());
2798 drop(guard);
2799 self.prev_subscription = combined.clone();
2800 Poll::Ready(combined)
2801 }
2802
2803 pub(super) fn weak(&self) -> TrackWeak {
2804 TrackWeak {
2805 name: self.name.clone(),
2806 state: self.state.weak(),
2807 }
2808 }
2809}
2810
2811#[cfg(test)]
2812use futures::FutureExt;
2813
2814#[cfg(test)]
2815#[allow(missing_docs)] impl Subscriber {
2817 pub fn assert_group(&mut self) -> group::Consumer {
2818 self.recv_group()
2819 .now_or_never()
2820 .expect("group would have blocked")
2821 .expect("would have errored")
2822 .expect("track was closed")
2823 }
2824
2825 pub fn assert_no_group(&mut self) {
2826 assert!(
2827 self.recv_group().now_or_never().is_none(),
2828 "recv_group would not have blocked"
2829 );
2830 }
2831
2832 pub fn assert_not_closed(&mut self) {
2833 assert!(self.finished().now_or_never().is_none(), "should not be closed");
2834 }
2835
2836 pub fn assert_closed(&mut self) {
2837 assert!(self.finished().now_or_never().is_some(), "should be closed");
2838 }
2839
2840 pub fn assert_error(&mut self) {
2842 assert!(
2843 self.finished().now_or_never().expect("should not block").is_err(),
2844 "should be error"
2845 );
2846 }
2847
2848 pub fn assert_is_clone(&self, other: &Self) {
2849 assert!(self.is_clone(other), "should be clone");
2850 }
2851
2852 pub fn assert_not_clone(&self, other: &Self) {
2853 assert!(!self.is_clone(other), "should not be clone");
2854 }
2855}
2856
2857#[cfg(test)]
2858mod test {
2859 use super::*;
2860 use crate::model::test_tracing::count_drop_warnings;
2861
2862 fn track_producer(name: impl Into<Arc<str>>, info: impl Into<Option<Info>>) -> Producer {
2865 Producer::new(Arc::new(broadcast::Info::default()), name, info)
2866 }
2867
2868 fn live_groups(state: &TrackState) -> usize {
2870 state.lookup.len()
2871 }
2872
2873 fn first_live_sequence(state: &TrackState) -> u64 {
2875 state
2876 .arrival
2877 .iter()
2878 .find(|(sequence, stamp)| state.lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp))
2879 .map(|(sequence, _)| *sequence)
2880 .unwrap()
2881 }
2882
2883 fn recv_datagram(dg: &mut Subscriber) -> Datagram {
2885 dg.recv_datagram()
2886 .now_or_never()
2887 .expect("datagram would have blocked")
2888 .expect("would have errored")
2889 .expect("track was closed")
2890 }
2891
2892 #[tokio::test]
2893 async fn append_datagram_shares_group_sequence() {
2894 let mut producer = track_producer("test", None);
2895 let ts = Timestamp::from_millis(10).unwrap();
2896
2897 assert_eq!(producer.append_group().unwrap().sequence, 0);
2899 assert_eq!(producer.append_datagram(ts, &b"a"[..]).unwrap(), 1);
2900 assert_eq!(producer.append_group().unwrap().sequence, 2);
2901 assert_eq!(producer.append_datagram(ts, &b"b"[..]).unwrap(), 3);
2902 assert_eq!(producer.latest(), Some(3));
2903 }
2904
2905 #[tokio::test]
2906 async fn append_datagram_roundtrip() {
2907 let mut producer = track_producer("test", None);
2908 let mut dg = producer.subscribe(None);
2909
2910 let ts = Timestamp::from_millis(42).unwrap();
2911 let seq = producer.append_datagram(ts, &b"hello"[..]).unwrap();
2912
2913 let got = recv_datagram(&mut dg);
2914 assert_eq!(got.sequence, seq);
2915 assert_eq!(got.timestamp, ts);
2916 assert_eq!(&got.payload[..], b"hello");
2917 }
2918
2919 #[tokio::test]
2920 async fn write_datagram_preserves_sequence() {
2921 let mut producer = track_producer("test", None);
2922 let mut dg = producer.subscribe(None);
2923
2924 let ts = Timestamp::from_millis(5).unwrap();
2925 producer
2927 .write_datagram(Datagram {
2928 sequence: 100,
2929 timestamp: ts,
2930 payload: bytes::Bytes::from_static(b"x"),
2931 })
2932 .unwrap();
2933
2934 assert_eq!(recv_datagram(&mut dg).sequence, 100);
2935 assert_eq!(producer.append_group().unwrap().sequence, 101);
2937 }
2938
2939 #[tokio::test]
2940 async fn recv_datagram_advances_ordered_group_cursor() {
2941 let mut producer = track_producer("test", None);
2942 let mut subscriber = producer.subscribe(None);
2943 let ts = Timestamp::from_millis(5).unwrap();
2944
2945 producer
2946 .write_datagram(Datagram {
2947 sequence: 5,
2948 timestamp: ts,
2949 payload: bytes::Bytes::from_static(b"x"),
2950 })
2951 .unwrap();
2952 assert_eq!(recv_datagram(&mut subscriber).sequence, 5);
2953
2954 producer.create_group(group::Info { sequence: 3 }).unwrap();
2955 producer.create_group(group::Info { sequence: 6 }).unwrap();
2956
2957 let group = subscriber
2958 .next_group()
2959 .now_or_never()
2960 .expect("group would have blocked")
2961 .expect("would have errored")
2962 .expect("track was closed");
2963 assert_eq!(group.sequence, 6);
2964 }
2965
2966 #[tokio::test]
2967 async fn datagram_normalized_to_track_timescale() {
2968 let info = Info::default().with_timescale(Timescale::MICRO);
2969 let mut producer = track_producer("test", info);
2970 let mut dg = producer.subscribe(None);
2971
2972 producer
2974 .append_datagram(Timestamp::from_millis(2).unwrap(), &b"z"[..])
2975 .unwrap();
2976 let got = recv_datagram(&mut dg);
2977 assert_eq!(got.timestamp.scale(), Timescale::MICRO);
2978 assert_eq!(got.timestamp.value(), 2_000);
2979 }
2980
2981 #[tokio::test]
2982 async fn datagram_rejects_oversized() {
2983 let mut producer = track_producer("test", None);
2984 let big = bytes::Bytes::from(vec![0u8; crate::model::datagram::MAX_DATAGRAM_PAYLOAD + 1]);
2985 let ts = Timestamp::from_millis(0).unwrap();
2986 assert!(matches!(
2987 producer.append_datagram(ts, big.clone()),
2988 Err(Error::FrameTooLarge)
2989 ));
2990 assert!(matches!(
2991 producer.write_datagram(Datagram {
2992 sequence: 0,
2993 timestamp: ts,
2994 payload: big,
2995 }),
2996 Err(Error::FrameTooLarge)
2997 ));
2998 }
2999
3000 #[tokio::test]
3001 async fn datagram_fanout_to_subscribers() {
3002 let mut producer = track_producer("test", None);
3003 let mut a = producer.subscribe(None);
3005 let mut b = producer.subscribe(None);
3006 let ts = Timestamp::from_millis(1).unwrap();
3007
3008 producer.append_datagram(ts, &b"first"[..]).unwrap();
3009 producer.append_datagram(ts, &b"second"[..]).unwrap();
3010
3011 assert_eq!(&recv_datagram(&mut a).payload[..], b"first");
3013 assert_eq!(&recv_datagram(&mut a).payload[..], b"second");
3014 assert_eq!(&recv_datagram(&mut b).payload[..], b"first");
3015 assert_eq!(&recv_datagram(&mut b).payload[..], b"second");
3016 }
3017
3018 #[tokio::test]
3019 async fn datagram_evicts_stale() {
3020 tokio::time::pause();
3021
3022 let mut producer = track_producer("test", None);
3023 let mut dg = producer.subscribe(None);
3024 let ts = Timestamp::from_millis(0).unwrap();
3025
3026 producer.append_datagram(ts, &b"old"[..]).unwrap(); tokio::time::advance(MAX_DATAGRAM_AGE + Duration::from_millis(10)).await;
3030 producer.append_datagram(ts, &b"new"[..]).unwrap(); let got = recv_datagram(&mut dg);
3034 assert_eq!(got.sequence, 1);
3035 assert_eq!(&got.payload[..], b"new");
3036 }
3037
3038 #[tokio::test]
3039 async fn datagram_recv_pends_until_written() {
3040 let mut producer = track_producer("test", None);
3041 let mut dg = producer.subscribe(None);
3042
3043 assert!(
3044 dg.recv_datagram().now_or_never().is_none(),
3045 "should block with no datagrams"
3046 );
3047
3048 producer
3049 .append_datagram(Timestamp::from_millis(0).unwrap(), &b"go"[..])
3050 .unwrap();
3051 assert_eq!(&recv_datagram(&mut dg).payload[..], b"go");
3052 }
3053
3054 #[tokio::test]
3058 async fn datagram_wire_roundtrip_between_tracks() {
3059 use crate::coding::{Decode, Encode};
3060 use crate::lite;
3061
3062 let version = lite::Version::Lite05;
3063
3064 let mut origin = track_producer("test", None);
3066 let mut origin_dg = origin.subscribe(None);
3067 let ts = Timestamp::from_millis(7).unwrap();
3068 let seq = origin.append_datagram(ts, &b"payload"[..]).unwrap();
3069
3070 let d = recv_datagram(&mut origin_dg);
3071 let body = lite::Datagram {
3072 subscribe: 5,
3073 sequence: d.sequence,
3074 timestamp: d.timestamp.value(),
3075 payload: d.payload.clone(),
3076 }
3077 .encode_bytes(version)
3078 .unwrap();
3079
3080 let mut slice = &body[..];
3082 let wire = lite::Datagram::decode(&mut slice, version).unwrap();
3083 let mut downstream = track_producer("test", None);
3084 let mut downstream_dg = downstream.subscribe(None);
3085 downstream
3086 .write_datagram(Datagram {
3087 sequence: wire.sequence,
3088 timestamp: Timestamp::new(wire.timestamp, Timescale::MILLI).unwrap(),
3089 payload: wire.payload,
3090 })
3091 .unwrap();
3092
3093 let got = recv_datagram(&mut downstream_dg);
3094 assert_eq!(got.sequence, seq);
3095 assert_eq!(got.timestamp, ts);
3096 assert_eq!(&got.payload[..], b"payload");
3097 }
3098
3099 #[tokio::test]
3100 async fn evict_expired_groups() {
3101 tokio::time::pause();
3102
3103 let mut producer = track_producer("test", None);
3104
3105 producer.append_group().unwrap(); producer.append_group().unwrap(); producer.append_group().unwrap(); {
3111 let state = producer.state.read();
3112 assert_eq!(live_groups(&state), 3);
3113 assert_eq!(state.offset, 0);
3114 }
3115
3116 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3118
3119 producer.append_group().unwrap(); {
3126 let state = producer.state.read();
3127 assert_eq!(live_groups(&state), 1);
3128 assert_eq!(first_live_sequence(&state), 3);
3129 assert_eq!(state.offset, 3);
3130 assert!(!state.lookup.contains_key(&0));
3131 assert!(!state.lookup.contains_key(&1));
3132 assert!(!state.lookup.contains_key(&2));
3133 assert!(state.lookup.contains_key(&3));
3134 }
3135 }
3136
3137 #[tokio::test]
3141 async fn aging_out_a_finished_group_keeps_the_clean_end() {
3142 tokio::time::pause();
3143
3144 let mut producer = track_producer("test", None);
3145 let mut group = producer.create_group(group::Info { sequence: 0 }).unwrap();
3146 let mut consumer = group.consume();
3147
3148 group
3149 .write_frame(Timestamp::from_millis(0).unwrap(), b"hello".as_slice())
3150 .unwrap();
3151 assert_eq!(consumer.next_frame().await.unwrap().unwrap().size, 5);
3152
3153 tokio::time::advance(DEFAULT_LATENCY_MAX * 12).await;
3155 group.finish().unwrap();
3156 let _next = producer.create_group(group::Info { sequence: 1 }).unwrap();
3157
3158 assert!(consumer.next_frame().await.unwrap().is_none());
3159 }
3160
3161 #[tokio::test]
3165 async fn active_reader_survives_expiry() {
3166 tokio::time::pause();
3167
3168 let mut producer = track_producer("test", None);
3169 let mut subscriber = producer.subscribe(None);
3170
3171 let mut group = producer.create_group(0u64.into()).unwrap();
3173 for _ in 0..10 {
3174 group.write_frame(Timestamp::ZERO, b"x".as_slice()).unwrap();
3175 }
3176 group.finish().unwrap();
3177 let mut reading = subscriber.assert_group();
3178
3179 producer.create_group(1u64.into()).unwrap().finish().unwrap();
3181
3182 for seq in 2..12u64 {
3185 tokio::time::advance(DEFAULT_LATENCY_MAX / 2).await;
3186 let frame = reading.next_frame().await;
3187 assert!(
3188 matches!(frame, Ok(Some(_))),
3189 "an actively-read group must not expire mid-read (step {seq})"
3190 );
3191 producer.create_group(seq.into()).unwrap().finish().unwrap();
3192 }
3193
3194 let state = producer.state.read();
3195 assert!(state.lookup.contains_key(&0), "the read group survived");
3196 assert!(!state.lookup.contains_key(&1), "the unread group still expired");
3197 }
3198
3199 #[tokio::test]
3203 async fn slow_frame_reader_survives_expiry() {
3204 tokio::time::pause();
3205
3206 let mut producer = track_producer("test", None);
3207 let mut subscriber = producer.subscribe(None);
3208
3209 let mut group = producer.create_group(0u64.into()).unwrap();
3210 for _ in 0..20 {
3211 group.write_frame(Timestamp::ZERO, b"x".as_slice()).unwrap();
3212 }
3213 group.finish().unwrap();
3214 let mut reading = subscriber.assert_group();
3215
3216 for seq in 1..20u64 {
3218 tokio::time::advance(DEFAULT_LATENCY_MAX / 2).await;
3219 let frame = reading.read_frame().await;
3220 assert!(
3221 matches!(frame, Ok(Some(_))),
3222 "a slow reader must not expire mid-read (step {seq})"
3223 );
3224 producer.create_group(seq.into()).unwrap().finish().unwrap();
3225 }
3226 }
3227
3228 #[tokio::test]
3233 async fn slow_batch_reader_survives_expiry_with_keep_alive() {
3234 tokio::time::pause();
3235
3236 let mut producer = track_producer("test", None);
3237 let mut subscriber = producer.subscribe(None);
3238
3239 let mut group = producer.create_group(0u64.into()).unwrap();
3240 for _ in 0..20 {
3241 group.write_frame(Timestamp::ZERO, b"x".as_slice()).unwrap();
3242 }
3243 group.finish().unwrap();
3244 let mut reading = subscriber.assert_group();
3245
3246 let mut buf = crate::frame::Buffer::<8>::new();
3249 let count = reading.read_frames(&mut buf).await.unwrap().len();
3250 assert_eq!(count, 8, "the batch is bounded by the buffer");
3251
3252 for step in 0..8u64 {
3253 tokio::time::advance(DEFAULT_LATENCY_MAX / 2).await;
3254 reading.keep_alive();
3255 producer.create_group((step + 1).into()).unwrap().finish().unwrap();
3257 }
3258
3259 let rest = reading
3261 .read_frames(&mut buf)
3262 .await
3263 .expect("a batch reader that kept the group alive must not be expired");
3264 assert_eq!(rest.len(), 8, "the next batch picks up where the last one stopped");
3265 }
3266
3267 #[tokio::test]
3271 async fn delivery_restarts_the_expiry_clock() {
3272 tokio::time::pause();
3273
3274 let mut producer = track_producer("test", None);
3275 let mut subscriber = producer.subscribe(None);
3276
3277 let mut group = producer.create_group(0u64.into()).unwrap();
3278 group.write_frame(Timestamp::ZERO, b"x".as_slice()).unwrap();
3279 group.finish().unwrap();
3280 producer.create_group(1u64.into()).unwrap().finish().unwrap();
3282
3283 tokio::time::advance(DEFAULT_LATENCY_MAX - Duration::from_secs(1)).await;
3285 let mut reading = subscriber.assert_group();
3286
3287 tokio::time::advance(DEFAULT_LATENCY_MAX - Duration::from_secs(1)).await;
3290 producer.create_group(2u64.into()).unwrap().finish().unwrap();
3291
3292 let frame = reading.read_frame().await.unwrap();
3293 assert!(frame.is_some(), "a just-delivered group must not expire unread");
3294 }
3295
3296 #[tokio::test]
3300 async fn streaming_frame_writes_keep_the_group_alive() {
3301 tokio::time::pause();
3302
3303 let mut producer = track_producer("test", None);
3304 let mut straggler = producer.create_group(0u64.into()).unwrap();
3305 producer.create_group(1u64.into()).unwrap().finish().unwrap();
3307
3308 let mut frame = straggler
3309 .create_frame(frame::Info {
3310 size: 10,
3311 timestamp: Timestamp::ZERO,
3312 })
3313 .unwrap();
3314 for seq in 2..12u64 {
3317 tokio::time::advance(DEFAULT_LATENCY_MAX / 2).await;
3318 frame.write(bytes::Bytes::from_static(b"x")).unwrap();
3319 producer.create_group(seq.into()).unwrap().finish().unwrap();
3320 }
3321 frame.finish().unwrap();
3322 straggler.finish().unwrap();
3323
3324 let state = producer.state.read();
3325 assert!(
3326 state.lookup.contains_key(&0),
3327 "a group streaming a frame survives expiry"
3328 );
3329 }
3330
3331 #[tokio::test]
3334 async fn parked_reoffer_restarts_the_expiry_clock() {
3335 tokio::time::pause();
3336
3337 let mut producer = track_producer("test", None);
3338 let mut subscriber = producer.subscribe(None);
3339 subscriber.end_at(0);
3340
3341 for seq in 0..2u64 {
3342 let mut group = producer.create_group(seq.into()).unwrap();
3343 group.write_frame(Timestamp::ZERO, b"x".as_slice()).unwrap();
3344 group.finish().unwrap();
3345 }
3346
3347 assert_eq!(subscriber.assert_group().sequence, 0);
3349 subscriber.assert_no_group();
3350
3351 tokio::time::advance(DEFAULT_LATENCY_MAX - Duration::from_secs(1)).await;
3353 subscriber.end_at(1);
3354 let mut reading = subscriber.assert_group();
3355 assert_eq!(reading.sequence, 1);
3356
3357 tokio::time::advance(DEFAULT_LATENCY_MAX - Duration::from_secs(1)).await;
3360 producer.create_group(2u64.into()).unwrap().finish().unwrap();
3361
3362 let frame = reading.read_frame().await.unwrap();
3363 assert!(frame.is_some(), "a just-re-offered group must not expire unread");
3364 }
3365
3366 #[tokio::test]
3367 async fn evict_keeps_max_sequence() {
3368 tokio::time::pause();
3369
3370 let mut producer = track_producer("test", None);
3371 producer.append_group().unwrap(); tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3375
3376 producer.append_group().unwrap(); {
3380 let state = producer.state.read();
3381 assert_eq!(live_groups(&state), 1);
3382 assert_eq!(first_live_sequence(&state), 1);
3383 assert_eq!(state.offset, 1);
3384 }
3385 }
3386
3387 #[tokio::test]
3388 async fn no_eviction_when_fresh() {
3389 tokio::time::pause();
3390
3391 let mut producer = track_producer("test", None);
3392 producer.append_group().unwrap(); producer.append_group().unwrap(); producer.append_group().unwrap(); {
3397 let state = producer.state.read();
3398 assert_eq!(live_groups(&state), 3);
3399 assert_eq!(state.offset, 0);
3400 }
3401 }
3402
3403 #[tokio::test]
3404 async fn consumer_skips_evicted_groups() {
3405 tokio::time::pause();
3406
3407 let mut producer = track_producer("test", None);
3408 producer.append_group().unwrap(); let mut consumer = producer.subscribe(None);
3411
3412 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3413 producer.append_group().unwrap(); let group = consumer.assert_group();
3417 assert_eq!(group.sequence, 1);
3418 }
3419
3420 #[tokio::test]
3421 async fn cache_age_controls_eviction() {
3422 tokio::time::pause();
3423
3424 let mut producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(1)));
3426 producer.append_group().unwrap(); tokio::time::advance(Duration::from_secs(2)).await;
3430 producer.append_group().unwrap(); let state = producer.state.read();
3434 assert_eq!(live_groups(&state), 1);
3435 assert_eq!(first_live_sequence(&state), 1);
3436 }
3437
3438 #[test]
3439 fn latency_max_clamped_to_cache() {
3440 let producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(2)));
3441
3442 let mut subscriber = producer.subscribe(Subscription::default().with_latency_max(Duration::from_secs(10)));
3446 assert_eq!(subscriber.subscription().latency_max, Duration::from_secs(10));
3447 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3448
3449 subscriber
3451 .update(Subscription::default().with_latency_max(Duration::from_millis(500)))
3452 .unwrap();
3453 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_millis(500));
3454
3455 subscriber
3456 .update(Subscription::default().with_latency_max(Duration::ZERO))
3457 .unwrap();
3458 assert_eq!(producer.subscription().unwrap().latency_max, Duration::ZERO);
3459 }
3460
3461 fn track_producer_capped(name: impl Into<Arc<str>>, info: Info, cap: Duration) -> Producer {
3464 let origin = crate::origin::Info::default().with_cache_duration(cap);
3465 Producer::new(Arc::new(broadcast::Info { origin }), name, info)
3466 }
3467
3468 #[test]
3469 fn origin_cache_duration_clamps_latency_max() {
3470 let capped = track_producer_capped(
3473 "test",
3474 Info::default().with_latency_max(Duration::from_secs(60)),
3475 Duration::from_secs(1),
3476 );
3477 assert_eq!(capped.state.read().latency_bound(), Some(Duration::from_secs(1)));
3478 assert_eq!(capped.subscribe(None).info().latency_max, Duration::from_secs(1));
3479
3480 let under = track_producer_capped(
3481 "test",
3482 Info::default().with_latency_max(Duration::from_millis(500)),
3483 Duration::from_secs(1),
3484 );
3485 assert_eq!(under.state.read().latency_bound(), Some(Duration::from_millis(500)));
3486 }
3487
3488 #[tokio::test]
3489 async fn origin_cache_duration_caps_eviction() {
3490 tokio::time::pause();
3491
3492 let mut producer = track_producer_capped(
3494 "test",
3495 Info::default().with_latency_max(Duration::from_secs(60)),
3496 Duration::from_secs(1),
3497 );
3498 producer.append_group().unwrap(); tokio::time::advance(Duration::from_secs(2)).await;
3502 producer.append_group().unwrap(); let state = producer.state.read();
3506 assert_eq!(live_groups(&state), 1);
3507 assert_eq!(first_live_sequence(&state), 1);
3508 }
3509
3510 #[test]
3511 fn latency_max_clamped_via_every_update_path() {
3512 let producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(2)));
3513 let over = Subscription::default().with_latency_max(Duration::from_secs(10));
3514
3515 let mut subscriber = producer.subscribe(over.clone());
3518 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3519
3520 subscriber.control().update(over.clone()).unwrap();
3521 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3522
3523 subscriber.update(over).unwrap();
3524 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3525 }
3526
3527 #[test]
3528 fn latency_max_aggregate_clamps_the_max_across_subscribers() {
3529 let producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(2)));
3530
3531 let _a = producer.subscribe(Subscription::default().with_latency_max(Duration::from_millis(500)));
3534 let _b = producer.subscribe(Subscription::default().with_latency_max(Duration::from_secs(10)));
3535
3536 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3537 }
3538
3539 #[test]
3540 fn subscriber_control_updates_while_read_future_is_pending() {
3541 let producer = track_producer("test", None);
3542 let mut subscriber = producer.subscribe(None);
3543 let control = subscriber.control();
3544
3545 let mut recv = Box::pin(subscriber.recv_group());
3546 assert!(recv.as_mut().now_or_never().is_none());
3547
3548 control
3549 .update(Subscription::default().with_priority(7).with_ordered(false))
3550 .unwrap();
3551
3552 let aggregate = producer.subscription().expect("expected an active subscription");
3553 assert_eq!(aggregate.priority, 7);
3554 assert!(!aggregate.ordered);
3555 }
3556
3557 #[test]
3558 fn dropped_subscriber_leaves_no_ghost_in_aggregate() {
3559 let mut producer = track_producer("test", None);
3564 let a = producer.subscribe(Subscription::default().with_priority(5));
3565
3566 let waiter = kio::Waiter::noop();
3568 assert!(
3569 matches!(producer.poll_subscription_changed(&waiter), Poll::Ready(Ok(Some(_)))),
3570 "one live subscriber should aggregate to Some",
3571 );
3572
3573 drop(a);
3575
3576 assert!(
3578 matches!(producer.poll_subscription_changed(&waiter), Poll::Ready(Ok(None))),
3579 "a dropped subscriber must not linger in the aggregate",
3580 );
3581
3582 assert!(
3584 producer.subscription().is_none(),
3585 "snapshot must exclude a dropped subscriber",
3586 );
3587 }
3588
3589 #[test]
3590 fn dropped_subscriber_wakes_the_aggregate() {
3591 use std::sync::atomic::{AtomicBool, Ordering};
3598
3599 let mut producer = track_producer("test", None);
3600 let a = producer.subscribe(Subscription::default().with_priority(5));
3601
3602 let woken = Arc::new(AtomicBool::new(false));
3603 let waiter = kio::Waiter::new(futures::task::waker(Arc::new(FlagWake(woken.clone()))));
3604
3605 assert!(matches!(
3607 producer.poll_subscription_changed(&waiter),
3608 Poll::Ready(Ok(Some(_)))
3609 ));
3610 assert!(
3611 producer.poll_subscription_changed(&waiter).is_pending(),
3612 "the aggregate is unchanged, so this poll must park",
3613 );
3614 assert!(!woken.load(Ordering::SeqCst), "nothing happened yet");
3615
3616 drop(a);
3617 assert!(
3618 woken.load(Ordering::SeqCst),
3619 "the last subscriber leaving must wake the aggregate watcher",
3620 );
3621 }
3622
3623 struct FlagWake(Arc<std::sync::atomic::AtomicBool>);
3625
3626 impl futures::task::ArcWake for FlagWake {
3627 fn wake_by_ref(arc_self: &Arc<Self>) {
3628 arc_self.0.store(true, std::sync::atomic::Ordering::SeqCst);
3629 }
3630 }
3631
3632 #[tokio::test]
3633 async fn out_of_order_max_sequence_at_front() {
3634 tokio::time::pause();
3635
3636 let mut producer = track_producer("test", None);
3637
3638 producer.create_group(group::Info { sequence: 5 }).unwrap();
3640 producer.create_group(group::Info { sequence: 3 }).unwrap();
3641 producer.create_group(group::Info { sequence: 4 }).unwrap();
3642
3643 {
3645 let state = producer.state.read();
3646 assert_eq!(state.max_sequence, Some(5));
3647 }
3648
3649 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3651
3652 producer.append_group().unwrap(); {
3658 let state = producer.state.read();
3659 assert_eq!(live_groups(&state), 1);
3660 assert_eq!(first_live_sequence(&state), 6);
3661 assert!(!state.lookup.contains_key(&3));
3662 assert!(!state.lookup.contains_key(&4));
3663 assert!(!state.lookup.contains_key(&5));
3664 assert!(state.lookup.contains_key(&6));
3665 }
3666 }
3667
3668 #[tokio::test]
3669 async fn max_sequence_at_front_blocks_trim() {
3670 tokio::time::pause();
3671
3672 let mut producer = track_producer("test", None);
3673
3674 producer.create_group(group::Info { sequence: 5 }).unwrap();
3676
3677 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3678
3679 producer.create_group(group::Info { sequence: 3 }).unwrap();
3681
3682 {
3685 let state = producer.state.read();
3686 assert_eq!(live_groups(&state), 2);
3687 assert_eq!(state.offset, 0);
3688 }
3689
3690 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3692
3693 producer.create_group(group::Info { sequence: 2 }).unwrap();
3695
3696 {
3701 let state = producer.state.read();
3702 assert_eq!(live_groups(&state), 2);
3703 assert_eq!(state.offset, 0);
3704 assert!(state.lookup.contains_key(&5));
3705 assert!(!state.lookup.contains_key(&3));
3706 assert!(state.lookup.contains_key(&2));
3707 }
3708
3709 let mut consumer = producer.subscribe(None);
3711 let group = consumer.assert_group();
3712 assert_eq!(group.sequence, 5);
3714 }
3715
3716 #[tokio::test]
3717 async fn abort_clears_cached_groups() {
3718 let mut producer = track_producer("test", None);
3719 producer.append_group().unwrap();
3720 producer.append_group().unwrap();
3721
3722 let mut consumer = producer.subscribe(None);
3724 assert_eq!(live_groups(&producer.state.read()), 2);
3725
3726 producer.clone().abort(Error::Cancel).unwrap();
3727
3728 {
3729 let state = producer.state.read();
3730 assert!(state.lookup.is_empty(), "cached groups should be dropped on abort");
3731 assert!(state.arrival.is_empty());
3732 assert!(state.evict.is_empty());
3733 }
3734
3735 let result = consumer.recv_group().now_or_never().expect("should not block");
3737 assert!(matches!(result, Err(Error::Cancel)));
3738 }
3739
3740 #[tokio::test]
3741 async fn drop_unfinished_clears_cached_groups() {
3742 let producer = track_producer("test", None);
3743 let mut writer = producer.clone();
3744 writer.append_group().unwrap();
3745
3746 let mut consumer = producer.subscribe(None);
3748 assert_eq!(live_groups(&producer.state.read()), 1);
3749
3750 drop(writer);
3752 drop(producer);
3753
3754 let result = consumer.recv_group().now_or_never().expect("should not block");
3755 assert!(matches!(result, Err(Error::Dropped)));
3756 }
3757
3758 #[tokio::test]
3759 async fn drop_after_abort_does_not_warn() {
3760 let warns = count_drop_warnings("track::Producer dropped without finish", || {
3763 let producer = track_producer("test", None);
3764 let keep = producer.clone();
3765 let mut writer = producer.clone();
3766 let mut group = writer.append_group().unwrap();
3767 group.finish().unwrap();
3768 let _consumer = producer.subscribe(None);
3769 writer.abort(Error::Cancel).unwrap();
3770 drop(keep);
3771 });
3772 assert_eq!(warns, 0, "abort-then-drop must not emit unfinished-producer WARN");
3773 }
3774
3775 #[tokio::test]
3776 async fn drop_unfinished_warns() {
3777 let warns = count_drop_warnings("track::Producer dropped without finish", || {
3778 let producer = track_producer("test", None);
3779 let mut writer = producer.clone();
3780 writer.append_group().unwrap();
3781 let _consumer = producer.subscribe(None);
3782 drop(writer);
3783 drop(producer);
3784 });
3785 assert!(warns >= 1, "unfinished drop must emit unfinished-producer WARN");
3786 }
3787
3788 #[tokio::test]
3789 async fn drop_finished_keeps_cached_groups() {
3790 let mut producer = track_producer("test", None);
3791 producer.append_group().unwrap();
3792 producer.finish().unwrap();
3793
3794 let mut consumer = producer.subscribe(None);
3795 drop(producer);
3796
3797 assert_eq!(consumer.assert_group().sequence, 0);
3799 let done = consumer.recv_group().now_or_never().expect("should not block").unwrap();
3800 assert!(done.is_none(), "consumer should drain then see clean finish");
3801 }
3802
3803 #[test]
3804 fn append_finish_cannot_be_rewritten() {
3805 let mut producer = track_producer("test", None);
3806
3807 assert!(producer.finish().is_ok());
3809 assert!(producer.finish().is_err());
3810 assert!(producer.append_group().is_err());
3811 }
3812
3813 #[test]
3814 fn finish_after_groups() {
3815 let mut producer = track_producer("test", None);
3816
3817 producer.append_group().unwrap();
3818 assert!(producer.finish().is_ok());
3819 assert!(producer.finish().is_err());
3820 assert!(producer.append_group().is_err());
3821 }
3822
3823 #[test]
3824 fn finish_at_rejects_a_boundary_at_or_below_the_live_edge() {
3825 let mut producer = track_producer("test", None);
3826 producer.create_group(group::Info { sequence: 5 }).unwrap();
3827
3828 assert!(producer.finish_at(4).is_err());
3831 assert!(producer.finish_at(5).is_err());
3832 assert!(producer.finish_at(6).is_ok());
3833
3834 {
3835 let state = producer.state.read();
3836 assert_eq!(state.final_sequence, Some(6));
3837 }
3838
3839 assert!(producer.finish_at(6).is_err());
3841 assert!(producer.create_group(group::Info { sequence: 4 }).is_ok());
3842 assert!(producer.create_group(group::Info { sequence: 6 }).is_err());
3843 }
3844
3845 #[test]
3846 fn final_sequence_reports_the_declared_boundary() {
3847 let mut producer = track_producer("test", None);
3848 assert_eq!(producer.final_sequence(), None);
3849
3850 producer.create_group(group::Info { sequence: 5 }).unwrap();
3851 assert_eq!(producer.final_sequence(), None, "a group does not declare a boundary");
3852
3853 producer.finish_at(9).unwrap();
3854 assert_eq!(producer.final_sequence(), Some(9));
3855
3856 assert!(producer.finish().is_err());
3858 }
3859
3860 #[test]
3861 fn final_sequence_reports_the_live_edge_after_finish() {
3862 let mut producer = track_producer("test", None);
3863 producer.create_group(group::Info { sequence: 5 }).unwrap();
3864 producer.finish().unwrap();
3865 assert_eq!(producer.final_sequence(), Some(6));
3866 }
3867
3868 #[tokio::test]
3869 async fn finish_at_declares_a_future_boundary() {
3870 let mut producer = track_producer("test", None);
3871 producer.create_group(group::Info { sequence: 5 }).unwrap();
3872
3873 producer.finish_at(7).unwrap();
3875
3876 let mut consumer = producer.subscribe(None);
3877 assert_eq!(consumer.assert_group().sequence, 5);
3878
3879 let boundary = consumer
3882 .finished()
3883 .now_or_never()
3884 .expect("boundary is known immediately")
3885 .expect("would have errored");
3886 assert_eq!(boundary, 7);
3887 assert!(
3888 consumer.recv_group().now_or_never().is_none(),
3889 "should wait for the outstanding group"
3890 );
3891
3892 producer.create_group(group::Info { sequence: 6 }).unwrap();
3894 assert_eq!(consumer.assert_group().sequence, 6);
3895 let done = consumer
3896 .recv_group()
3897 .now_or_never()
3898 .expect("should not block")
3899 .expect("would have errored");
3900 assert!(done.is_none(), "track completes once the boundary is reached");
3901 }
3902
3903 #[tokio::test]
3904 async fn recv_group_finishes_without_waiting_for_gaps() {
3905 let mut producer = track_producer("test", None);
3906 producer.create_group(group::Info { sequence: 1 }).unwrap();
3907 producer.finish().unwrap();
3908
3909 let mut consumer = producer.subscribe(None);
3910 assert_eq!(consumer.assert_group().sequence, 1);
3911
3912 let done = consumer
3913 .recv_group()
3914 .now_or_never()
3915 .expect("should not block")
3916 .expect("would have errored");
3917 assert!(done.is_none(), "track should finish without waiting for gaps");
3918 }
3919
3920 #[tokio::test]
3921 async fn next_group_skips_late_arrivals() {
3922 let mut producer = track_producer("test", None);
3923 let mut consumer = producer.subscribe(None);
3924
3925 producer.create_group(group::Info { sequence: 5 }).unwrap();
3927 let group = consumer
3928 .next_group()
3929 .now_or_never()
3930 .expect("should not block")
3931 .expect("would have errored")
3932 .expect("track should not be closed");
3933 assert_eq!(group.sequence, 5);
3934
3935 producer.create_group(group::Info { sequence: 3 }).unwrap();
3937 producer.create_group(group::Info { sequence: 4 }).unwrap();
3939 producer.create_group(group::Info { sequence: 7 }).unwrap();
3941
3942 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, 7);
3949
3950 assert!(
3952 consumer.next_group().now_or_never().is_none(),
3953 "should block waiting for a higher sequence"
3954 );
3955 }
3956
3957 #[tokio::test]
3958 async fn next_group_returns_arrivals_in_order() {
3959 let mut producer = track_producer("test", None);
3960 let mut consumer = producer.subscribe(None);
3961
3962 producer.create_group(group::Info { sequence: 3 }).unwrap();
3964 producer.create_group(group::Info { sequence: 5 }).unwrap();
3965
3966 let group = consumer
3967 .next_group()
3968 .now_or_never()
3969 .expect("should not block")
3970 .expect("would have errored")
3971 .expect("track should not be closed");
3972 assert_eq!(group.sequence, 3);
3973
3974 let group = consumer
3975 .next_group()
3976 .now_or_never()
3977 .expect("should not block")
3978 .expect("would have errored")
3979 .expect("track should not be closed");
3980 assert_eq!(group.sequence, 5);
3981 }
3982
3983 #[tokio::test]
3984 async fn next_group_and_recv_group_use_independent_cursors() {
3985 let mut producer = track_producer("test", None);
3986 let mut consumer = producer.subscribe(None);
3987
3988 producer.create_group(group::Info { sequence: 5 }).unwrap();
3990 producer.create_group(group::Info { sequence: 3 }).unwrap();
3991
3992 let group = consumer
3995 .next_group()
3996 .now_or_never()
3997 .expect("should not block")
3998 .expect("would have errored")
3999 .expect("track should not be closed");
4000 assert_eq!(group.sequence, 3);
4001
4002 assert_eq!(consumer.assert_group().sequence, 5);
4005 }
4006
4007 #[tokio::test]
4008 async fn end_at_caps_next_group() {
4009 let mut producer = track_producer("test", None);
4010 let mut consumer = producer.subscribe(None);
4011
4012 for s in 0..6 {
4013 producer.create_group(group::Info { sequence: s }).unwrap();
4014 }
4015
4016 consumer.end_at(2);
4017
4018 assert_eq!(
4020 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4021 0
4022 );
4023 assert_eq!(
4024 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4025 1
4026 );
4027 assert_eq!(
4028 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4029 2
4030 );
4031
4032 assert!(
4034 consumer.next_group().now_or_never().is_none(),
4035 "capped consumer must block instead of returning out-of-range groups"
4036 );
4037 }
4038
4039 #[tokio::test]
4040 async fn end_at_release_drains_cached_groups() {
4041 let mut producer = track_producer("test", None);
4042 let mut consumer = producer.subscribe(None);
4043
4044 for s in 0..6 {
4045 producer.create_group(group::Info { sequence: s }).unwrap();
4046 }
4047
4048 consumer.end_at(1);
4049 assert_eq!(
4050 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4051 0
4052 );
4053 assert_eq!(
4054 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4055 1
4056 );
4057 assert!(consumer.next_group().now_or_never().is_none(), "capped at 1");
4058
4059 consumer.end_at(4);
4061 assert_eq!(
4062 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4063 2
4064 );
4065 assert_eq!(
4066 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4067 3
4068 );
4069 assert_eq!(
4070 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4071 4
4072 );
4073 assert!(consumer.next_group().now_or_never().is_none(), "capped at 4");
4074
4075 consumer.end_at(None);
4077 assert_eq!(
4078 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4079 5
4080 );
4081 assert!(consumer.next_group().now_or_never().is_none(), "no more groups");
4082 }
4083
4084 #[tokio::test]
4085 async fn end_at_lower_than_cursor_parks_consumer() {
4086 let mut producer = track_producer("test", None);
4087 let mut consumer = producer.subscribe(None);
4088
4089 for s in 0..3 {
4090 producer.create_group(group::Info { sequence: s }).unwrap();
4091 }
4092
4093 assert_eq!(
4095 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4096 0
4097 );
4098 assert_eq!(
4099 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4100 1
4101 );
4102 assert_eq!(
4103 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4104 2
4105 );
4106
4107 consumer.end_at(1);
4109 producer.create_group(group::Info { sequence: 3 }).unwrap();
4110 producer.create_group(group::Info { sequence: 4 }).unwrap();
4111 assert!(
4112 consumer.next_group().now_or_never().is_none(),
4113 "cap is below cursor; nothing returnable until cap rises"
4114 );
4115
4116 consumer.end_at(None);
4118 assert_eq!(
4119 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4120 3
4121 );
4122 assert_eq!(
4123 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4124 4
4125 );
4126 }
4127
4128 #[tokio::test]
4129 async fn end_at_toggling_around_late_arrivals() {
4130 let mut producer = track_producer("test", None);
4131 let mut consumer = producer.subscribe(None);
4132
4133 consumer.end_at(5);
4134
4135 producer.create_group(group::Info { sequence: 2 }).unwrap();
4137 producer.create_group(group::Info { sequence: 5 }).unwrap();
4138 producer.create_group(group::Info { sequence: 3 }).unwrap();
4139 producer.create_group(group::Info { sequence: 8 }).unwrap();
4141 producer.create_group(group::Info { sequence: 4 }).unwrap();
4142
4143 assert_eq!(
4145 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4146 2
4147 );
4148 assert_eq!(
4149 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4150 3
4151 );
4152 assert_eq!(
4153 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4154 4
4155 );
4156 assert_eq!(
4157 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4158 5
4159 );
4160 assert!(consumer.next_group().now_or_never().is_none());
4162
4163 consumer.end_at(10);
4165 assert_eq!(
4166 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4167 8
4168 );
4169 }
4170
4171 #[tokio::test]
4175 async fn end_at_parks_recv_group() {
4176 let mut producer = track_producer("test", None);
4177 let mut consumer = producer.subscribe(None);
4178
4179 for s in 0..3 {
4180 producer.create_group(group::Info { sequence: s }).unwrap();
4181 }
4182
4183 consumer.end_at(1);
4184 assert_eq!(
4185 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4186 0
4187 );
4188 assert_eq!(
4189 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4190 1
4191 );
4192 assert!(consumer.recv_group().now_or_never().is_none(), "capped at 1");
4193
4194 producer.finish().unwrap();
4196 assert!(
4197 consumer.recv_group().now_or_never().is_none(),
4198 "still parked after finish"
4199 );
4200
4201 consumer.end_at(None);
4202 assert_eq!(
4203 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4204 2
4205 );
4206 assert!(
4207 matches!(consumer.recv_group().now_or_never(), Some(Ok(None))),
4208 "finished once the parked group drains"
4209 );
4210 }
4211
4212 #[tokio::test]
4215 async fn recv_group_serves_arrivals_behind_the_cap() {
4216 let mut producer = track_producer("test", None);
4217 let mut consumer = producer.subscribe(None);
4218
4219 consumer.end_at(1);
4220
4221 producer.create_group(group::Info { sequence: 2 }).unwrap();
4223 producer.create_group(group::Info { sequence: 0 }).unwrap();
4224 producer.create_group(group::Info { sequence: 1 }).unwrap();
4225
4226 assert_eq!(
4227 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4228 0
4229 );
4230 assert_eq!(
4231 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4232 1
4233 );
4234 assert!(consumer.recv_group().now_or_never().is_none(), "capped at 1");
4235
4236 consumer.end_at(2);
4237 assert_eq!(
4238 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4239 2
4240 );
4241 }
4242
4243 #[tokio::test]
4246 async fn start_at_drops_parked_recv_groups() {
4247 let mut producer = track_producer("test", None);
4248 let mut consumer = producer.subscribe(None);
4249
4250 consumer.end_at(0);
4251 producer.create_group(group::Info { sequence: 1 }).unwrap();
4252 assert!(
4253 consumer.recv_group().now_or_never().is_none(),
4254 "group 1 parked at the cap"
4255 );
4256
4257 consumer.start_at(2);
4258 consumer.end_at(None);
4259 producer.create_group(group::Info { sequence: 2 }).unwrap();
4260 assert_eq!(
4261 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4262 2,
4263 "the overtaken parked group is dropped, not re-offered"
4264 );
4265 }
4266
4267 #[tokio::test]
4271 async fn evicted_parked_recv_groups_are_dropped() {
4272 let mut producer = track_producer("test", None);
4273 let mut consumer = producer.subscribe(None);
4274
4275 producer.create_group(group::Info { sequence: 0 }).unwrap();
4276 assert_eq!(
4277 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4278 0
4279 );
4280
4281 consumer.end_at(0);
4282 let straggler = producer.create_group(group::Info { sequence: 1 }).unwrap();
4283 assert!(
4284 consumer.recv_group().now_or_never().is_none(),
4285 "group 1 parked at the cap"
4286 );
4287
4288 straggler.abort(Error::Old).unwrap();
4290 producer.finish().unwrap();
4291
4292 consumer.end_at(None);
4293 assert!(
4294 matches!(consumer.recv_group().now_or_never(), Some(Ok(None))),
4295 "a dead parked group must not be delivered or hold the stream open"
4296 );
4297 }
4298
4299 #[tokio::test]
4303 async fn evicted_parked_group_wakes_the_clean_end() {
4304 use std::sync::atomic::{AtomicUsize, Ordering};
4305 use std::task::{Context, Wake};
4306
4307 struct CountWaker(AtomicUsize);
4310 impl Wake for CountWaker {
4311 fn wake(self: std::sync::Arc<Self>) {
4312 self.0.fetch_add(1, Ordering::SeqCst);
4313 }
4314 }
4315
4316 let mut producer = track_producer("test", None);
4317 let mut consumer = producer.subscribe(None);
4318
4319 producer.create_group(group::Info { sequence: 0 }).unwrap();
4320 assert_eq!(
4321 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4322 0
4323 );
4324
4325 consumer.end_at(0);
4326 let straggler = producer.create_group(group::Info { sequence: 1 }).unwrap();
4327 assert!(consumer.recv_group().now_or_never().is_none(), "parked at the cap");
4328 producer.finish().unwrap();
4329
4330 let counter = std::sync::Arc::new(CountWaker(AtomicUsize::new(0)));
4331 let waker = std::task::Waker::from(counter.clone());
4332 let mut cx = Context::from_waker(&waker);
4333 let mut fut = std::pin::pin!(consumer.recv_group());
4334 assert!(
4335 fut.as_mut().poll(&mut cx).is_pending(),
4336 "the parked group holds it open"
4337 );
4338
4339 straggler.abort(Error::Old).unwrap();
4340 assert!(counter.0.load(Ordering::SeqCst) > 0, "the eviction wakeup was lost");
4341 assert!(matches!(fut.as_mut().poll(&mut cx), Poll::Ready(Ok(None))));
4342 }
4343
4344 #[tokio::test]
4345 async fn read_frame_returns_single_frame_per_group() {
4346 let mut producer = track_producer("test", None);
4347 let mut consumer = producer.subscribe(None);
4348
4349 producer.write_frame(Timestamp::ZERO, b"hello".as_slice()).unwrap();
4350 producer.write_frame(Timestamp::ZERO, b"world".as_slice()).unwrap();
4351
4352 let frame = consumer
4353 .read_frame()
4354 .now_or_never()
4355 .expect("should not block")
4356 .expect("would have errored")
4357 .expect("track should not be closed");
4358 assert_eq!(&frame.payload[..], b"hello");
4359
4360 let frame = consumer
4361 .read_frame()
4362 .now_or_never()
4363 .expect("should not block")
4364 .expect("would have errored")
4365 .expect("track should not be closed");
4366 assert_eq!(&frame.payload[..], b"world");
4367 }
4368
4369 #[test]
4370 fn write_frame_rejects_an_oversized_frame_before_appending_its_group() {
4371 let mut producer = track_producer("test", None);
4372 let frame = bytes::Bytes::from(vec![0; group::MAX_CACHE_BYTES as usize + 1]);
4373
4374 assert!(matches!(
4375 producer.write_frame(Timestamp::ZERO, frame),
4376 Err(Error::FrameTooLarge)
4377 ));
4378 assert_eq!(producer.latest(), None, "the rejected frame did not publish a group");
4379 }
4380
4381 #[tokio::test]
4382 async fn read_frame_preserves_timestamp() {
4383 let mut producer = track_producer("test", None);
4384 let mut consumer = producer.subscribe(None);
4385
4386 producer
4387 .write_frame(Timestamp::from_micros(20_000).unwrap(), b"hello".as_slice())
4388 .unwrap();
4389
4390 let frame = consumer
4391 .read_frame()
4392 .now_or_never()
4393 .expect("should not block")
4394 .expect("would have errored")
4395 .expect("track should not be closed");
4396 assert_eq!(frame.timestamp.as_micros(), 20_000);
4397 assert_eq!(&frame.payload[..], b"hello");
4398 }
4399
4400 #[tokio::test]
4401 async fn read_frame_skips_stalled_group_for_newer_ready_frame() {
4402 let mut producer = track_producer("test", None);
4403 let mut consumer = producer.subscribe(None);
4404
4405 let _stalled = producer.create_group(group::Info { sequence: 3 }).unwrap();
4407 let mut g5 = producer.create_group(group::Info { sequence: 5 }).unwrap();
4409 g5.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"later"))
4410 .unwrap();
4411 g5.finish().unwrap();
4412
4413 let frame = consumer
4415 .read_frame()
4416 .now_or_never()
4417 .expect("should not block on stalled earlier group")
4418 .expect("would have errored")
4419 .expect("track should not be closed");
4420 assert_eq!(&frame.payload[..], b"later");
4421 }
4422
4423 #[tokio::test]
4424 async fn read_frame_discards_rest_of_multi_frame_group() {
4425 let mut producer = track_producer("test", None);
4426 let mut consumer = producer.subscribe(None);
4427
4428 let mut g0 = producer.create_group(group::Info { sequence: 0 }).unwrap();
4430 g0.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"one"))
4431 .unwrap();
4432 g0.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"two"))
4433 .unwrap();
4434 g0.finish().unwrap();
4435
4436 producer.write_frame(Timestamp::ZERO, b"next".as_slice()).unwrap();
4438
4439 let frame = consumer
4440 .read_frame()
4441 .now_or_never()
4442 .expect("should not block")
4443 .expect("would have errored")
4444 .expect("track should not be closed");
4445 assert_eq!(&frame.payload[..], b"one");
4446
4447 let frame = consumer
4449 .read_frame()
4450 .now_or_never()
4451 .expect("should not block")
4452 .expect("would have errored")
4453 .expect("track should not be closed");
4454 assert_eq!(&frame.payload[..], b"next");
4455 }
4456
4457 #[tokio::test]
4458 async fn read_frame_waits_for_pending_group_after_finish() {
4459 let mut producer = track_producer("test", None);
4462 let mut consumer = producer.subscribe(None);
4463
4464 let mut g0 = producer.create_group(group::Info { sequence: 0 }).unwrap();
4465 producer.finish().unwrap();
4466
4467 assert!(
4469 consumer.read_frame().now_or_never().is_none(),
4470 "read_frame must block on a pending group even after finish()"
4471 );
4472
4473 g0.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"late"))
4475 .unwrap();
4476 let frame = consumer
4477 .read_frame()
4478 .now_or_never()
4479 .expect("should not block once a frame is written")
4480 .expect("would have errored")
4481 .expect("track should not be closed");
4482 assert_eq!(&frame.payload[..], b"late");
4483 }
4484
4485 #[tokio::test]
4486 async fn read_frame_respects_start_at() {
4487 let mut producer = track_producer("test", None);
4490 let mut consumer = producer.subscribe(None);
4491 consumer.start_at(5);
4492
4493 let mut g3 = producer.create_group(group::Info { sequence: 3 }).unwrap();
4495 g3.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"skip-me"))
4496 .unwrap();
4497 g3.finish().unwrap();
4498
4499 let mut g5 = producer.create_group(group::Info { sequence: 5 }).unwrap();
4500 g5.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"keep"))
4501 .unwrap();
4502 g5.finish().unwrap();
4503
4504 let frame = consumer
4505 .read_frame()
4506 .now_or_never()
4507 .expect("should not block")
4508 .expect("would have errored")
4509 .expect("track should not be closed");
4510 assert_eq!(&frame.payload[..], b"keep");
4511 }
4512
4513 #[tokio::test]
4514 async fn read_frame_returns_none_when_finished() {
4515 let mut producer = track_producer("test", None);
4516 let mut consumer = producer.subscribe(None);
4517
4518 producer.write_frame(Timestamp::ZERO, b"only".as_slice()).unwrap();
4519 producer.finish().unwrap();
4520
4521 let frame = consumer
4522 .read_frame()
4523 .now_or_never()
4524 .expect("should not block")
4525 .expect("would have errored")
4526 .expect("track should not be closed");
4527 assert_eq!(&frame.payload[..], b"only");
4528
4529 let done = consumer
4530 .read_frame()
4531 .now_or_never()
4532 .expect("should not block")
4533 .expect("would have errored");
4534 assert!(done.is_none());
4535 }
4536
4537 #[test]
4538 fn append_group_returns_bounds_exceeded_on_sequence_overflow() {
4539 let mut producer = track_producer("test", None);
4540 {
4541 let mut state = producer.state.write().ok().unwrap();
4542 state.max_sequence = Some(u64::MAX);
4543 }
4544
4545 assert!(matches!(producer.append_group(), Err(Error::BoundsExceeded(_))));
4546 }
4547
4548 #[tokio::test]
4549 async fn fetch_cache_hit() {
4550 let mut producer = track_producer("test", None);
4551
4552 let mut group = producer.append_group().unwrap(); group
4555 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hello"))
4556 .unwrap();
4557 group.finish().unwrap();
4558
4559 let dynamic = producer.dynamic();
4562 let consumer = producer.consume();
4563 assert!(consumer.peek_group(0).is_some());
4564 let mut g = consumer.fetch_group(0, None).await.unwrap();
4565 assert_eq!(g.sequence, 0);
4566 assert_eq!(&g.read_frame().await.unwrap().unwrap().payload[..], b"hello");
4567
4568 assert!(dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending());
4570 }
4571
4572 #[tokio::test]
4573 async fn fetch_miss_signals_dynamic() {
4574 let producer = track_producer("test", None);
4575 let dynamic = producer.dynamic();
4576 let consumer = producer.consume();
4577
4578 assert!(consumer.peek_group(5).is_none());
4582 let pending = consumer.fetch_group(5, group::Fetch::default().with_priority(7));
4583 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
4584
4585 let req = dynamic
4586 .requested_group()
4587 .now_or_never()
4588 .expect("should not block")
4589 .unwrap();
4590 assert_eq!(req.sequence(), 5);
4591 assert_eq!(req.priority(), 7);
4592
4593 let mut group = req.accept(None).unwrap();
4595 group
4596 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hi"))
4597 .unwrap();
4598 group.finish().unwrap();
4599
4600 let mut g = pending.await.unwrap();
4601 assert_eq!(g.sequence, 5);
4602 assert_eq!(&g.read_frame().await.unwrap().unwrap().payload[..], b"hi");
4603 }
4604
4605 #[tokio::test]
4606 async fn fetch_miss_rejects() {
4607 let producer = track_producer("test", None);
4608 let dynamic = producer.dynamic();
4609 let consumer = producer.consume();
4610
4611 let pending = consumer.fetch_group(5, None);
4612 let req = dynamic
4613 .requested_group()
4614 .now_or_never()
4615 .expect("should not block")
4616 .unwrap();
4617
4618 req.reject(Error::Cancel);
4619 assert!(matches!(pending.await, Err(Error::Cancel)));
4620 let fetch = producer.state.read().fetch.clone();
4621 assert!(fetch.read().is_empty());
4622 }
4623
4624 #[tokio::test]
4625 async fn fetch_miss_drop_rejects() {
4626 let producer = track_producer("test", None);
4627 let dynamic = producer.dynamic();
4628 let consumer = producer.consume();
4629
4630 let pending = consumer.fetch_group(5, None);
4631 let req = dynamic
4632 .requested_group()
4633 .now_or_never()
4634 .expect("should not block")
4635 .unwrap();
4636
4637 drop(req);
4638 assert!(matches!(pending.await, Err(Error::Dropped)));
4639 }
4640
4641 #[tokio::test]
4642 async fn fetch_reject_does_not_poison_retry() {
4643 let producer = track_producer("test", None);
4644 let dynamic = producer.dynamic();
4645 let consumer = producer.consume();
4646
4647 let pending = consumer.fetch_group(5, None);
4648 let req = dynamic
4649 .requested_group()
4650 .now_or_never()
4651 .expect("should not block")
4652 .unwrap();
4653 req.reject(Error::Cancel);
4654 assert!(matches!(pending.await, Err(Error::Cancel)));
4655
4656 let retry = consumer.fetch_group(5, None);
4657 let req = dynamic
4658 .requested_group()
4659 .now_or_never()
4660 .expect("should not block")
4661 .unwrap();
4662 let mut group = req.accept(None).unwrap();
4663 group
4664 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"retry"))
4665 .unwrap();
4666 group.finish().unwrap();
4667
4668 let mut group = retry.await.unwrap();
4669 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"retry");
4670 }
4671
4672 #[tokio::test]
4673 async fn fetch_coalesces_concurrent() {
4674 let producer = track_producer("test", None);
4675 let dynamic = producer.dynamic();
4676 let consumer = producer.consume();
4677
4678 let first = consumer.fetch_group(5, group::Fetch::default().with_priority(1));
4681 let second = consumer.fetch_group(5, group::Fetch::default().with_priority(7));
4682 assert!(kio::Pollable::poll(&*first, &kio::Waiter::noop()).is_pending());
4683
4684 let req = dynamic
4685 .requested_group()
4686 .now_or_never()
4687 .expect("should not block")
4688 .unwrap();
4689 assert_eq!(req.sequence(), 5);
4690 assert_eq!(req.priority(), 7);
4691 assert!(
4692 dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending(),
4693 "the second fetch queued a duplicate request"
4694 );
4695
4696 let third = consumer.fetch_group(5, None);
4698
4699 let mut group = req.accept(None).unwrap();
4701 group
4702 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hi"))
4703 .unwrap();
4704 group.finish().unwrap();
4705
4706 assert_eq!(first.await.unwrap().sequence, 5);
4707 assert_eq!(second.await.unwrap().sequence, 5);
4708 assert_eq!(third.await.unwrap().sequence, 5);
4709 }
4710
4711 #[tokio::test]
4712 async fn fetch_coalesced_reject_fails_all() {
4713 let producer = track_producer("test", None);
4714 let dynamic = producer.dynamic();
4715 let consumer = producer.consume();
4716
4717 let first = consumer.fetch_group(5, None);
4718 let second = consumer.fetch_group(5, None);
4719 let req = dynamic
4720 .requested_group()
4721 .now_or_never()
4722 .expect("should not block")
4723 .unwrap();
4724 req.reject(Error::Cancel);
4725
4726 assert!(matches!(first.await, Err(Error::Cancel)));
4727 assert!(matches!(second.await, Err(Error::Cancel)));
4728
4729 let retry = consumer.fetch_group(5, None);
4731 assert!(kio::Pollable::poll(&*retry, &kio::Waiter::noop()).is_pending());
4732 let req = dynamic
4733 .requested_group()
4734 .now_or_never()
4735 .expect("should not block")
4736 .unwrap();
4737 assert_eq!(req.sequence(), 5);
4738 }
4739
4740 #[tokio::test]
4741 async fn fetch_queued_fails_when_handlers_leave() {
4742 let producer = track_producer("test", None);
4743 let dynamic = producer.dynamic();
4744 let consumer = producer.consume();
4745
4746 let pending = consumer.fetch_group(5, None);
4748 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
4749 drop(dynamic);
4750 assert!(matches!(pending.await, Err(Error::NotFound)));
4751
4752 let fetch = producer.state.read().fetch.clone();
4754 assert!(fetch.read().is_empty());
4755 }
4756
4757 #[tokio::test]
4758 async fn fetch_miss_no_dynamic_not_found() {
4759 let mut producer = track_producer("test", None);
4762 producer.append_group().unwrap(); let consumer = producer.consume();
4764 assert!(matches!(consumer.fetch_group(5, None).await, Err(Error::NotFound)));
4765 }
4766
4767 #[tokio::test]
4768 async fn fetch_past_final_not_found() {
4769 let mut producer = track_producer("test", None);
4770 producer.append_group().unwrap(); producer.finish().unwrap(); let dynamic = producer.dynamic();
4776 let consumer = producer.consume();
4777 assert!(matches!(consumer.fetch_group(5, None).await, Err(Error::NotFound)));
4778
4779 assert!(dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending());
4781 }
4782
4783 fn pooled_producer(capacity: u64) -> (Producer, cache::Pool) {
4785 let pool = cache::Pool::new(capacity);
4786 let broadcast = broadcast::Info {
4787 origin: crate::origin::Info::default().with_pool(pool.clone()),
4788 ..Default::default()
4789 };
4790 let producer = Producer::new(Arc::new(broadcast), "test", None);
4791 (producer, pool)
4792 }
4793
4794 fn finished_group(producer: &mut Producer, size: usize) -> u64 {
4795 let mut group = producer.append_group().unwrap();
4796 group
4797 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; size]))
4798 .unwrap();
4799 group.finish().unwrap();
4800 group.sequence
4801 }
4802
4803 #[tokio::test]
4806 async fn debt_evicts_oldest_group() {
4807 tokio::time::pause();
4808
4809 let (mut producer, pool) = pooled_producer(10_000);
4811
4812 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let consumer = producer.consume();
4817 assert!(consumer.peek_group(0).is_none(), "oldest group is evicted");
4818 assert!(consumer.peek_group(2).is_some(), "latest group survives");
4819 assert!(
4822 pool.used() <= 2 * (10_000 + cache::ENTRY_OVERHEAD),
4823 "usage hovers near capacity: {}",
4824 pool.used()
4825 );
4826
4827 let mut subscriber = producer.subscribe(None);
4829 assert!(subscriber.assert_group().sequence > 0, "evicted group is not delivered");
4830 }
4831
4832 #[tokio::test]
4834 async fn latest_group_never_evicted() {
4835 tokio::time::pause();
4836
4837 let (mut producer, pool) = pooled_producer(100);
4839 finished_group(&mut producer, 1000); assert!(pool.used() > 100, "the latest may exceed the budget");
4841
4842 finished_group(&mut producer, 1000); finished_group(&mut producer, 1000); let consumer = producer.consume();
4847 assert!(consumer.peek_group(0).is_none());
4848 let mut group = consumer.peek_group(2).expect("latest survives");
4849 assert_eq!(group.read_frame().await.unwrap().unwrap().payload.len(), 1000);
4850 }
4851
4852 #[tokio::test]
4856 async fn fetch_refresh_survives_eviction() {
4857 tokio::time::pause();
4858
4859 let (mut producer, _pool) = pooled_producer(10_000);
4860 let consumer = producer.consume();
4861
4862 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_secs(1)).await;
4864 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_secs(1)).await;
4866 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_millis(500)).await;
4868
4869 let mut fetched = consumer.fetch_group(0, None).await.unwrap();
4871 assert_eq!(fetched.read_frame().await.unwrap().unwrap().payload.len(), 3_000);
4872 tokio::time::advance(Duration::from_millis(500)).await;
4873
4874 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_secs(1)).await;
4878 finished_group(&mut producer, 3_000); assert!(consumer.peek_group(0).is_some(), "refreshed group survives");
4881 assert!(consumer.peek_group(1).is_none(), "unread group is evicted instead");
4882 }
4883
4884 #[tokio::test]
4887 async fn eviction_aborts_readers() {
4888 tokio::time::pause();
4889
4890 let (mut producer, _pool) = pooled_producer(10_000);
4891 let mut subscriber = producer.subscribe(None);
4892
4893 finished_group(&mut producer, 10_000); let mut group0 = subscriber.assert_group();
4895
4896 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let read = group0.read_frame().await;
4900 assert!(matches!(read, Err(Error::Evicted)), "expected Evicted, got {read:?}");
4901 }
4902
4903 #[tokio::test]
4907 async fn small_writes_carry_debt() {
4908 tokio::time::pause();
4909
4910 let (mut producer, pool) = pooled_producer(22_000);
4911 let consumer = producer.consume();
4912
4913 finished_group(&mut producer, 20_000); for _ in 0..3 {
4918 finished_group(&mut producer, 1_000);
4919 }
4920 assert!(consumer.peek_group(0).is_some(), "debt smaller than the victim carries");
4921
4922 for _ in 0..20 {
4924 finished_group(&mut producer, 1_000);
4925 }
4926 assert!(
4927 consumer.peek_group(0).is_none(),
4928 "accumulated debt evicts the large group"
4929 );
4930 assert!(pool.used() <= 24_000, "usage hovers near capacity: {}", pool.used());
4933 }
4934
4935 #[tokio::test]
4939 async fn payment_capped_per_write() {
4940 tokio::time::pause();
4941
4942 let (mut producer, pool) = pooled_producer(1 << 40);
4943 for _ in 0..10 {
4944 finished_group(&mut producer, 1_000);
4945 }
4946
4947 pool.resize(100);
4949 let before = pool.used();
4950
4951 finished_group(&mut producer, 1_000);
4953
4954 let consumer = producer.consume();
4955 assert!(consumer.peek_group(0).is_none(), "the oldest groups are evicted");
4956 assert!(consumer.peek_group(1).is_none());
4957 assert!(consumer.peek_group(2).is_some(), "the backlog drains gradually");
4958 assert!(pool.used() > before - 4_000, "one write must not dump the backlog");
4959 }
4960
4961 #[tokio::test]
4965 async fn accept_preserves_write_accounting() {
4966 tokio::time::pause();
4967
4968 let pool = cache::Pool::new(12_000);
4969 let broadcast = broadcast::Info {
4970 origin: crate::origin::Info::default().with_pool(pool.clone()),
4971 ..Default::default()
4972 };
4973 let request = Request::new(Arc::new(broadcast), "test");
4974 let dynamic = request.dynamic();
4975 let consumer = request.consume();
4976
4977 let pending = consumer.fetch_group(0, None);
4979 let req = dynamic
4980 .requested_group()
4981 .now_or_never()
4982 .expect("should not block")
4983 .unwrap();
4984 let mut backfill = req.accept(None).unwrap();
4985 pending.await.unwrap();
4986 backfill
4987 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 30_000]))
4988 .unwrap();
4989
4990 let mut producer = request.accept(None);
4993 producer.append_group().unwrap().finish().unwrap();
4994 producer.append_group().unwrap().finish().unwrap();
4995
4996 assert!(
4997 producer.consume().peek_group(0).is_none(),
4998 "pre-accept backfill growth is reclaimed after accept"
4999 );
5000 assert!(pool.used() <= 13_000, "usage converges: {}", pool.used());
5001 }
5002
5003 #[tokio::test]
5006 async fn recreated_sequence_bounds_eviction_hints() {
5007 let (mut producer, _pool) = pooled_producer(1 << 40);
5008 producer.create_group(5u64.into()).unwrap().finish().unwrap();
5009
5010 for _ in 0..200 {
5011 let group = producer.create_group(1u64.into()).unwrap();
5012 group.abort(Error::Cancel).unwrap();
5013 }
5014
5015 let state = producer.state.read();
5016 assert!(
5017 state.evict.len() <= 2 * state.lookup.len() + EVICT_SLACK,
5018 "stale hints are compacted: {} entries for {} slots",
5019 state.evict.len(),
5020 state.lookup.len()
5021 );
5022 }
5023
5024 #[tokio::test]
5027 async fn same_tick_write_outranks_inserted() {
5028 tokio::time::pause();
5029
5030 let (mut producer, _pool) = pooled_producer(10_000);
5032
5033 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();
5040 assert!(consumer.peek_group(0).is_none(), "insert-only content pays first");
5041 assert!(consumer.peek_group(1).is_some(), "same-tick written content survives");
5042 }
5043
5044 #[tokio::test]
5047 async fn frame_only_writer_pays() {
5048 tokio::time::pause();
5049
5050 let (mut producer, pool) = pooled_producer(2_000);
5051 let mut demoted = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); demoted
5057 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 300_000]))
5058 .unwrap();
5059
5060 assert!(
5061 pool.used() <= 5_000,
5062 "the frame write settled the debt: {}",
5063 pool.used()
5064 );
5065 assert!(matches!(demoted.finish(), Err(Error::Evicted)));
5066 }
5067
5068 #[tokio::test]
5071 async fn each_track_owns_its_account() {
5072 let broadcast = Arc::new(broadcast::Info::default());
5073 let info = Info::default();
5074 let a = Producer::new(broadcast.clone(), "a", info.clone());
5075 let b = Producer::new(broadcast, "b", info);
5076
5077 let a = a.state.read().cache.clone();
5078 let b = b.state.read().cache.clone();
5079 assert!(!Arc::ptr_eq(&a, &b), "each track owns its account");
5080 }
5081
5082 #[tokio::test]
5085 async fn a_dynamic_defers_teardown() {
5086 let (mut producer, pool) = pooled_producer(1 << 40);
5087 let dynamic = producer.dynamic();
5088 finished_group(&mut producer, 100);
5089
5090 drop(producer);
5091 assert!(pool.used() > 0, "the handler still serves the cache");
5092
5093 drop(dynamic);
5094 assert_eq!(pool.used(), 0, "the last handle tears it down");
5095 }
5096
5097 #[tokio::test]
5103 async fn finished_track_frees_its_cache() {
5104 let (mut producer, pool) = pooled_producer(1 << 40);
5105 finished_group(&mut producer, 100);
5106 producer.finish().unwrap();
5107
5108 let state = producer.state.downgrade();
5109 drop(producer);
5110
5111 assert!(state.upgrade().is_none(), "the track state is freed");
5112 assert_eq!(pool.used(), 0, "so are its cached bytes");
5113 }
5114
5115 #[tokio::test]
5119 async fn teardown_ignores_a_settling_group() {
5120 let (mut producer, pool) = pooled_producer(1 << 40);
5121 finished_group(&mut producer, 100);
5122
5123 let settling = producer.state.downgrade().upgrade().expect("open");
5125 drop(producer);
5126
5127 assert_eq!(pool.used(), 0, "the abrupt teardown still released the cache");
5128 drop(settling);
5129 }
5130
5131 #[tokio::test]
5134 async fn cached_group_outlives_its_track() {
5135 let (mut producer, pool) = pooled_producer(1 << 40);
5136 let sequence = finished_group(&mut producer, 100);
5137 let group = producer.consume().peek_group(sequence).expect("cached");
5138 producer.finish().unwrap();
5139
5140 let state = producer.state.downgrade();
5141 drop(producer);
5142 assert!(state.upgrade().is_none(), "the track state is freed");
5143 assert!(pool.used() > 0, "the retained group keeps its own bytes");
5144
5145 drop(group);
5146 assert_eq!(pool.used(), 0, "which it releases when dropped");
5147 }
5148
5149 #[tokio::test]
5153 async fn pre_accept_backfill_settles_late_writes() {
5154 tokio::time::pause();
5155
5156 let pool = cache::Pool::new(2_000);
5157 let broadcast = broadcast::Info {
5158 origin: crate::origin::Info::default().with_pool(pool.clone()),
5159 ..Default::default()
5160 };
5161 let request = Request::new(Arc::new(broadcast), "test");
5162 let dynamic = request.dynamic();
5163 let consumer = request.consume();
5164
5165 let pending = consumer.fetch_group(0, None);
5167 let req = dynamic
5168 .requested_group()
5169 .now_or_never()
5170 .expect("should not block")
5171 .unwrap();
5172 let mut backfill = req.accept(None).unwrap();
5173 pending.await.unwrap();
5174
5175 let mut producer = request.accept(None);
5177 producer.append_group().unwrap().finish().unwrap();
5178
5179 backfill
5182 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 300_000]))
5183 .unwrap();
5184
5185 assert!(
5186 pool.used() <= 5_000,
5187 "the frame write settled the debt: {}",
5188 pool.used()
5189 );
5190 }
5191
5192 #[tokio::test]
5196 async fn write_restarts_retention_clock() {
5197 tokio::time::pause();
5198
5199 let (mut producer, _pool) = pooled_producer(1 << 40);
5200 let mut straggler = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
5205 straggler
5206 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 100]))
5207 .unwrap();
5208 producer.append_group().unwrap().finish().unwrap(); let consumer = producer.consume();
5211 assert!(consumer.peek_group(0).is_some(), "the write restarted the clock");
5212
5213 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
5215 producer.append_group().unwrap().finish().unwrap(); assert!(consumer.peek_group(0).is_none(), "idle content still expires");
5217 }
5218
5219 #[tokio::test]
5222 async fn refreshed_front_does_not_starve_expiry() {
5223 tokio::time::pause();
5224
5225 let (mut producer, _pool) = pooled_producer(1 << 40);
5226 let dynamic = producer.dynamic();
5227 let consumer = producer.consume();
5228
5229 producer.create_group(10u64.into()).unwrap().finish().unwrap();
5230 for sequence in 1..=5u64 {
5231 let pending = consumer.fetch_group(sequence, None);
5232 let req = dynamic
5233 .requested_group()
5234 .now_or_never()
5235 .expect("should not block")
5236 .unwrap();
5237 let mut group = req.accept(None).unwrap();
5238 group
5239 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 100]))
5240 .unwrap();
5241 group.finish().unwrap();
5242 pending.await.unwrap();
5243 }
5244
5245 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
5248 for sequence in 1..=4u64 {
5249 consumer.fetch_group(sequence, None).await.unwrap();
5250 }
5251
5252 for _ in 0..3 {
5254 producer.append_group().unwrap().finish().unwrap();
5255 }
5256 assert!(consumer.peek_group(5).is_none(), "expired backfill is reclaimed");
5257 assert!(consumer.peek_group(1).is_some(), "refreshed backfill survives");
5258 }
5259
5260 #[tokio::test]
5263 async fn recreated_sequence_delivered_once() {
5264 let (mut producer, _pool) = pooled_producer(1 << 40);
5265
5266 producer.create_group(0u64.into()).unwrap().finish().unwrap();
5267 let aborted = producer.create_group(1u64.into()).unwrap();
5268 aborted.abort(Error::Cancel).unwrap();
5269 producer.create_group(2u64.into()).unwrap().finish().unwrap();
5270 producer.create_group(1u64.into()).unwrap().finish().unwrap();
5271
5272 let mut subscriber = producer.subscribe(None);
5273 assert_eq!(subscriber.assert_group().sequence, 0);
5274 assert_eq!(subscriber.assert_group().sequence, 2);
5275 assert_eq!(
5276 subscriber.assert_group().sequence,
5277 1,
5278 "replacement arrives at its own position"
5279 );
5280 subscriber.assert_no_group();
5281 }
5282
5283 #[tokio::test]
5287 async fn datagrams_do_not_block_eviction() {
5288 tokio::time::pause();
5289
5290 let (mut producer, pool) = pooled_producer(1_000);
5291 for _ in 0..10 {
5292 finished_group(&mut producer, 1_000);
5293 producer.append_datagram(Timestamp::ZERO, &b"beat"[..]).unwrap();
5294 }
5295
5296 let consumer = producer.consume();
5297 assert!(consumer.peek_group(0).is_none(), "old groups still evict");
5298 assert!(
5299 pool.used() < 4 * 1_256,
5300 "interleaved datagrams must not bypass the budget: {}",
5301 pool.used()
5302 );
5303 }
5304
5305 #[tokio::test]
5309 async fn aborted_group_leaves_no_ghost_sample() {
5310 tokio::time::pause();
5311
5312 let (mut producer, pool) = pooled_producer(1 << 40);
5313 let group0 = producer.append_group().unwrap();
5314 producer.append_group().unwrap(); assert!(pool.average().is_some(), "demoted group is sampled");
5317 group0.abort(Error::Cancel).unwrap();
5318 assert_eq!(pool.average(), None, "the abort must remove the sample");
5319 }
5320
5321 #[tokio::test]
5324 async fn empty_groups_repay_overhead() {
5325 tokio::time::pause();
5326
5327 let (mut producer, pool) = pooled_producer(1_000);
5328 for _ in 0..100 {
5329 let mut group = producer.append_group().unwrap();
5330 group.finish().unwrap();
5331 }
5332
5333 assert!(
5334 pool.used() <= 3_000,
5335 "empty-group overhead must stay near the budget: {}",
5336 pool.used()
5337 );
5338 }
5339
5340 #[tokio::test]
5343 async fn growth_on_demoted_group_is_billed() {
5344 tokio::time::pause();
5345
5346 let (mut producer, pool) = pooled_producer(2_000);
5347 let mut straggler = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); straggler
5352 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 10_000]))
5353 .unwrap();
5354
5355 producer.append_group().unwrap().finish().unwrap(); let consumer = producer.consume();
5359 assert!(consumer.peek_group(0).is_none(), "the ballooned group is evicted");
5360 assert!(pool.used() <= 3_000, "growth is reclaimed: {}", pool.used());
5361 }
5362
5363 #[tokio::test]
5366 async fn refilled_sequence_stays_out_of_subscriptions() {
5367 let (mut producer, _pool) = pooled_producer(1 << 40);
5368 let dynamic = producer.dynamic();
5369 let consumer = producer.consume();
5370
5371 producer.create_group(0u64.into()).unwrap().finish().unwrap();
5372 let aborted = producer.create_group(1u64.into()).unwrap();
5373 aborted.abort(Error::Cancel).unwrap();
5374 producer.create_group(2u64.into()).unwrap().finish().unwrap();
5375
5376 let pending = consumer.fetch_group(1, None);
5379 let req = dynamic
5380 .requested_group()
5381 .now_or_never()
5382 .expect("should not block")
5383 .unwrap();
5384 let mut group = req.accept(None).unwrap();
5385 group
5386 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"backfill"))
5387 .unwrap();
5388 group.finish().unwrap();
5389 pending.await.unwrap();
5390
5391 assert!(consumer.peek_group(1).is_some());
5393 let mut subscriber = producer.subscribe(None);
5394 assert_eq!(subscriber.assert_group().sequence, 0);
5395 assert_eq!(subscriber.assert_group().sequence, 2);
5396 subscriber.assert_no_group();
5397 }
5398
5399 #[tokio::test]
5402 async fn expired_backfill_behind_refreshed_reclaimed() {
5403 tokio::time::pause();
5404
5405 let (mut producer, _pool) = pooled_producer(1 << 40);
5406 let dynamic = producer.dynamic();
5407 let consumer = producer.consume();
5408
5409 producer.create_group(5u64.into()).unwrap().finish().unwrap();
5410 for sequence in [2u64, 3u64] {
5411 let pending = consumer.fetch_group(sequence, None);
5412 let req = dynamic
5413 .requested_group()
5414 .now_or_never()
5415 .expect("should not block")
5416 .unwrap();
5417 let mut group = req.accept(None).unwrap();
5418 group
5419 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 100]))
5420 .unwrap();
5421 group.finish().unwrap();
5422 pending.await.unwrap();
5423 }
5424
5425 tokio::time::advance(Duration::from_secs(4)).await;
5427 consumer.fetch_group(2, None).await.unwrap();
5428 tokio::time::advance(DEFAULT_LATENCY_MAX - Duration::from_secs(2)).await;
5429 producer.create_group(6u64.into()).unwrap().finish().unwrap();
5430
5431 let consumer = producer.consume();
5432 assert!(consumer.peek_group(2).is_some(), "refreshed backfill survives");
5433 assert!(consumer.peek_group(3).is_none(), "expired backfill is reclaimed");
5434 }
5435
5436 #[tokio::test]
5439 async fn same_tick_fetch_protects() {
5440 tokio::time::pause();
5441
5442 let (mut producer, _pool) = pooled_producer(10_000);
5444 let consumer = producer.consume();
5445
5446 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();
5451
5452 finished_group(&mut producer, 3_000); finished_group(&mut producer, 3_000); assert!(consumer.peek_group(0).is_some(), "same-tick refresh protects");
5456 assert!(consumer.peek_group(1).is_none(), "the unread group dies instead");
5457 }
5458
5459 #[tokio::test]
5463 async fn refetched_latest_stays_protected() {
5464 tokio::time::pause();
5465
5466 let (mut producer, _pool) = pooled_producer(10_000);
5467 let dynamic = producer.dynamic();
5468 let consumer = producer.consume();
5469
5470 let straggler = producer.append_group().unwrap(); let latest = producer.append_group().unwrap(); latest.abort(Error::Cancel).unwrap();
5475
5476 let pending = consumer.fetch_group(1, None);
5478 let req = dynamic
5479 .requested_group()
5480 .now_or_never()
5481 .expect("should not block")
5482 .unwrap();
5483 let mut group = req.accept(None).unwrap();
5484 group
5485 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 1000]))
5486 .unwrap();
5487 group.finish().unwrap();
5488 pending.await.unwrap();
5489
5490 {
5493 let state = producer.state.read();
5494 assert!(state.lookup.contains_key(&1), "refetched group is cached");
5495 assert!(
5496 state.evict.iter().all(|(sequence, _)| *sequence != 1),
5497 "the live edge must not be an eviction candidate"
5498 );
5499 }
5500 drop(straggler);
5501 }
5502
5503 #[tokio::test]
5506 async fn eviction_allows_refetch() {
5507 tokio::time::pause();
5508
5509 let (mut producer, _pool) = pooled_producer(10_000);
5510 let dynamic = producer.dynamic();
5511
5512 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let consumer = producer.consume();
5517 assert!(consumer.peek_group(0).is_none());
5518 let pending = consumer.fetch_group(0, None);
5519
5520 let req = dynamic
5521 .requested_group()
5522 .now_or_never()
5523 .expect("should not block")
5524 .unwrap();
5525 assert_eq!(req.sequence(), 0);
5526
5527 let mut group = req.accept(None).unwrap();
5528 group
5529 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"refetched"))
5530 .unwrap();
5531 group.finish().unwrap();
5532
5533 let mut group = pending.await.unwrap();
5534 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"refetched");
5535 }
5536
5537 #[tokio::test]
5540 async fn fetched_backfill_not_subscribed() {
5541 let (mut producer, _pool) = pooled_producer(1 << 40);
5542 let dynamic = producer.dynamic();
5543 let consumer = producer.consume();
5544
5545 producer.create_group(5u64.into()).unwrap().finish().unwrap();
5547 producer.create_group(6u64.into()).unwrap().finish().unwrap();
5548
5549 let pending = consumer.fetch_group(2, None);
5551 let req = dynamic
5552 .requested_group()
5553 .now_or_never()
5554 .expect("should not block")
5555 .unwrap();
5556 let mut group = req.accept(None).unwrap();
5557 group
5558 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"backfill"))
5559 .unwrap();
5560 group.finish().unwrap();
5561 let mut fetched = pending.await.unwrap();
5562 assert_eq!(&fetched.read_frame().await.unwrap().unwrap().payload[..], b"backfill");
5563 assert!(consumer.peek_group(2).is_some(), "backfill is cached for later fetches");
5564
5565 let mut subscriber = producer.subscribe(None);
5567 assert_eq!(subscriber.assert_group().sequence, 5);
5568 assert_eq!(subscriber.assert_group().sequence, 6);
5569 subscriber.assert_no_group();
5570 }
5571
5572 #[tokio::test]
5575 async fn expired_backfill_reclaimed() {
5576 tokio::time::pause();
5577
5578 let (mut producer, pool) = pooled_producer(1 << 40);
5579 let dynamic = producer.dynamic();
5580 let consumer = producer.consume();
5581
5582 producer.create_group(5u64.into()).unwrap().finish().unwrap();
5583
5584 let pending = consumer.fetch_group(2, None);
5586 let req = dynamic
5587 .requested_group()
5588 .now_or_never()
5589 .expect("should not block")
5590 .unwrap();
5591 let mut group = req.accept(None).unwrap();
5592 group
5593 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 1000]))
5594 .unwrap();
5595 group.finish().unwrap();
5596 pending.await.unwrap();
5597 let used = pool.used();
5598
5599 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
5601 producer.create_group(6u64.into()).unwrap().finish().unwrap();
5602
5603 assert!(consumer.peek_group(2).is_none(), "expired backfill is reclaimed");
5604 assert!(pool.used() < used, "its bytes are released");
5605 }
5606
5607 #[tokio::test]
5608 async fn fetch_aborts_with_track() {
5609 let producer = track_producer("test", None);
5610 let dynamic = producer.dynamic();
5611 let consumer = producer.consume();
5612
5613 let pending = consumer.fetch_group(3, None);
5614 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
5615
5616 producer.abort(Error::Cancel).unwrap();
5617 assert!(pending.await.is_err());
5618 drop(dynamic);
5619 }
5620}