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::{HashMap, VecDeque},
25 sync::Arc,
26 sync::OnceLock,
27 sync::atomic::{AtomicBool, Ordering},
28 task::{Poll, ready},
29 time::Duration,
30};
31
32pub const DEFAULT_LATENCY_MAX: Duration = Duration::from_secs(5);
34
35const MAX_DATAGRAM_AGE: Duration = Duration::from_millis(50);
41
42const EVICT_SLACK: usize = 64;
45
46const EVICT_SCAN: usize = 4;
50
51#[derive(Clone, Debug)]
62#[non_exhaustive]
63pub struct Info {
64 pub timescale: Timescale,
71 pub latency_max: Duration,
80 pub priority: u8,
83 pub ordered: bool,
87}
88
89impl Default for Info {
90 fn default() -> Self {
91 Self {
92 timescale: Timescale::default(),
93 latency_max: DEFAULT_LATENCY_MAX,
94 priority: 0,
95 ordered: false,
96 }
97 }
98}
99
100impl Info {
101 pub fn with_timescale(mut self, timescale: Timescale) -> Self {
106 self.timescale = timescale;
107 self
108 }
109
110 pub fn with_latency_max(mut self, latency_max: Duration) -> Self {
112 self.latency_max = latency_max;
113 self
114 }
115
116 pub fn with_priority(mut self, priority: u8) -> Self {
118 self.priority = priority;
119 self
120 }
121
122 pub fn with_ordered(mut self, ordered: bool) -> Self {
126 self.ordered = ordered;
127 self
128 }
129}
130
131#[derive(Default)]
132pub(crate) struct TrackState {
133 info: Option<Info>,
136
137 broadcast: Arc<broadcast::Info>,
140
141 cache: Arc<cache::Track>,
145
146 lookup: HashMap<u64, Slot>,
150
151 arrival: VecDeque<(u64, u32)>,
156
157 evict: VecDeque<(u64, u32)>,
165
166 debt: u64,
171
172 datagrams: VecDeque<(Datagram, web_async::time::Instant)>,
176
177 datagram_offset: usize,
180
181 offset: usize,
184
185 max_sequence: Option<u64>,
188
189 latest_group: Option<u64>,
194
195 next_stamp: u32,
197
198 expire_cursor: usize,
201
202 final_sequence: Option<u64>,
204
205 abort: Option<Error>,
207
208 subscriptions: kio::Shared<Subscriptions>,
212
213 fetch: kio::Shared<FetchState>,
216}
217
218struct Slot {
224 group: group::Producer,
225
226 stamp: u32,
231}
232
233type Subscriptions = Vec<kio::Consumer<Subscription>>;
235
236type FetchState = Requests<u64, PendingFetch>;
241
242struct PendingFetch {
244 priority: u8,
246
247 result: kio::Producer<FetchOutcome>,
252}
253
254#[derive(Default)]
257struct FetchOutcome {
258 rejected: Option<Error>,
259}
260
261impl TrackState {
262 fn poll_info(&self) -> Poll<Result<Info>> {
263 if let Some(info) = &self.info {
264 Poll::Ready(Ok(info.clone()))
265 } else {
266 Poll::Pending
267 }
268 }
269
270 fn poll_recv_group(&self, index: usize, min_sequence: u64) -> Poll<Result<Option<(group::Consumer, usize)>>> {
274 let start = index.saturating_sub(self.offset);
275 for (i, (sequence, stamp)) in self.arrival.iter().enumerate().skip(start) {
276 if *sequence >= min_sequence
277 && let Some(slot) = self.lookup.get(sequence)
278 && slot.stamp == *stamp
279 && !slot.group.is_aborted()
280 {
281 return Poll::Ready(Ok(Some((slot.group.consume(), self.offset + i))));
282 }
283 }
284
285 if self.is_complete() {
287 Poll::Ready(Ok(None))
288 } else if let Some(err) = &self.abort {
289 Poll::Ready(Err(err.clone()))
290 } else {
291 Poll::Pending
292 }
293 }
294
295 fn poll_recv_datagram(&self, index: usize) -> Poll<Result<Option<(Datagram, usize)>>> {
301 let start = index.saturating_sub(self.datagram_offset);
302 if let Some((datagram, _)) = self.datagrams.get(start) {
303 return Poll::Ready(Ok(Some((datagram.clone(), self.datagram_offset + start))));
304 }
305
306 if self.is_complete() {
308 Poll::Ready(Ok(None))
309 } else if let Some(err) = &self.abort {
310 Poll::Ready(Err(err.clone()))
311 } else {
312 Poll::Pending
313 }
314 }
315
316 fn push_datagram(&mut self, datagram: Datagram) {
318 let now = web_async::time::Instant::now();
319 self.datagrams.push_back((datagram, now));
320 while let Some((_, at)) = self.datagrams.front() {
321 if now.duration_since(*at) <= MAX_DATAGRAM_AGE {
322 break;
323 }
324 self.datagrams.pop_front();
325 self.datagram_offset += 1;
326 }
327 }
328
329 fn poll_read_frame(
333 &self,
334 index: usize,
335 next_sequence: u64,
336 waiter: &kio::Waiter,
337 ) -> Poll<Result<Option<(frame::Frame, usize, u64)>>> {
338 let start = index.saturating_sub(self.offset);
339 let mut pending_seen = false;
340 for (i, (sequence, stamp)) in self.arrival.iter().enumerate().skip(start) {
341 if *sequence < next_sequence {
342 continue;
343 }
344 let Some(slot) = self.lookup.get(sequence) else {
345 continue;
346 };
347 if slot.stamp != *stamp {
348 continue;
351 }
352
353 let mut consumer = slot.group.consume();
354 match consumer.poll_read_frame(waiter) {
355 Poll::Ready(Ok(Some(frame))) => {
356 return Poll::Ready(Ok(Some((frame, self.offset + i, *sequence))));
357 }
358 Poll::Ready(Ok(None)) => continue,
359 Poll::Ready(Err(_)) => continue,
362 Poll::Pending => {
363 pending_seen = true;
364 continue;
365 }
366 }
367 }
368
369 if pending_seen {
372 Poll::Pending
373 } else if self.is_complete() {
374 Poll::Ready(Ok(None))
375 } else if let Some(err) = &self.abort {
376 Poll::Ready(Err(err.clone()))
377 } else {
378 Poll::Pending
379 }
380 }
381
382 fn poll_next_in_range(
392 &self,
393 next_sequence: u64,
394 end_sequence: Option<u64>,
395 ) -> Poll<Result<Option<group::Consumer>>> {
396 if let Some(end) = end_sequence
400 && end < next_sequence
401 {
402 if let Some(err) = &self.abort {
403 return Poll::Ready(Err(err.clone()));
404 }
405 return Poll::Pending;
406 }
407
408 let mut best: Option<&group::Producer> = None;
409 for slot in self.lookup.values() {
410 let group = &slot.group;
411 if group.sequence < next_sequence {
412 continue;
413 }
414 if let Some(end) = end_sequence
415 && group.sequence > end
416 {
417 continue;
418 }
419 if group.is_aborted() {
420 continue;
421 }
422 if best.is_none_or(|b| group.sequence < b.sequence) {
423 best = Some(group);
424 }
425 }
426
427 if let Some(group) = best {
428 return Poll::Ready(Ok(Some(group.consume())));
429 }
430
431 if let Some(err) = &self.abort {
433 return Poll::Ready(Err(err.clone()));
434 }
435 if let Some(fin) = self.final_sequence
438 && next_sequence >= fin
439 {
440 return Poll::Ready(Ok(None));
441 }
442 Poll::Pending
443 }
444
445 #[cfg(test)]
450 fn cached_group(&self, sequence: u64) -> Option<group::Consumer> {
451 let slot = self.lookup.get(&sequence)?;
452 if slot.group.is_aborted() {
453 return None;
454 }
455 Some(slot.group.consume())
456 }
457
458 fn latency_bound(&self) -> Option<Duration> {
461 self.info.as_ref().map(|info| info.latency_max)
462 }
463
464 fn poll_fetch_cached(&self, sequence: u64) -> Poll<Result<group::Consumer>> {
469 if let Some(slot) = self.lookup.get(&sequence)
470 && !slot.group.is_aborted()
471 {
472 slot.group.cache_refresh();
476 return Poll::Ready(Ok(slot.group.consume()));
477 }
478
479 if let Some(err) = &self.abort {
480 return Poll::Ready(Err(err.clone()));
481 }
482
483 if self.final_sequence.is_some_and(|fin| sequence >= fin) {
485 return Poll::Ready(Err(Error::NotFound));
486 }
487
488 Poll::Pending
489 }
490
491 fn evict_expired(&mut self, max_age: Duration) {
500 let now = self.cache.pool().now();
501 let max_ticks = cache::Pool::ticks(max_age);
502
503 let len = self.evict.len();
504 if len > 0 {
505 let start = self.expire_cursor % len;
506 for step in 0..len.min(EVICT_SCAN) {
507 let (sequence, stamp) = self.evict[(start + step) % len];
508 let Some(slot) = self.lookup.get(&sequence) else {
509 continue;
510 };
511 if slot.stamp != stamp {
512 continue;
514 }
515 if slot.group.is_aborted() {
518 self.lookup.remove(&sequence);
519 continue;
520 }
521 if Some(sequence) == self.latest_group || now.saturating_sub(slot.group.cache_accessed()) <= max_ticks {
522 continue;
523 }
524 let slot = self.lookup.remove(&sequence).unwrap();
528 let _ = slot.group.abort(Error::Old);
529 }
530 self.expire_cursor = (start + EVICT_SCAN) % len;
531 }
532
533 while let Some((sequence, stamp)) = self.arrival.front() {
536 if self.lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp) {
537 break;
538 }
539 self.arrival.pop_front();
540 self.offset += 1;
541 }
542
543 while let Some((sequence, stamp)) = self.evict.front() {
545 if self.lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp) {
546 break;
547 }
548 self.evict.pop_front();
549 }
550
551 if self.evict.len() > 2 * self.lookup.len() + EVICT_SLACK {
554 let lookup = &self.lookup;
555 self.evict
556 .retain(|(sequence, stamp)| lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp));
557 }
558 }
559
560 fn clear_cache(&mut self) {
563 self.lookup.clear();
564 self.arrival.clear();
565 self.evict.clear();
566 self.latest_group = None;
567 self.debt = 0;
568 }
569
570 fn install(&mut self, mut info: Info) {
576 info.latency_max = info.latency_max.min(self.broadcast.origin.cache_duration);
577 self.info = Some(info);
578 }
579
580 fn spawn(broadcast: Arc<broadcast::Info>) -> kio::Producer<Self> {
587 let state = kio::Producer::new(Self {
588 broadcast: broadcast.clone(),
589 ..Default::default()
590 });
591 let cache = cache::Track::new(broadcast.origin.pool.clone(), state.downgrade());
592 state.write().ok().expect("a new track is open").cache = cache;
593 state
594 }
595
596 fn claim_sequence(&mut self, sequence: u64) -> Result<()> {
602 if let Some(slot) = self.lookup.get(&sequence) {
603 if !slot.group.is_aborted() {
604 return Err(Error::Duplicate);
605 }
606 self.lookup.remove(&sequence);
607 }
608 Ok(())
609 }
610
611 fn insert_group(&mut self, group: &group::Producer, visible: bool) {
618 let sequence = group.sequence;
619 self.next_stamp = self.next_stamp.wrapping_add(1);
620 let stamp = self.next_stamp;
621
622 if self.latest_group.is_none_or(|latest| sequence >= latest) {
626 if let Some(latest) = self.latest_group
629 && sequence > latest
630 && let Some(prev) = self.lookup.get(&latest)
631 {
632 prev.group.cache_demote();
633 self.evict.push_back((latest, prev.stamp));
634 }
635 self.latest_group = Some(sequence);
636 } else {
637 group.cache_demote();
638 self.evict.push_back((sequence, stamp));
639 }
640
641 self.max_sequence = Some(self.max_sequence.map_or(sequence, |max| max.max(sequence)));
642 self.lookup.insert(
643 sequence,
644 Slot {
645 group: group.clone(),
646 stamp,
647 },
648 );
649 if visible {
650 self.arrival.push_back((sequence, stamp));
651 }
652 }
653
654 fn commit_group(&mut self, group: &group::Producer, visible: bool, latency_max: Duration) {
658 self.charge_debt();
659 self.insert_group(group, visible);
660 self.evict_expired(latency_max);
661 }
662
663 pub(super) fn charge_debt(&mut self) {
676 let written = self.cache.take_written();
677 let pool = self.cache.pool().clone();
678 match pool.accrue(written) {
679 Some(mut accrued) => {
680 if self.oldest_is_stale(&pool) {
681 accrued = accrued.saturating_mul(2);
682 }
683 self.debt = self.debt.saturating_add(accrued).min(pool.used());
686 self.pay_debt(&pool, written.saturating_mul(2));
689 }
690 None => self.debt = 0,
693 }
694 }
695
696 fn oldest_is_stale(&self, pool: &cache::Pool) -> bool {
700 let Some(average) = pool.average() else {
701 return false;
702 };
703 let Some((sequence, stamp)) = self.evict.front() else {
704 return false;
705 };
706 let Some(slot) = self.lookup.get(sequence) else {
707 return false;
708 };
709 slot.stamp == *stamp && !slot.group.is_aborted() && slot.group.cache_accessed() <= average
710 }
711
712 fn pay_debt(&mut self, pool: &cache::Pool, cap: u64) {
725 let average = pool.average().unwrap_or(0);
726 let mut paid = 0u64;
727 let mut scanned = 0usize;
728 for _ in 0..self.evict.len() {
729 if self.debt == 0 || paid >= cap || scanned >= EVICT_SCAN {
730 return;
731 }
732 let Some((sequence, stamp)) = self.evict.pop_front() else {
733 return;
734 };
735 let Some(slot) = self.lookup.get(&sequence) else {
736 continue;
738 };
739 if slot.stamp != stamp {
740 continue;
742 }
743 if slot.group.is_aborted() {
744 self.lookup.remove(&sequence);
746 continue;
747 }
748 if Some(sequence) == self.latest_group {
749 self.evict.push_back((sequence, stamp));
751 continue;
752 }
753
754 scanned += 1;
755 if slot.group.cache_accessed() > average {
759 self.evict.push_back((sequence, stamp));
760 continue;
761 }
762 let size = slot.group.cache_size();
765 if size > self.debt {
766 self.evict.push_front((sequence, stamp));
767 return;
768 }
769
770 self.debt -= size;
771 paid = paid.saturating_add(size);
772 let slot = self.lookup.remove(&sequence).unwrap();
773 let _ = slot.group.abort(Error::Evicted);
774 }
775 }
776
777 fn set_final(&mut self, final_sequence: u64) -> Result<()> {
780 if self.final_sequence.is_some() {
781 return Err(Error::Closed);
782 }
783 if let Some(max) = self.max_sequence
784 && final_sequence <= max
785 {
786 return Err(Error::ProtocolViolation);
787 }
788 self.final_sequence = Some(final_sequence);
789 Ok(())
790 }
791
792 fn is_complete(&self) -> bool {
798 self.final_sequence
799 .is_some_and(|fin| self.max_sequence.map_or(0, |max| max.saturating_add(1)) >= fin)
800 }
801
802 fn poll_finished(&self) -> Poll<Result<u64>> {
803 if let Some(fin) = self.final_sequence {
804 Poll::Ready(Ok(fin))
805 } else if let Some(err) = &self.abort {
806 Poll::Ready(Err(err.clone()))
807 } else {
808 Poll::Pending
809 }
810 }
811
812 fn modify(producer: &kio::Producer<Self>) -> Result<kio::Mut<'_, Self>> {
813 producer.write().map_err(|r| r.abort.clone().unwrap_or(Error::Dropped))
814 }
815
816 fn insert_group_request(&mut self, sequence: u64, info: Option<Info>) -> Result<group::Producer> {
822 if let Some(err) = &self.abort {
823 return Err(err.clone());
824 }
825 if let Some(fin) = self.final_sequence
826 && sequence >= fin
827 {
828 return Err(Error::Closed);
829 }
830
831 if self.info.is_none() {
835 self.install(info.unwrap_or_default());
836 }
837 let info = self.info.clone().unwrap();
838
839 self.claim_sequence(sequence)?;
841
842 let latency_max = info.latency_max;
843 let group = group::Producer::new(group::Info { sequence }, info, self.cache.clone());
844 group.cache_refresh();
849 self.commit_group(&group, false, latency_max);
850 Ok(group)
851 }
852}
853
854#[derive(Clone)]
856pub struct Producer {
857 name: Arc<str>,
858 broadcast: Arc<broadcast::Info>,
861 state: kio::Producer<TrackState>,
862 prev_subscription: Option<Subscription>,
863 alive: Arc<Alive>,
865 stats: stats::Scope,
869}
870
871impl Producer {
872 pub(crate) fn new(
880 broadcast: Arc<broadcast::Info>,
881 name: impl Into<Arc<str>>,
882 info: impl Into<Option<Info>>,
883 ) -> Self {
884 let name = name.into();
885 let state = TrackState::spawn(broadcast.clone());
886 state
887 .write()
888 .ok()
889 .expect("a new track is open")
890 .install(info.into().unwrap_or_default());
891 let alive = Alive::new(name.clone(), state.clone());
892 alive.publish(None);
893 Self {
894 name,
895 state,
896 broadcast,
897 prev_subscription: None,
898 alive,
899 stats: stats::Scope::default(),
900 }
901 }
902
903 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
907 self.alive.publish(Some(&scope));
908 self.stats = scope;
909 self
910 }
911
912 pub fn name(&self) -> &str {
914 &self.name
915 }
916
917 pub fn broadcast(&self) -> &broadcast::Info {
919 &self.broadcast
920 }
921
922 pub fn create_group(&mut self, group: group::Info) -> Result<group::Producer> {
924 let mut state = self.modify()?;
925 if let Some(fin) = state.final_sequence
926 && group.sequence >= fin
927 {
928 return Err(Error::Closed);
929 }
930 let track = state.info.clone().unwrap();
931 let latency_max = track.latency_max;
932
933 state.claim_sequence(group.sequence)?;
935
936 let group = group::Producer::new(group, track, state.cache.clone()).with_meter(self.stats.meter());
937 state.commit_group(&group, true, latency_max);
938
939 Ok(group)
940 }
941
942 pub fn append_group(&mut self) -> Result<group::Producer> {
944 let mut state = self.modify()?;
945 let sequence = match state.max_sequence {
946 Some(s) => s.checked_add(1).ok_or(coding::BoundsExceeded)?,
947 None => 0,
948 };
949 if let Some(fin) = state.final_sequence
950 && sequence >= fin
951 {
952 return Err(Error::Closed);
953 }
954
955 let track = state.info.clone().unwrap();
956 let latency_max = track.latency_max;
957
958 let group =
959 group::Producer::new(group::Info { sequence }, track, state.cache.clone()).with_meter(self.stats.meter());
960 state.commit_group(&group, true, latency_max);
961
962 Ok(group)
963 }
964
965 pub fn append_datagram<B: crate::IntoBytes>(&mut self, timestamp: Timestamp, payload: B) -> Result<u64> {
977 let payload = payload.into_bytes();
978 if payload.len() > super::datagram::MAX_DATAGRAM_PAYLOAD {
979 return Err(Error::FrameTooLarge);
980 }
981 let meter = self.stats.meter();
983 let mut state = self.modify()?;
984 let timescale = state.info.as_ref().unwrap().timescale;
986 let timestamp = timestamp.convert(timescale).map_err(|_| Error::TimestampMismatch)?;
987 let sequence = match state.max_sequence {
988 Some(s) => s.checked_add(1).ok_or(coding::BoundsExceeded)?,
989 None => 0,
990 };
991 if let Some(fin) = state.final_sequence
992 && sequence >= fin
993 {
994 return Err(Error::Closed);
995 }
996 state.max_sequence = Some(sequence);
997 meter.datagram(payload.len() as u64);
998 state.push_datagram(Datagram {
999 sequence,
1000 timestamp,
1001 payload,
1002 });
1003 Ok(sequence)
1004 }
1005
1006 pub fn write_datagram(&mut self, mut datagram: Datagram) -> Result<()> {
1012 if datagram.payload.len() > super::datagram::MAX_DATAGRAM_PAYLOAD {
1013 return Err(Error::FrameTooLarge);
1014 }
1015 let meter = self.stats.meter();
1017 let mut state = self.modify()?;
1018 let timescale = state.info.as_ref().unwrap().timescale;
1020 datagram.timestamp = datagram
1021 .timestamp
1022 .convert(timescale)
1023 .map_err(|_| Error::TimestampMismatch)?;
1024 if let Some(fin) = state.final_sequence
1025 && datagram.sequence >= fin
1026 {
1027 return Err(Error::Closed);
1028 }
1029 state.max_sequence = Some(state.max_sequence.unwrap_or(0).max(datagram.sequence));
1030 meter.datagram(datagram.payload.len() as u64);
1031 state.push_datagram(datagram);
1032 Ok(())
1033 }
1034
1035 pub fn write_frame<B: crate::IntoBytes>(&mut self, timestamp: Timestamp, frame: B) -> Result<()> {
1040 let mut group = self.append_group()?;
1041 group.write_frame(timestamp, frame)?;
1042 group.finish()?;
1043 Ok(())
1044 }
1045
1046 pub fn finish(&mut self) -> Result<()> {
1052 let mut state = self.modify()?;
1053 let final_sequence = match state.max_sequence {
1054 Some(max) => max.checked_add(1).ok_or(coding::BoundsExceeded)?,
1055 None => 0,
1056 };
1057 state.set_final(final_sequence)
1058 }
1059
1060 pub fn finish_at(&mut self, final_sequence: u64) -> Result<()> {
1073 self.modify()?.set_final(final_sequence)
1074 }
1075
1076 pub fn final_sequence(&self) -> Option<u64> {
1081 self.state.read().final_sequence
1082 }
1083
1084 pub fn abort(self, err: Error) -> Result<()> {
1095 let mut guard = self.modify()?;
1096 guard.abort = Some(err);
1097 guard.clear_cache();
1098 guard.datagrams.clear();
1099 guard.close();
1100 Ok(())
1101 }
1102
1103 pub async fn unused(&self) -> Result<()> {
1105 self.state.unused().await.map_err(|_| self.abort_reason())
1106 }
1107
1108 pub async fn used(&self) -> Result<()> {
1110 self.state.used().await.map_err(|_| self.abort_reason())
1111 }
1112
1113 pub async fn closed(&self) -> Error {
1115 kio::wait(|waiter| self.poll_closed(waiter)).await
1116 }
1117
1118 pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<Error> {
1120 self.state.poll_closed(waiter).map(|()| self.abort_reason())
1121 }
1122
1123 fn abort_reason(&self) -> Error {
1125 self.state.read().abort.clone().unwrap_or(Error::Dropped)
1126 }
1127
1128 pub fn is_closed(&self) -> bool {
1130 self.state.read().is_closed()
1131 }
1132
1133 pub fn latest(&self) -> Option<u64> {
1135 self.state.read().max_sequence
1136 }
1137
1138 pub fn is_clone(&self, other: &Self) -> bool {
1140 self.state.same_channel(&other.state)
1141 }
1142
1143 pub(crate) fn weak(&self) -> TrackWeak {
1145 TrackWeak {
1146 name: self.name.clone(),
1147 state: self.state.weak(),
1148 }
1149 }
1150
1151 pub fn demand(&self) -> Demand {
1159 Demand {
1160 name: self.name.clone(),
1161 state: self.state.weak(),
1162 }
1163 }
1164
1165 pub fn consume(&self) -> Consumer {
1170 Consumer::plain(self.name.clone(), self.state.consume())
1171 }
1172
1173 pub fn subscribe(&self, subscription: impl Into<Option<Subscription>>) -> Subscriber {
1178 let preferences = subscription.into().unwrap_or_default();
1179
1180 let info = self.state.read().info.clone().expect("producer always has info");
1185 let subscription = kio::Producer::new(preferences);
1186 register_subscription(self.state.read(), &subscription);
1187
1188 Subscriber {
1189 name: self.name.clone(),
1190 info,
1191 inner: SubscriberKind::Plain(PlainSubscriber {
1192 state: self.state.consume(),
1193 subscription,
1194 index: 0,
1195 datagram_index: 0,
1196 min_sequence: 0,
1197 next_sequence: 0,
1198 end_sequence: None,
1199 }),
1200 stats: stats::Scope::default(),
1202 _stats_sub: stats::Subscription::default(),
1203 }
1204 }
1205
1206 pub async fn subscription_changed(&mut self) -> Result<Option<Subscription>> {
1212 kio::wait(|waiter| self.poll_subscription_changed(waiter)).await
1213 }
1214
1215 pub fn subscription(&self) -> Option<Subscription> {
1223 let state = self.state.read();
1224 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
1225 drop(state);
1226 snapshot_subscription(&subs, bound)
1227 }
1228
1229 pub fn poll_subscription_changed(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Subscription>>> {
1233 if self.state.poll_closed(waiter).is_ready() {
1236 let abort = self.state.read().abort.clone();
1237 return Poll::Ready(Err(abort.unwrap_or(Error::Dropped)));
1238 }
1239
1240 let state = self.state.read();
1242 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
1243 drop(state);
1244
1245 let prev = &self.prev_subscription;
1246 let mut combined = None;
1247 let mut guard = match subs.poll(waiter, |subs| {
1248 let next = combined_subscription(subs, bound, waiter);
1249 if &next == prev {
1250 Poll::Pending
1251 } else {
1252 combined = next;
1253 Poll::Ready(())
1254 }
1255 }) {
1256 Poll::Ready(guard) => guard,
1257 Poll::Pending => return Poll::Pending,
1258 };
1259 guard.retain(|sub| !sub.is_closed());
1261 drop(guard);
1262 self.prev_subscription = combined.clone();
1263 Poll::Ready(Ok(combined))
1264 }
1265
1266 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
1268 self.state.poll_unused(waiter).map(|_| ())
1269 }
1270
1271 pub fn dynamic(&self) -> Dynamic {
1275 Dynamic::new(self.name.clone(), self.state.clone(), self.alive.clone())
1276 }
1277
1278 fn modify(&self) -> Result<kio::Mut<'_, TrackState>> {
1279 TrackState::modify(&self.state)
1280 }
1281}
1282
1283fn poll_requested_group(
1287 state: &kio::Producer<TrackState>,
1288 fetch: &kio::Shared<FetchState>,
1289 waiter: &kio::Waiter,
1290) -> Poll<Result<GroupRequest>> {
1291 if let Poll::Ready(mut guard) = fetch.poll(waiter, |fetch| {
1293 if fetch.has_queued() {
1294 Poll::Ready(())
1295 } else {
1296 Poll::Pending
1297 }
1298 }) {
1299 let sequence = guard.pop().expect("predicate guaranteed a request");
1300 let pending = guard.get(&sequence).expect("popped key must be pending");
1304 let priority = pending.priority;
1305 let result = pending.result.clone();
1306 drop(guard);
1307 return Poll::Ready(Ok(GroupRequest {
1308 state: state.clone(),
1309 fetch: fetch.clone(),
1310 sequence,
1311 priority,
1312 result,
1313 done: false,
1314 }));
1315 }
1316
1317 match state.poll_ref(waiter, |state| match &state.abort {
1319 Some(err) => Poll::Ready(err.clone()),
1320 None => Poll::Pending,
1321 }) {
1322 Poll::Ready(Ok(err)) => Poll::Ready(Err(err)),
1323 Poll::Ready(Err(closed)) => Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped))),
1324 Poll::Pending => Poll::Pending,
1325 }
1326}
1327
1328pub struct Dynamic {
1338 name: Arc<str>,
1339 state: kio::Producer<TrackState>,
1341 fetch: kio::Shared<FetchState>,
1343 alive: Arc<Alive>,
1346}
1347
1348impl Dynamic {
1349 fn new(name: Arc<str>, state: kio::Producer<TrackState>, alive: Arc<Alive>) -> Self {
1350 let fetch = state.read().fetch.clone();
1351 fetch.lock().add_handler();
1352 Self {
1353 name,
1354 state,
1355 fetch,
1356 alive,
1357 }
1358 }
1359
1360 pub fn name(&self) -> &str {
1362 &self.name
1363 }
1364
1365 pub async fn requested_group(&self) -> Result<GroupRequest> {
1371 kio::wait(|waiter| self.poll_requested_group(waiter)).await
1372 }
1373
1374 pub fn poll_requested_group(&self, waiter: &kio::Waiter) -> Poll<Result<GroupRequest>> {
1376 poll_requested_group(&self.state, &self.fetch, waiter)
1377 }
1378
1379 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
1381 self.state.poll_unused(waiter).map(|_| ())
1382 }
1383}
1384
1385impl Clone for Dynamic {
1386 fn clone(&self) -> Self {
1387 self.fetch.lock().add_handler();
1389 Self {
1390 name: self.name.clone(),
1391 state: self.state.clone(),
1392 fetch: self.fetch.clone(),
1393 alive: self.alive.clone(),
1394 }
1395 }
1396}
1397
1398impl Drop for Dynamic {
1399 fn drop(&mut self) {
1400 let mut fetch = self.fetch.lock();
1406 if fetch.remove_handler() {
1407 fetch.drain_queued();
1408 }
1409 }
1410}
1411
1412struct Alive {
1421 name: Arc<str>,
1422 state: kio::Producer<TrackState>,
1423
1424 published: AtomicBool,
1427
1428 stats: OnceLock<stats::Subscription>,
1431}
1432
1433impl Alive {
1434 fn new(name: Arc<str>, state: kio::Producer<TrackState>) -> Arc<Self> {
1435 Arc::new(Self {
1436 name,
1437 state,
1438 published: Default::default(),
1439 stats: Default::default(),
1440 })
1441 }
1442
1443 fn publish(&self, stats: Option<&stats::Scope>) {
1447 self.published.store(true, Ordering::Relaxed);
1448 if let Some(scope) = stats {
1449 let _ = self.stats.set(scope.subscribe());
1452 }
1453 }
1454}
1455
1456impl Drop for Alive {
1457 fn drop(&mut self) {
1458 if !self.published.load(Ordering::Relaxed) {
1460 return;
1461 }
1462 if let Ok(mut state) = self.state.write()
1467 && state.final_sequence.is_none()
1468 {
1469 tracing::warn!(
1473 track = %self.name,
1474 "track::Producer dropped without finish() or abort()"
1475 );
1476 state.clear_cache();
1477 state.datagrams.clear();
1478 }
1479 }
1480}
1481
1482fn combined_subscription(subs: &Subscriptions, bound: Option<Duration>, waiter: &kio::Waiter) -> Option<Subscription> {
1488 let mut combined = None;
1489 for sub in subs.iter() {
1490 if sub.is_closed() {
1495 continue;
1496 }
1497 let _ = sub.poll_closed(waiter);
1502 if let Poll::Ready(Ok(sub)) = sub.poll(waiter, |sub| sub.poll_combined(&combined)) {
1503 combined = Some(sub);
1504 }
1505 }
1506 clamp_combined(combined, bound)
1507}
1508
1509fn snapshot_subscription(subs: &kio::Shared<Subscriptions>, bound: Option<Duration>) -> Option<Subscription> {
1511 let mut combined: Option<Subscription> = None;
1512 for sub in subs.read().iter() {
1513 if sub.is_closed() {
1515 continue;
1516 }
1517 if let Poll::Ready(merged) = sub.read().poll_combined(&combined) {
1518 combined = Some(merged);
1519 }
1520 }
1521 clamp_combined(combined, bound)
1522}
1523
1524fn clamp_combined(combined: Option<Subscription>, bound: Option<Duration>) -> Option<Subscription> {
1532 let mut combined = combined?;
1533 if let Some(bound) = bound {
1534 combined.latency_max = combined.latency_max.min(bound);
1535 }
1536 Some(combined)
1537}
1538
1539fn register_subscription(state: kio::Ref<'_, TrackState>, subscription: &kio::Producer<Subscription>) {
1543 if state.is_closed() {
1544 return;
1545 }
1546 let subs = state.subscriptions.clone();
1547 drop(state);
1548 subs.lock().push(subscription.consume());
1549}
1550
1551#[derive(Clone)]
1553pub(crate) struct TrackWeak {
1554 name: Arc<str>,
1555 state: kio::ProducerWeak<TrackState>,
1556}
1557
1558impl TrackWeak {
1559 pub fn consume(&self) -> Consumer {
1560 Consumer::plain(self.name.clone(), self.state.consume())
1561 }
1562
1563 pub(crate) fn name(&self) -> &Arc<str> {
1566 &self.name
1567 }
1568
1569 pub(crate) fn is_used(&self) -> bool {
1572 !self.state.is_closed() && self.state.is_used()
1573 }
1574
1575 pub(crate) fn poll_used(&self, waiter: &kio::Waiter) {
1578 let _ = self.state.poll_used(waiter);
1579 }
1580
1581 pub(crate) fn poll_unused(&self, waiter: &kio::Waiter) {
1584 let _ = self.state.poll_unused(waiter);
1585 }
1586}
1587
1588impl super::WeakEntry for TrackWeak {
1589 fn is_closed(&self) -> bool {
1590 self.state.is_closed()
1591 }
1592
1593 fn same_channel(&self, other: &Self) -> bool {
1594 self.state.same_channel(&other.state)
1595 }
1596}
1597
1598#[derive(Clone)]
1607pub struct Demand {
1608 name: Arc<str>,
1609 state: kio::ProducerWeak<TrackState>,
1610}
1611
1612impl Demand {
1613 pub fn name(&self) -> &str {
1615 &self.name
1616 }
1617
1618 pub async fn used(&self) -> Result<()> {
1620 self.state.used().await.map_err(|_| self.abort_reason())
1621 }
1622
1623 pub async fn unused(&self) -> Result<()> {
1625 self.state.unused().await.map_err(|_| self.abort_reason())
1626 }
1627
1628 pub async fn closed(&self) -> Error {
1630 self.state.closed().await;
1631 self.abort_reason()
1632 }
1633
1634 fn abort_reason(&self) -> Error {
1636 self.state.read().abort.clone().unwrap_or(Error::Dropped)
1637 }
1638}
1639
1640#[derive(Clone)]
1651pub struct Consumer {
1652 name: Arc<str>,
1653 inner: ConsumerKind,
1654 stats: stats::Scope,
1657}
1658
1659#[derive(Clone)]
1660enum ConsumerKind {
1661 Plain(kio::Consumer<TrackState>),
1662 Spliced(super::resume::Consumer),
1663}
1664
1665impl Consumer {
1666 fn plain(name: Arc<str>, state: kio::Consumer<TrackState>) -> Self {
1667 Self {
1668 name,
1669 inner: ConsumerKind::Plain(state),
1670 stats: stats::Scope::default(),
1671 }
1672 }
1673
1674 pub(crate) fn spliced(name: Arc<str>, resume: super::resume::Consumer) -> Self {
1676 Self {
1677 name,
1678 inner: ConsumerKind::Spliced(resume),
1679 stats: stats::Scope::default(),
1680 }
1681 }
1682
1683 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
1686 self.stats = scope;
1687 self
1688 }
1689
1690 pub fn name(&self) -> &str {
1692 &self.name
1693 }
1694
1695 pub fn subscribe(&self, subscription: impl Into<Option<Subscription>>) -> kio::Pending<Subscribing> {
1701 let subscription = kio::Producer::new(subscription.into().unwrap_or_default());
1702
1703 let inner = match &self.inner {
1704 ConsumerKind::Plain(state) => {
1705 register_subscription(state.read(), &subscription);
1708 SubscribingKind::Plain(state.clone())
1709 }
1710 ConsumerKind::Spliced(resume) => SubscribingKind::Spliced(resume.clone()),
1712 };
1713
1714 kio::Pending::new(Subscribing {
1715 name: self.name.clone(),
1716 inner,
1717 subscription,
1718 stats: self.stats.clone(),
1719 })
1720 }
1721
1722 #[cfg(test)]
1726 pub(crate) fn peek_group(&self, sequence: u64) -> Option<group::Consumer> {
1727 match &self.inner {
1728 ConsumerKind::Plain(state) => state.read().cached_group(sequence),
1729 ConsumerKind::Spliced(_) => None,
1732 }
1733 }
1734
1735 pub fn fetch_group(&self, sequence: u64, options: impl Into<Option<group::Fetch>>) -> kio::Pending<Fetching> {
1747 let options = options.into().unwrap_or_default();
1748
1749 self.stats.fetch();
1753
1754 let state = match &self.inner {
1755 ConsumerKind::Plain(state) => state,
1756 ConsumerKind::Spliced(resume) => {
1759 return kio::Pending::new(Fetching {
1760 inner: FetchingKind::Spliced(resume.fetch_group(sequence, options)),
1761 stats: self.stats.clone(),
1762 });
1763 }
1764 };
1765
1766 let mut result = None;
1767
1768 let (fetch, unresolved) = {
1772 let state = state.read();
1773 (state.fetch.clone(), state.poll_fetch_cached(sequence).is_pending())
1774 };
1775
1776 if unresolved {
1777 let mut fetch = fetch.lock();
1778 if let Some(pending) = fetch.join(&sequence) {
1779 pending.priority = pending.priority.max(options.priority);
1782 result = Some(pending.result.consume());
1783 } else {
1784 let producer = kio::Producer::<FetchOutcome>::default();
1788 let consumer = producer.consume();
1789 let attempt = PendingFetch {
1790 priority: options.priority,
1791 result: producer,
1792 };
1793 if fetch.insert(sequence, attempt).is_ok() {
1794 result = Some(consumer);
1795 }
1796 }
1797 }
1798
1799 kio::Pending::new(Fetching {
1800 inner: FetchingKind::Plain {
1801 state: state.clone(),
1802 fetch,
1803 sequence,
1804 result,
1805 },
1806 stats: self.stats.clone(),
1807 })
1808 }
1809
1810 pub fn info(&self) -> kio::Pending<Querying> {
1817 kio::Pending::new(Querying {
1818 inner: match &self.inner {
1819 ConsumerKind::Plain(state) => QueryingKind::Plain(state.clone()),
1820 ConsumerKind::Spliced(resume) => QueryingKind::Spliced(resume.clone()),
1821 },
1822 })
1823 }
1824
1825 pub fn latest(&self) -> Option<u64> {
1827 match &self.inner {
1828 ConsumerKind::Plain(state) => state.read().max_sequence,
1829 ConsumerKind::Spliced(resume) => resume.latest(),
1830 }
1831 }
1832
1833 pub(crate) fn poll_complete(&self, waiter: &kio::Waiter) -> Poll<Result<()>> {
1838 let ConsumerKind::Plain(state) = &self.inner else {
1839 return Poll::Pending;
1841 };
1842 match ready!(state.poll(waiter, |state| {
1843 if state.is_complete() {
1844 Poll::Ready(())
1845 } else {
1846 Poll::Pending
1847 }
1848 })) {
1849 Ok(_) => Poll::Ready(Ok(())),
1850 Err(closed) => Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped))),
1853 }
1854 }
1855}
1856
1857pub struct Subscribing {
1860 name: Arc<str>,
1861 inner: SubscribingKind,
1862 subscription: kio::Producer<Subscription>,
1863 stats: stats::Scope,
1864}
1865
1866enum SubscribingKind {
1867 Plain(kio::Consumer<TrackState>),
1868 Spliced(super::resume::Consumer),
1869}
1870
1871impl Subscribing {
1872 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<Subscriber>> {
1875 match &self.inner {
1876 SubscribingKind::Plain(state) => {
1877 let info = ready!(state.poll(waiter, |state| state.poll_info()))
1879 .map_err(|e| e.abort.clone().unwrap_or(Error::Dropped))??;
1880
1881 Poll::Ready(Ok(Subscriber {
1882 name: self.name.clone(),
1883 info,
1884 inner: SubscriberKind::Plain(PlainSubscriber {
1885 state: state.clone(),
1886 subscription: self.subscription.clone(),
1887 index: 0,
1888 datagram_index: 0,
1889 min_sequence: 0,
1890 next_sequence: 0,
1891 end_sequence: None,
1892 }),
1893 stats: self.stats.clone(),
1894 _stats_sub: self.stats.subscribe(),
1895 }))
1896 }
1897 SubscribingKind::Spliced(resume) => {
1898 let info = ready!(resume.poll_info(waiter))?;
1901
1902 Poll::Ready(Ok(Subscriber {
1903 name: self.name.clone(),
1904 info,
1905 inner: SubscriberKind::Spliced(Box::new(resume.subscribe_shared(self.subscription.clone()))),
1906 stats: self.stats.clone(),
1907 _stats_sub: self.stats.subscribe(),
1908 }))
1909 }
1910 }
1911 }
1912
1913 pub fn update(&mut self, subscription: Subscription) -> Result<()> {
1918 let mut state = self.subscription.write().map_err(|_| Error::Closed)?;
1919 *state = subscription;
1920 Ok(())
1921 }
1922}
1923
1924impl kio::Pollable for Subscribing {
1925 type Output = Result<Subscriber>;
1926
1927 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
1928 self.poll_ok(waiter)
1929 }
1930}
1931
1932pub struct Querying {
1935 inner: QueryingKind,
1936}
1937
1938enum QueryingKind {
1939 Plain(kio::Consumer<TrackState>),
1940 Spliced(super::resume::Consumer),
1941}
1942
1943impl Querying {
1944 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<Info>> {
1946 match &self.inner {
1947 QueryingKind::Plain(state) => {
1948 let info = ready!(state.poll(waiter, |state| state.poll_info()))
1950 .map_err(|e| e.abort.clone().unwrap_or(Error::Dropped))??;
1951 Poll::Ready(Ok(info))
1952 }
1953 QueryingKind::Spliced(resume) => resume.poll_info(waiter),
1954 }
1955 }
1956}
1957
1958impl kio::Pollable for Querying {
1959 type Output = Result<Info>;
1960
1961 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
1962 self.poll_ok(waiter)
1963 }
1964}
1965
1966pub struct GroupRequest {
1975 state: kio::Producer<TrackState>,
1976 fetch: kio::Shared<FetchState>,
1978 sequence: u64,
1979 priority: u8,
1980 result: kio::Producer<FetchOutcome>,
1982 done: bool,
1983}
1984
1985impl GroupRequest {
1986 pub fn sequence(&self) -> u64 {
1988 self.sequence
1989 }
1990
1991 pub fn priority(&self) -> u8 {
1993 self.priority
1994 }
1995
1996 pub fn accept(mut self, info: impl Into<Option<Info>>) -> Result<group::Producer> {
2004 self.done = true;
2005 let res = TrackState::modify(&self.state)
2009 .and_then(|mut state| state.insert_group_request(self.sequence, info.into()));
2010 self.remove();
2011 res
2012 }
2013
2014 pub fn reject(mut self, err: Error) {
2016 self.done = true;
2017 self.remove();
2020 if let Ok(mut outcome) = self.result.write() {
2021 outcome.rejected = Some(err);
2022 }
2023 }
2024
2025 fn remove(&self) {
2028 self.fetch
2029 .lock()
2030 .remove_if(&self.sequence, |pending| pending.result.same_channel(&self.result));
2031 }
2032}
2033
2034impl Drop for GroupRequest {
2035 fn drop(&mut self) {
2036 if self.done {
2037 return;
2038 }
2039 self.remove();
2040 if let Ok(mut outcome) = self.result.write() {
2041 outcome.rejected = Some(Error::Dropped);
2042 }
2043 }
2044}
2045
2046pub struct Fetching {
2052 inner: FetchingKind,
2053 stats: stats::Scope,
2056}
2057
2058enum FetchingKind {
2059 Plain {
2060 state: kio::Consumer<TrackState>,
2061 fetch: kio::Shared<FetchState>,
2062 sequence: u64,
2063 result: Option<kio::Consumer<FetchOutcome>>,
2065 },
2066 Spliced(kio::Pending<super::resume::Fetching>),
2068}
2069
2070impl kio::Pollable for Fetching {
2071 type Output = Result<group::Consumer>;
2072
2073 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
2074 let (state, fetch, sequence, result) = match &self.inner {
2075 FetchingKind::Plain {
2076 state,
2077 fetch,
2078 sequence,
2079 result,
2080 } => (state, fetch, *sequence, result.as_ref()),
2081 FetchingKind::Spliced(spliced) => {
2082 return kio::Pollable::poll(&**spliced, waiter)
2085 .map(|res| res.map(|group| group.with_meter(self.stats.meter())));
2086 }
2087 };
2088
2089 match state.poll(waiter, |state| state.poll_fetch_cached(sequence)) {
2092 Poll::Ready(Ok(res)) => return Poll::Ready(res.map(|group| group.with_meter(self.stats.meter()))),
2093 Poll::Ready(Err(closed)) => {
2094 return Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped)));
2095 }
2096 Poll::Pending => {}
2097 }
2098
2099 let Some(result) = result else {
2101 return match fetch.poll(waiter, |fetch| match fetch.has_handlers() {
2104 false => Poll::Ready(()),
2105 true => Poll::Pending,
2106 }) {
2107 Poll::Ready(_guard) => Poll::Ready(Err(Error::NotFound)),
2108 Poll::Pending => Poll::Pending,
2109 };
2110 };
2111
2112 match result.poll(waiter, |outcome| match &outcome.rejected {
2115 Some(err) => Poll::Ready(err.clone()),
2116 None => Poll::Pending,
2117 }) {
2118 Poll::Ready(Ok(err)) => Poll::Ready(Err(err)),
2119 Poll::Ready(Err(_closed)) => Poll::Ready(Err(Error::NotFound)),
2120 Poll::Pending => Poll::Pending,
2121 }
2122 }
2123}
2124
2125pub struct Subscriber {
2148 name: Arc<str>,
2149 info: Info,
2150 inner: SubscriberKind,
2151 stats: stats::Scope,
2154 _stats_sub: stats::Subscription,
2157}
2158
2159enum SubscriberKind {
2160 Plain(PlainSubscriber),
2161 Spliced(Box<super::resume::Subscriber>),
2163}
2164
2165struct PlainSubscriber {
2167 state: kio::Consumer<TrackState>,
2168
2169 subscription: kio::Producer<Subscription>,
2170 index: usize,
2172 datagram_index: usize,
2174 min_sequence: u64,
2176 next_sequence: u64,
2179 end_sequence: Option<u64>,
2184}
2185
2186impl PlainSubscriber {
2187 fn poll<F, R>(&self, waiter: &kio::Waiter, f: F) -> Poll<Result<R>>
2189 where
2190 F: Fn(&kio::Ref<'_, TrackState>) -> Poll<Result<R>>,
2191 {
2192 Poll::Ready(match ready!(self.state.poll(waiter, f)) {
2193 Ok(res) => res,
2194 Err(state) => Err(state.abort.clone().unwrap_or(Error::Dropped)),
2196 })
2197 }
2198
2199 fn poll_recv_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2200 let Some((consumer, found_index)) =
2201 ready!(self.poll(waiter, |state| state.poll_recv_group(self.index, self.min_sequence))?)
2202 else {
2203 return Poll::Ready(Ok(None));
2204 };
2205
2206 self.index = found_index + 1;
2207 Poll::Ready(Ok(Some(consumer)))
2208 }
2209
2210 fn poll_recv_datagram(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Datagram>>> {
2211 let Some((datagram, found_index)) =
2212 ready!(self.poll(waiter, |state| state.poll_recv_datagram(self.datagram_index))?)
2213 else {
2214 return Poll::Ready(Ok(None));
2215 };
2216
2217 self.datagram_index = found_index + 1;
2218 self.next_sequence = self.next_sequence.max(datagram.sequence.saturating_add(1));
2219 Poll::Ready(Ok(Some(datagram)))
2220 }
2221
2222 fn poll_next_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2223 let floor = self.next_sequence.max(self.min_sequence);
2224 let Some(group) = ready!(self.poll(waiter, |state| state.poll_next_in_range(floor, self.end_sequence))?) else {
2225 return Poll::Ready(Ok(None));
2226 };
2227 self.next_sequence = group.sequence.saturating_add(1);
2228 Poll::Ready(Ok(Some(group)))
2229 }
2230
2231 fn poll_read_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Frame>>> {
2232 let lower = self.min_sequence.max(self.next_sequence);
2233 let Some((frame, found_index, sequence)) =
2234 ready!(self.poll(waiter, |state| { state.poll_read_frame(self.index, lower, waiter) })?)
2235 else {
2236 return Poll::Ready(Ok(None));
2237 };
2238
2239 self.index = found_index + 1;
2240 self.next_sequence = sequence.saturating_add(1);
2241 Poll::Ready(Ok(Some(frame)))
2242 }
2243}
2244
2245#[derive(Clone)]
2251pub struct SubscriberControl {
2252 subscription: kio::Producer<Subscription>,
2253}
2254
2255impl SubscriberControl {
2256 pub fn subscription(&self) -> Subscription {
2258 self.subscription.read().clone()
2259 }
2260
2261 pub fn update(&self, subscription: Subscription) -> Result<()> {
2266 let mut state = self.subscription.write().map_err(|_| Error::Closed)?;
2267 *state = subscription;
2268 Ok(())
2269 }
2270}
2271
2272impl Subscriber {
2273 pub fn info(&self) -> &Info {
2278 &self.info
2279 }
2280
2281 pub fn name(&self) -> &str {
2283 &self.name
2284 }
2285
2286 pub fn control(&self) -> SubscriberControl {
2288 SubscriberControl {
2289 subscription: match &self.inner {
2290 SubscriberKind::Plain(plain) => plain.subscription.clone(),
2291 SubscriberKind::Spliced(spliced) => spliced.prefs(),
2292 },
2293 }
2294 }
2295
2296 pub fn poll_recv_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2307 let meter = self.stats.meter();
2308 let res = match &mut self.inner {
2309 SubscriberKind::Plain(plain) => plain.poll_recv_group(waiter),
2310 SubscriberKind::Spliced(spliced) => spliced.poll_recv_group(waiter),
2311 };
2312 res.map(|res| res.map(|group| group.map(|group| group.with_meter(meter))))
2313 }
2314
2315 pub async fn recv_group(&mut self) -> Result<Option<group::Consumer>> {
2321 kio::wait(|waiter| self.poll_recv_group(waiter)).await
2322 }
2323
2324 pub fn poll_recv_datagram(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Datagram>>> {
2335 let meter = self.stats.meter();
2336 let res = match &mut self.inner {
2337 SubscriberKind::Plain(plain) => plain.poll_recv_datagram(waiter),
2338 SubscriberKind::Spliced(spliced) => spliced.poll_recv_datagram(waiter),
2339 };
2340 if let Poll::Ready(Ok(Some(datagram))) = &res {
2343 meter.datagram(datagram.payload.len() as u64);
2344 }
2345 res
2346 }
2347
2348 pub async fn recv_datagram(&mut self) -> Result<Option<Datagram>> {
2355 kio::wait(|waiter| self.poll_recv_datagram(waiter)).await
2356 }
2357
2358 pub fn poll_next_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2367 let meter = self.stats.meter();
2368 let res = match &mut self.inner {
2369 SubscriberKind::Plain(plain) => plain.poll_next_group(waiter),
2370 SubscriberKind::Spliced(spliced) => spliced.poll_next_group(waiter),
2371 };
2372 res.map(|res| res.map(|group| group.map(|group| group.with_meter(meter))))
2373 }
2374
2375 pub async fn next_group(&mut self) -> Result<Option<group::Consumer>> {
2381 kio::wait(|waiter| self.poll_next_group(waiter)).await
2382 }
2383
2384 pub fn poll_read_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Frame>>> {
2388 let meter = self.stats.meter();
2389 let res = match &mut self.inner {
2390 SubscriberKind::Plain(plain) => plain.poll_read_frame(waiter),
2391 SubscriberKind::Spliced(spliced) => spliced.poll_read_frame(waiter),
2392 };
2393 if let Poll::Ready(Ok(Some(frame))) = &res {
2396 meter.group();
2397 meter.frames(1);
2398 meter.bytes(frame.payload.len() as u64);
2399 }
2400 res
2401 }
2402
2403 pub async fn read_frame(&mut self) -> Result<Option<frame::Frame>> {
2408 kio::wait(|waiter| self.poll_read_frame(waiter)).await
2409 }
2410
2411 pub fn is_clone(&self, other: &Self) -> bool {
2413 match (&self.inner, &other.inner) {
2414 (SubscriberKind::Plain(a), SubscriberKind::Plain(b)) => a.state.same_channel(&b.state),
2415 (SubscriberKind::Spliced(a), SubscriberKind::Spliced(b)) => a.is_clone(b),
2416 _ => false,
2417 }
2418 }
2419
2420 pub fn poll_finished(&mut self, waiter: &kio::Waiter) -> Poll<Result<u64>> {
2422 match &mut self.inner {
2423 SubscriberKind::Plain(plain) => plain.poll(waiter, |state| state.poll_finished()),
2424 SubscriberKind::Spliced(spliced) => spliced.poll_finished(waiter),
2425 }
2426 }
2427
2428 pub async fn finished(&mut self) -> Result<u64> {
2436 kio::wait(|waiter| self.poll_finished(waiter)).await
2437 }
2438
2439 pub fn start_at(&mut self, sequence: u64) {
2446 match &mut self.inner {
2447 SubscriberKind::Plain(plain) => plain.min_sequence = sequence,
2448 SubscriberKind::Spliced(spliced) => spliced.start_at(sequence),
2449 }
2450 }
2451
2452 pub fn end_at(&mut self, sequence: impl Into<Option<u64>>) {
2465 match &mut self.inner {
2466 SubscriberKind::Plain(plain) => plain.end_sequence = sequence.into(),
2467 SubscriberKind::Spliced(spliced) => spliced.end_at(sequence),
2468 }
2469 }
2470
2471 pub fn subscription(&self) -> Subscription {
2473 self.control().subscription()
2474 }
2475
2476 pub fn update(&mut self, subscription: Subscription) -> Result<()> {
2482 match &mut self.inner {
2483 SubscriberKind::Plain(plain) => {
2484 let mut state = plain.subscription.write().map_err(|_| Error::Closed)?;
2485 *state = subscription;
2486 }
2487 SubscriberKind::Spliced(spliced) => spliced.update(subscription),
2488 }
2489 Ok(())
2490 }
2491
2492 pub fn latest(&self) -> Option<u64> {
2494 match &self.inner {
2495 SubscriberKind::Plain(plain) => plain.state.read().max_sequence,
2496 SubscriberKind::Spliced(spliced) => spliced.latest(),
2497 }
2498 }
2499}
2500
2501pub struct Request {
2513 name: Arc<str>,
2514 broadcast: Arc<broadcast::Info>,
2516 state: kio::Producer<TrackState>,
2517
2518 prev_subscription: Option<Subscription>,
2520
2521 alive: Arc<Alive>,
2524
2525 _dynamic: Dynamic,
2530
2531 stats: stats::Scope,
2534}
2535
2536impl Request {
2537 pub(crate) fn new(broadcast: Arc<broadcast::Info>, name: impl Into<Arc<str>>) -> Self {
2538 let name = name.into();
2539 let state = TrackState::spawn(broadcast.clone());
2540 let alive = Alive::new(name.clone(), state.clone());
2541 let dynamic = Dynamic::new(name.clone(), state.clone(), alive.clone());
2542 Self {
2543 name,
2544 broadcast,
2545 state,
2546 prev_subscription: None,
2547 alive,
2548 _dynamic: dynamic,
2549 stats: stats::Scope::default(),
2550 }
2551 }
2552
2553 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
2556 self.stats = scope;
2557 self
2558 }
2559
2560 pub fn name(&self) -> &str {
2562 &self.name
2563 }
2564
2565 pub fn consume(&self) -> Consumer {
2567 Consumer::plain(self.name.clone(), self.state.consume())
2568 }
2569
2570 pub fn dynamic(&self) -> Dynamic {
2574 Dynamic::new(self.name.clone(), self.state.clone(), self.alive.clone())
2575 }
2576
2577 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
2580 self.state.poll_unused(waiter).map(|_| ())
2581 }
2582
2583 pub fn accept(self, info: impl Into<Option<Info>>) -> Producer {
2590 if let Ok(mut state) = self.state.write() {
2593 state.install(info.into().unwrap_or_default());
2594 }
2595 self.alive.publish(Some(&self.stats));
2598 Producer {
2599 name: self.name,
2600 broadcast: self.broadcast,
2601 state: self.state,
2602 prev_subscription: None,
2603 alive: self.alive,
2604 stats: self.stats,
2605 }
2606 }
2607
2608 pub fn reject(self, err: Error) {
2610 if let Ok(mut state) = self.state.write() {
2611 state.abort = Some(err);
2612 }
2613 }
2614
2615 pub fn subscription(&self) -> Option<Subscription> {
2618 let state = self.state.read();
2619 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
2620 drop(state);
2621 snapshot_subscription(&subs, bound)
2622 }
2623
2624 pub async fn subscription_changed(&mut self) -> Option<Subscription> {
2627 kio::wait(|waiter| self.poll_subscription_changed(waiter)).await
2628 }
2629
2630 pub fn poll_subscription_changed(&mut self, waiter: &kio::Waiter) -> Poll<Option<Subscription>> {
2632 let state = self.state.read();
2633 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
2634 drop(state);
2635
2636 let prev = &self.prev_subscription;
2637 let mut combined = None;
2638 let mut guard = ready!(subs.poll(waiter, |subs| {
2639 let next = combined_subscription(subs, bound, waiter);
2640 if &next == prev {
2641 Poll::Pending
2642 } else {
2643 combined = next;
2644 Poll::Ready(())
2645 }
2646 }));
2647 guard.retain(|sub| !sub.is_closed());
2649 drop(guard);
2650 self.prev_subscription = combined.clone();
2651 Poll::Ready(combined)
2652 }
2653
2654 pub(super) fn weak(&self) -> TrackWeak {
2655 TrackWeak {
2656 name: self.name.clone(),
2657 state: self.state.weak(),
2658 }
2659 }
2660}
2661
2662#[cfg(test)]
2663use futures::FutureExt;
2664
2665#[cfg(test)]
2666#[allow(missing_docs)] impl Subscriber {
2668 pub fn assert_group(&mut self) -> group::Consumer {
2669 self.recv_group()
2670 .now_or_never()
2671 .expect("group would have blocked")
2672 .expect("would have errored")
2673 .expect("track was closed")
2674 }
2675
2676 pub fn assert_no_group(&mut self) {
2677 assert!(
2678 self.recv_group().now_or_never().is_none(),
2679 "recv_group would not have blocked"
2680 );
2681 }
2682
2683 pub fn assert_not_closed(&mut self) {
2684 assert!(self.finished().now_or_never().is_none(), "should not be closed");
2685 }
2686
2687 pub fn assert_closed(&mut self) {
2688 assert!(self.finished().now_or_never().is_some(), "should be closed");
2689 }
2690
2691 pub fn assert_error(&mut self) {
2693 assert!(
2694 self.finished().now_or_never().expect("should not block").is_err(),
2695 "should be error"
2696 );
2697 }
2698
2699 pub fn assert_is_clone(&self, other: &Self) {
2700 assert!(self.is_clone(other), "should be clone");
2701 }
2702
2703 pub fn assert_not_clone(&self, other: &Self) {
2704 assert!(!self.is_clone(other), "should not be clone");
2705 }
2706}
2707
2708#[cfg(test)]
2709mod test {
2710 use super::*;
2711
2712 fn track_producer(name: impl Into<Arc<str>>, info: impl Into<Option<Info>>) -> Producer {
2715 Producer::new(Arc::new(broadcast::Info::default()), name, info)
2716 }
2717
2718 fn live_groups(state: &TrackState) -> usize {
2720 state.lookup.len()
2721 }
2722
2723 fn first_live_sequence(state: &TrackState) -> u64 {
2725 state
2726 .arrival
2727 .iter()
2728 .find(|(sequence, stamp)| state.lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp))
2729 .map(|(sequence, _)| *sequence)
2730 .unwrap()
2731 }
2732
2733 fn recv_datagram(dg: &mut Subscriber) -> Datagram {
2735 dg.recv_datagram()
2736 .now_or_never()
2737 .expect("datagram would have blocked")
2738 .expect("would have errored")
2739 .expect("track was closed")
2740 }
2741
2742 #[tokio::test]
2743 async fn append_datagram_shares_group_sequence() {
2744 let mut producer = track_producer("test", None);
2745 let ts = Timestamp::from_millis(10).unwrap();
2746
2747 assert_eq!(producer.append_group().unwrap().sequence, 0);
2749 assert_eq!(producer.append_datagram(ts, &b"a"[..]).unwrap(), 1);
2750 assert_eq!(producer.append_group().unwrap().sequence, 2);
2751 assert_eq!(producer.append_datagram(ts, &b"b"[..]).unwrap(), 3);
2752 assert_eq!(producer.latest(), Some(3));
2753 }
2754
2755 #[tokio::test]
2756 async fn append_datagram_roundtrip() {
2757 let mut producer = track_producer("test", None);
2758 let mut dg = producer.subscribe(None);
2759
2760 let ts = Timestamp::from_millis(42).unwrap();
2761 let seq = producer.append_datagram(ts, &b"hello"[..]).unwrap();
2762
2763 let got = recv_datagram(&mut dg);
2764 assert_eq!(got.sequence, seq);
2765 assert_eq!(got.timestamp, ts);
2766 assert_eq!(&got.payload[..], b"hello");
2767 }
2768
2769 #[tokio::test]
2770 async fn write_datagram_preserves_sequence() {
2771 let mut producer = track_producer("test", None);
2772 let mut dg = producer.subscribe(None);
2773
2774 let ts = Timestamp::from_millis(5).unwrap();
2775 producer
2777 .write_datagram(Datagram {
2778 sequence: 100,
2779 timestamp: ts,
2780 payload: bytes::Bytes::from_static(b"x"),
2781 })
2782 .unwrap();
2783
2784 assert_eq!(recv_datagram(&mut dg).sequence, 100);
2785 assert_eq!(producer.append_group().unwrap().sequence, 101);
2787 }
2788
2789 #[tokio::test]
2790 async fn recv_datagram_advances_ordered_group_cursor() {
2791 let mut producer = track_producer("test", None);
2792 let mut subscriber = producer.subscribe(None);
2793 let ts = Timestamp::from_millis(5).unwrap();
2794
2795 producer
2796 .write_datagram(Datagram {
2797 sequence: 5,
2798 timestamp: ts,
2799 payload: bytes::Bytes::from_static(b"x"),
2800 })
2801 .unwrap();
2802 assert_eq!(recv_datagram(&mut subscriber).sequence, 5);
2803
2804 producer.create_group(group::Info { sequence: 3 }).unwrap();
2805 producer.create_group(group::Info { sequence: 6 }).unwrap();
2806
2807 let group = subscriber
2808 .next_group()
2809 .now_or_never()
2810 .expect("group would have blocked")
2811 .expect("would have errored")
2812 .expect("track was closed");
2813 assert_eq!(group.sequence, 6);
2814 }
2815
2816 #[tokio::test]
2817 async fn datagram_normalized_to_track_timescale() {
2818 let info = Info::default().with_timescale(Timescale::MICRO);
2819 let mut producer = track_producer("test", info);
2820 let mut dg = producer.subscribe(None);
2821
2822 producer
2824 .append_datagram(Timestamp::from_millis(2).unwrap(), &b"z"[..])
2825 .unwrap();
2826 let got = recv_datagram(&mut dg);
2827 assert_eq!(got.timestamp.scale(), Timescale::MICRO);
2828 assert_eq!(got.timestamp.value(), 2_000);
2829 }
2830
2831 #[tokio::test]
2832 async fn datagram_rejects_oversized() {
2833 let mut producer = track_producer("test", None);
2834 let big = bytes::Bytes::from(vec![0u8; crate::model::datagram::MAX_DATAGRAM_PAYLOAD + 1]);
2835 let ts = Timestamp::from_millis(0).unwrap();
2836 assert!(matches!(
2837 producer.append_datagram(ts, big.clone()),
2838 Err(Error::FrameTooLarge)
2839 ));
2840 assert!(matches!(
2841 producer.write_datagram(Datagram {
2842 sequence: 0,
2843 timestamp: ts,
2844 payload: big,
2845 }),
2846 Err(Error::FrameTooLarge)
2847 ));
2848 }
2849
2850 #[tokio::test]
2851 async fn datagram_fanout_to_subscribers() {
2852 let mut producer = track_producer("test", None);
2853 let mut a = producer.subscribe(None);
2855 let mut b = producer.subscribe(None);
2856 let ts = Timestamp::from_millis(1).unwrap();
2857
2858 producer.append_datagram(ts, &b"first"[..]).unwrap();
2859 producer.append_datagram(ts, &b"second"[..]).unwrap();
2860
2861 assert_eq!(&recv_datagram(&mut a).payload[..], b"first");
2863 assert_eq!(&recv_datagram(&mut a).payload[..], b"second");
2864 assert_eq!(&recv_datagram(&mut b).payload[..], b"first");
2865 assert_eq!(&recv_datagram(&mut b).payload[..], b"second");
2866 }
2867
2868 #[tokio::test]
2869 async fn datagram_evicts_stale() {
2870 tokio::time::pause();
2871
2872 let mut producer = track_producer("test", None);
2873 let mut dg = producer.subscribe(None);
2874 let ts = Timestamp::from_millis(0).unwrap();
2875
2876 producer.append_datagram(ts, &b"old"[..]).unwrap(); tokio::time::advance(MAX_DATAGRAM_AGE + Duration::from_millis(10)).await;
2880 producer.append_datagram(ts, &b"new"[..]).unwrap(); let got = recv_datagram(&mut dg);
2884 assert_eq!(got.sequence, 1);
2885 assert_eq!(&got.payload[..], b"new");
2886 }
2887
2888 #[tokio::test]
2889 async fn datagram_recv_pends_until_written() {
2890 let mut producer = track_producer("test", None);
2891 let mut dg = producer.subscribe(None);
2892
2893 assert!(
2894 dg.recv_datagram().now_or_never().is_none(),
2895 "should block with no datagrams"
2896 );
2897
2898 producer
2899 .append_datagram(Timestamp::from_millis(0).unwrap(), &b"go"[..])
2900 .unwrap();
2901 assert_eq!(&recv_datagram(&mut dg).payload[..], b"go");
2902 }
2903
2904 #[tokio::test]
2908 async fn datagram_wire_roundtrip_between_tracks() {
2909 use crate::coding::{Decode, Encode};
2910 use crate::lite;
2911
2912 let version = lite::Version::Lite05;
2913
2914 let mut origin = track_producer("test", None);
2916 let mut origin_dg = origin.subscribe(None);
2917 let ts = Timestamp::from_millis(7).unwrap();
2918 let seq = origin.append_datagram(ts, &b"payload"[..]).unwrap();
2919
2920 let d = recv_datagram(&mut origin_dg);
2921 let body = lite::Datagram {
2922 subscribe: 5,
2923 sequence: d.sequence,
2924 timestamp: d.timestamp.value(),
2925 payload: d.payload.clone(),
2926 }
2927 .encode_bytes(version)
2928 .unwrap();
2929
2930 let mut slice = &body[..];
2932 let wire = lite::Datagram::decode(&mut slice, version).unwrap();
2933 let mut downstream = track_producer("test", None);
2934 let mut downstream_dg = downstream.subscribe(None);
2935 downstream
2936 .write_datagram(Datagram {
2937 sequence: wire.sequence,
2938 timestamp: Timestamp::new(wire.timestamp, Timescale::MILLI).unwrap(),
2939 payload: wire.payload,
2940 })
2941 .unwrap();
2942
2943 let got = recv_datagram(&mut downstream_dg);
2944 assert_eq!(got.sequence, seq);
2945 assert_eq!(got.timestamp, ts);
2946 assert_eq!(&got.payload[..], b"payload");
2947 }
2948
2949 #[tokio::test]
2950 async fn evict_expired_groups() {
2951 tokio::time::pause();
2952
2953 let mut producer = track_producer("test", None);
2954
2955 producer.append_group().unwrap(); producer.append_group().unwrap(); producer.append_group().unwrap(); {
2961 let state = producer.state.read();
2962 assert_eq!(live_groups(&state), 3);
2963 assert_eq!(state.offset, 0);
2964 }
2965
2966 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
2968
2969 producer.append_group().unwrap(); {
2976 let state = producer.state.read();
2977 assert_eq!(live_groups(&state), 1);
2978 assert_eq!(first_live_sequence(&state), 3);
2979 assert_eq!(state.offset, 3);
2980 assert!(!state.lookup.contains_key(&0));
2981 assert!(!state.lookup.contains_key(&1));
2982 assert!(!state.lookup.contains_key(&2));
2983 assert!(state.lookup.contains_key(&3));
2984 }
2985 }
2986
2987 #[tokio::test]
2991 async fn aging_out_a_finished_group_keeps_the_clean_end() {
2992 tokio::time::pause();
2993
2994 let mut producer = track_producer("test", None);
2995 let mut group = producer.create_group(group::Info { sequence: 0 }).unwrap();
2996 let mut consumer = group.consume();
2997
2998 group
2999 .write_frame(Timestamp::from_millis(0).unwrap(), b"hello".as_slice())
3000 .unwrap();
3001 assert_eq!(consumer.next_frame().await.unwrap().unwrap().size, 5);
3002
3003 tokio::time::advance(DEFAULT_LATENCY_MAX * 12).await;
3005 group.finish().unwrap();
3006 let _next = producer.create_group(group::Info { sequence: 1 }).unwrap();
3007
3008 assert!(consumer.next_frame().await.unwrap().is_none());
3009 }
3010
3011 #[tokio::test]
3012 async fn evict_keeps_max_sequence() {
3013 tokio::time::pause();
3014
3015 let mut producer = track_producer("test", None);
3016 producer.append_group().unwrap(); tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3020
3021 producer.append_group().unwrap(); {
3025 let state = producer.state.read();
3026 assert_eq!(live_groups(&state), 1);
3027 assert_eq!(first_live_sequence(&state), 1);
3028 assert_eq!(state.offset, 1);
3029 }
3030 }
3031
3032 #[tokio::test]
3033 async fn no_eviction_when_fresh() {
3034 tokio::time::pause();
3035
3036 let mut producer = track_producer("test", None);
3037 producer.append_group().unwrap(); producer.append_group().unwrap(); producer.append_group().unwrap(); {
3042 let state = producer.state.read();
3043 assert_eq!(live_groups(&state), 3);
3044 assert_eq!(state.offset, 0);
3045 }
3046 }
3047
3048 #[tokio::test]
3049 async fn consumer_skips_evicted_groups() {
3050 tokio::time::pause();
3051
3052 let mut producer = track_producer("test", None);
3053 producer.append_group().unwrap(); let mut consumer = producer.subscribe(None);
3056
3057 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3058 producer.append_group().unwrap(); let group = consumer.assert_group();
3062 assert_eq!(group.sequence, 1);
3063 }
3064
3065 #[tokio::test]
3066 async fn cache_age_controls_eviction() {
3067 tokio::time::pause();
3068
3069 let mut producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(1)));
3071 producer.append_group().unwrap(); tokio::time::advance(Duration::from_secs(2)).await;
3075 producer.append_group().unwrap(); let state = producer.state.read();
3079 assert_eq!(live_groups(&state), 1);
3080 assert_eq!(first_live_sequence(&state), 1);
3081 }
3082
3083 #[test]
3084 fn latency_max_clamped_to_cache() {
3085 let producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(2)));
3086
3087 let mut subscriber = producer.subscribe(Subscription::default().with_latency_max(Duration::from_secs(10)));
3091 assert_eq!(subscriber.subscription().latency_max, Duration::from_secs(10));
3092 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3093
3094 subscriber
3096 .update(Subscription::default().with_latency_max(Duration::from_millis(500)))
3097 .unwrap();
3098 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_millis(500));
3099
3100 subscriber
3101 .update(Subscription::default().with_latency_max(Duration::ZERO))
3102 .unwrap();
3103 assert_eq!(producer.subscription().unwrap().latency_max, Duration::ZERO);
3104 }
3105
3106 fn track_producer_capped(name: impl Into<Arc<str>>, info: Info, cap: Duration) -> Producer {
3109 let origin = crate::origin::Info::default().with_cache_duration(cap);
3110 Producer::new(Arc::new(broadcast::Info { origin }), name, info)
3111 }
3112
3113 #[test]
3114 fn origin_cache_duration_clamps_latency_max() {
3115 let capped = track_producer_capped(
3118 "test",
3119 Info::default().with_latency_max(Duration::from_secs(60)),
3120 Duration::from_secs(1),
3121 );
3122 assert_eq!(capped.state.read().latency_bound(), Some(Duration::from_secs(1)));
3123
3124 let under = track_producer_capped(
3125 "test",
3126 Info::default().with_latency_max(Duration::from_millis(500)),
3127 Duration::from_secs(1),
3128 );
3129 assert_eq!(under.state.read().latency_bound(), Some(Duration::from_millis(500)));
3130 }
3131
3132 #[tokio::test]
3133 async fn origin_cache_duration_caps_eviction() {
3134 tokio::time::pause();
3135
3136 let mut producer = track_producer_capped(
3138 "test",
3139 Info::default().with_latency_max(Duration::from_secs(60)),
3140 Duration::from_secs(1),
3141 );
3142 producer.append_group().unwrap(); tokio::time::advance(Duration::from_secs(2)).await;
3146 producer.append_group().unwrap(); let state = producer.state.read();
3150 assert_eq!(live_groups(&state), 1);
3151 assert_eq!(first_live_sequence(&state), 1);
3152 }
3153
3154 #[test]
3155 fn latency_max_clamped_via_every_update_path() {
3156 let producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(2)));
3157 let over = Subscription::default().with_latency_max(Duration::from_secs(10));
3158
3159 let mut subscriber = producer.subscribe(over.clone());
3162 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3163
3164 subscriber.control().update(over.clone()).unwrap();
3165 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3166
3167 subscriber.update(over).unwrap();
3168 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3169 }
3170
3171 #[test]
3172 fn latency_max_aggregate_clamps_the_max_across_subscribers() {
3173 let producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(2)));
3174
3175 let _a = producer.subscribe(Subscription::default().with_latency_max(Duration::from_millis(500)));
3178 let _b = producer.subscribe(Subscription::default().with_latency_max(Duration::from_secs(10)));
3179
3180 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3181 }
3182
3183 #[test]
3184 fn subscriber_control_updates_while_read_future_is_pending() {
3185 let producer = track_producer("test", None);
3186 let mut subscriber = producer.subscribe(None);
3187 let control = subscriber.control();
3188
3189 let mut recv = Box::pin(subscriber.recv_group());
3190 assert!(recv.as_mut().now_or_never().is_none());
3191
3192 control
3193 .update(Subscription::default().with_priority(7).with_ordered(false))
3194 .unwrap();
3195
3196 let aggregate = producer.subscription().expect("expected an active subscription");
3197 assert_eq!(aggregate.priority, 7);
3198 assert!(!aggregate.ordered);
3199 }
3200
3201 #[test]
3202 fn dropped_subscriber_leaves_no_ghost_in_aggregate() {
3203 let mut producer = track_producer("test", None);
3208 let a = producer.subscribe(Subscription::default().with_priority(5));
3209
3210 let waiter = kio::Waiter::noop();
3212 assert!(
3213 matches!(producer.poll_subscription_changed(&waiter), Poll::Ready(Ok(Some(_)))),
3214 "one live subscriber should aggregate to Some",
3215 );
3216
3217 drop(a);
3219
3220 assert!(
3222 matches!(producer.poll_subscription_changed(&waiter), Poll::Ready(Ok(None))),
3223 "a dropped subscriber must not linger in the aggregate",
3224 );
3225
3226 assert!(
3228 producer.subscription().is_none(),
3229 "snapshot must exclude a dropped subscriber",
3230 );
3231 }
3232
3233 #[test]
3234 fn dropped_subscriber_wakes_the_aggregate() {
3235 use std::sync::atomic::{AtomicBool, Ordering};
3242
3243 let mut producer = track_producer("test", None);
3244 let a = producer.subscribe(Subscription::default().with_priority(5));
3245
3246 let woken = Arc::new(AtomicBool::new(false));
3247 let waiter = kio::Waiter::new(futures::task::waker(Arc::new(FlagWake(woken.clone()))));
3248
3249 assert!(matches!(
3251 producer.poll_subscription_changed(&waiter),
3252 Poll::Ready(Ok(Some(_)))
3253 ));
3254 assert!(
3255 producer.poll_subscription_changed(&waiter).is_pending(),
3256 "the aggregate is unchanged, so this poll must park",
3257 );
3258 assert!(!woken.load(Ordering::SeqCst), "nothing happened yet");
3259
3260 drop(a);
3261 assert!(
3262 woken.load(Ordering::SeqCst),
3263 "the last subscriber leaving must wake the aggregate watcher",
3264 );
3265 }
3266
3267 struct FlagWake(Arc<std::sync::atomic::AtomicBool>);
3269
3270 impl futures::task::ArcWake for FlagWake {
3271 fn wake_by_ref(arc_self: &Arc<Self>) {
3272 arc_self.0.store(true, std::sync::atomic::Ordering::SeqCst);
3273 }
3274 }
3275
3276 #[tokio::test]
3277 async fn out_of_order_max_sequence_at_front() {
3278 tokio::time::pause();
3279
3280 let mut producer = track_producer("test", None);
3281
3282 producer.create_group(group::Info { sequence: 5 }).unwrap();
3284 producer.create_group(group::Info { sequence: 3 }).unwrap();
3285 producer.create_group(group::Info { sequence: 4 }).unwrap();
3286
3287 {
3289 let state = producer.state.read();
3290 assert_eq!(state.max_sequence, Some(5));
3291 }
3292
3293 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3295
3296 producer.append_group().unwrap(); {
3302 let state = producer.state.read();
3303 assert_eq!(live_groups(&state), 1);
3304 assert_eq!(first_live_sequence(&state), 6);
3305 assert!(!state.lookup.contains_key(&3));
3306 assert!(!state.lookup.contains_key(&4));
3307 assert!(!state.lookup.contains_key(&5));
3308 assert!(state.lookup.contains_key(&6));
3309 }
3310 }
3311
3312 #[tokio::test]
3313 async fn max_sequence_at_front_blocks_trim() {
3314 tokio::time::pause();
3315
3316 let mut producer = track_producer("test", None);
3317
3318 producer.create_group(group::Info { sequence: 5 }).unwrap();
3320
3321 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3322
3323 producer.create_group(group::Info { sequence: 3 }).unwrap();
3325
3326 {
3329 let state = producer.state.read();
3330 assert_eq!(live_groups(&state), 2);
3331 assert_eq!(state.offset, 0);
3332 }
3333
3334 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3336
3337 producer.create_group(group::Info { sequence: 2 }).unwrap();
3339
3340 {
3345 let state = producer.state.read();
3346 assert_eq!(live_groups(&state), 2);
3347 assert_eq!(state.offset, 0);
3348 assert!(state.lookup.contains_key(&5));
3349 assert!(!state.lookup.contains_key(&3));
3350 assert!(state.lookup.contains_key(&2));
3351 }
3352
3353 let mut consumer = producer.subscribe(None);
3355 let group = consumer.assert_group();
3356 assert_eq!(group.sequence, 5);
3358 }
3359
3360 #[tokio::test]
3361 async fn abort_clears_cached_groups() {
3362 let mut producer = track_producer("test", None);
3363 producer.append_group().unwrap();
3364 producer.append_group().unwrap();
3365
3366 let mut consumer = producer.subscribe(None);
3368 assert_eq!(live_groups(&producer.state.read()), 2);
3369
3370 producer.clone().abort(Error::Cancel).unwrap();
3371
3372 {
3373 let state = producer.state.read();
3374 assert!(state.lookup.is_empty(), "cached groups should be dropped on abort");
3375 assert!(state.arrival.is_empty());
3376 assert!(state.evict.is_empty());
3377 }
3378
3379 let result = consumer.recv_group().now_or_never().expect("should not block");
3381 assert!(matches!(result, Err(Error::Cancel)));
3382 }
3383
3384 #[tokio::test]
3385 async fn drop_unfinished_clears_cached_groups() {
3386 let producer = track_producer("test", None);
3387 let mut writer = producer.clone();
3388 writer.append_group().unwrap();
3389
3390 let mut consumer = producer.subscribe(None);
3392 assert_eq!(live_groups(&producer.state.read()), 1);
3393
3394 drop(writer);
3396 drop(producer);
3397
3398 let result = consumer.recv_group().now_or_never().expect("should not block");
3399 assert!(matches!(result, Err(Error::Dropped)));
3400 }
3401
3402 #[tokio::test]
3403 async fn drop_finished_keeps_cached_groups() {
3404 let mut producer = track_producer("test", None);
3405 producer.append_group().unwrap();
3406 producer.finish().unwrap();
3407
3408 let mut consumer = producer.subscribe(None);
3409 drop(producer);
3410
3411 assert_eq!(consumer.assert_group().sequence, 0);
3413 let done = consumer.recv_group().now_or_never().expect("should not block").unwrap();
3414 assert!(done.is_none(), "consumer should drain then see clean finish");
3415 }
3416
3417 #[test]
3418 fn append_finish_cannot_be_rewritten() {
3419 let mut producer = track_producer("test", None);
3420
3421 assert!(producer.finish().is_ok());
3423 assert!(producer.finish().is_err());
3424 assert!(producer.append_group().is_err());
3425 }
3426
3427 #[test]
3428 fn finish_after_groups() {
3429 let mut producer = track_producer("test", None);
3430
3431 producer.append_group().unwrap();
3432 assert!(producer.finish().is_ok());
3433 assert!(producer.finish().is_err());
3434 assert!(producer.append_group().is_err());
3435 }
3436
3437 #[test]
3438 fn finish_at_rejects_a_boundary_at_or_below_the_live_edge() {
3439 let mut producer = track_producer("test", None);
3440 producer.create_group(group::Info { sequence: 5 }).unwrap();
3441
3442 assert!(producer.finish_at(4).is_err());
3445 assert!(producer.finish_at(5).is_err());
3446 assert!(producer.finish_at(6).is_ok());
3447
3448 {
3449 let state = producer.state.read();
3450 assert_eq!(state.final_sequence, Some(6));
3451 }
3452
3453 assert!(producer.finish_at(6).is_err());
3455 assert!(producer.create_group(group::Info { sequence: 4 }).is_ok());
3456 assert!(producer.create_group(group::Info { sequence: 6 }).is_err());
3457 }
3458
3459 #[test]
3460 fn final_sequence_reports_the_declared_boundary() {
3461 let mut producer = track_producer("test", None);
3462 assert_eq!(producer.final_sequence(), None);
3463
3464 producer.create_group(group::Info { sequence: 5 }).unwrap();
3465 assert_eq!(producer.final_sequence(), None, "a group does not declare a boundary");
3466
3467 producer.finish_at(9).unwrap();
3468 assert_eq!(producer.final_sequence(), Some(9));
3469
3470 assert!(producer.finish().is_err());
3472 }
3473
3474 #[test]
3475 fn final_sequence_reports_the_live_edge_after_finish() {
3476 let mut producer = track_producer("test", None);
3477 producer.create_group(group::Info { sequence: 5 }).unwrap();
3478 producer.finish().unwrap();
3479 assert_eq!(producer.final_sequence(), Some(6));
3480 }
3481
3482 #[tokio::test]
3483 async fn finish_at_declares_a_future_boundary() {
3484 let mut producer = track_producer("test", None);
3485 producer.create_group(group::Info { sequence: 5 }).unwrap();
3486
3487 producer.finish_at(7).unwrap();
3489
3490 let mut consumer = producer.subscribe(None);
3491 assert_eq!(consumer.assert_group().sequence, 5);
3492
3493 let boundary = consumer
3496 .finished()
3497 .now_or_never()
3498 .expect("boundary is known immediately")
3499 .expect("would have errored");
3500 assert_eq!(boundary, 7);
3501 assert!(
3502 consumer.recv_group().now_or_never().is_none(),
3503 "should wait for the outstanding group"
3504 );
3505
3506 producer.create_group(group::Info { sequence: 6 }).unwrap();
3508 assert_eq!(consumer.assert_group().sequence, 6);
3509 let done = consumer
3510 .recv_group()
3511 .now_or_never()
3512 .expect("should not block")
3513 .expect("would have errored");
3514 assert!(done.is_none(), "track completes once the boundary is reached");
3515 }
3516
3517 #[tokio::test]
3518 async fn recv_group_finishes_without_waiting_for_gaps() {
3519 let mut producer = track_producer("test", None);
3520 producer.create_group(group::Info { sequence: 1 }).unwrap();
3521 producer.finish().unwrap();
3522
3523 let mut consumer = producer.subscribe(None);
3524 assert_eq!(consumer.assert_group().sequence, 1);
3525
3526 let done = consumer
3527 .recv_group()
3528 .now_or_never()
3529 .expect("should not block")
3530 .expect("would have errored");
3531 assert!(done.is_none(), "track should finish without waiting for gaps");
3532 }
3533
3534 #[tokio::test]
3535 async fn next_group_skips_late_arrivals() {
3536 let mut producer = track_producer("test", None);
3537 let mut consumer = producer.subscribe(None);
3538
3539 producer.create_group(group::Info { sequence: 5 }).unwrap();
3541 let group = consumer
3542 .next_group()
3543 .now_or_never()
3544 .expect("should not block")
3545 .expect("would have errored")
3546 .expect("track should not be closed");
3547 assert_eq!(group.sequence, 5);
3548
3549 producer.create_group(group::Info { sequence: 3 }).unwrap();
3551 producer.create_group(group::Info { sequence: 4 }).unwrap();
3553 producer.create_group(group::Info { sequence: 7 }).unwrap();
3555
3556 let group = consumer
3557 .next_group()
3558 .now_or_never()
3559 .expect("should not block")
3560 .expect("would have errored")
3561 .expect("track should not be closed");
3562 assert_eq!(group.sequence, 7);
3563
3564 assert!(
3566 consumer.next_group().now_or_never().is_none(),
3567 "should block waiting for a higher sequence"
3568 );
3569 }
3570
3571 #[tokio::test]
3572 async fn next_group_returns_arrivals_in_order() {
3573 let mut producer = track_producer("test", None);
3574 let mut consumer = producer.subscribe(None);
3575
3576 producer.create_group(group::Info { sequence: 3 }).unwrap();
3578 producer.create_group(group::Info { sequence: 5 }).unwrap();
3579
3580 let group = consumer
3581 .next_group()
3582 .now_or_never()
3583 .expect("should not block")
3584 .expect("would have errored")
3585 .expect("track should not be closed");
3586 assert_eq!(group.sequence, 3);
3587
3588 let group = consumer
3589 .next_group()
3590 .now_or_never()
3591 .expect("should not block")
3592 .expect("would have errored")
3593 .expect("track should not be closed");
3594 assert_eq!(group.sequence, 5);
3595 }
3596
3597 #[tokio::test]
3598 async fn next_group_and_recv_group_use_independent_cursors() {
3599 let mut producer = track_producer("test", None);
3600 let mut consumer = producer.subscribe(None);
3601
3602 producer.create_group(group::Info { sequence: 5 }).unwrap();
3604 producer.create_group(group::Info { sequence: 3 }).unwrap();
3605
3606 let group = consumer
3609 .next_group()
3610 .now_or_never()
3611 .expect("should not block")
3612 .expect("would have errored")
3613 .expect("track should not be closed");
3614 assert_eq!(group.sequence, 3);
3615
3616 assert_eq!(consumer.assert_group().sequence, 5);
3619 }
3620
3621 #[tokio::test]
3622 async fn end_at_caps_next_group() {
3623 let mut producer = track_producer("test", None);
3624 let mut consumer = producer.subscribe(None);
3625
3626 for s in 0..6 {
3627 producer.create_group(group::Info { sequence: s }).unwrap();
3628 }
3629
3630 consumer.end_at(2);
3631
3632 assert_eq!(
3634 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3635 0
3636 );
3637 assert_eq!(
3638 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3639 1
3640 );
3641 assert_eq!(
3642 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3643 2
3644 );
3645
3646 assert!(
3648 consumer.next_group().now_or_never().is_none(),
3649 "capped consumer must block instead of returning out-of-range groups"
3650 );
3651 }
3652
3653 #[tokio::test]
3654 async fn end_at_release_drains_cached_groups() {
3655 let mut producer = track_producer("test", None);
3656 let mut consumer = producer.subscribe(None);
3657
3658 for s in 0..6 {
3659 producer.create_group(group::Info { sequence: s }).unwrap();
3660 }
3661
3662 consumer.end_at(1);
3663 assert_eq!(
3664 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3665 0
3666 );
3667 assert_eq!(
3668 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3669 1
3670 );
3671 assert!(consumer.next_group().now_or_never().is_none(), "capped at 1");
3672
3673 consumer.end_at(4);
3675 assert_eq!(
3676 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3677 2
3678 );
3679 assert_eq!(
3680 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3681 3
3682 );
3683 assert_eq!(
3684 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3685 4
3686 );
3687 assert!(consumer.next_group().now_or_never().is_none(), "capped at 4");
3688
3689 consumer.end_at(None);
3691 assert_eq!(
3692 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3693 5
3694 );
3695 assert!(consumer.next_group().now_or_never().is_none(), "no more groups");
3696 }
3697
3698 #[tokio::test]
3699 async fn end_at_lower_than_cursor_parks_consumer() {
3700 let mut producer = track_producer("test", None);
3701 let mut consumer = producer.subscribe(None);
3702
3703 for s in 0..3 {
3704 producer.create_group(group::Info { sequence: s }).unwrap();
3705 }
3706
3707 assert_eq!(
3709 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3710 0
3711 );
3712 assert_eq!(
3713 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3714 1
3715 );
3716 assert_eq!(
3717 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3718 2
3719 );
3720
3721 consumer.end_at(1);
3723 producer.create_group(group::Info { sequence: 3 }).unwrap();
3724 producer.create_group(group::Info { sequence: 4 }).unwrap();
3725 assert!(
3726 consumer.next_group().now_or_never().is_none(),
3727 "cap is below cursor; nothing returnable until cap rises"
3728 );
3729
3730 consumer.end_at(None);
3732 assert_eq!(
3733 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3734 3
3735 );
3736 assert_eq!(
3737 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3738 4
3739 );
3740 }
3741
3742 #[tokio::test]
3743 async fn end_at_toggling_around_late_arrivals() {
3744 let mut producer = track_producer("test", None);
3745 let mut consumer = producer.subscribe(None);
3746
3747 consumer.end_at(5);
3748
3749 producer.create_group(group::Info { sequence: 2 }).unwrap();
3751 producer.create_group(group::Info { sequence: 5 }).unwrap();
3752 producer.create_group(group::Info { sequence: 3 }).unwrap();
3753 producer.create_group(group::Info { sequence: 8 }).unwrap();
3755 producer.create_group(group::Info { sequence: 4 }).unwrap();
3756
3757 assert_eq!(
3759 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3760 2
3761 );
3762 assert_eq!(
3763 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3764 3
3765 );
3766 assert_eq!(
3767 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3768 4
3769 );
3770 assert_eq!(
3771 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3772 5
3773 );
3774 assert!(consumer.next_group().now_or_never().is_none());
3776
3777 consumer.end_at(10);
3779 assert_eq!(
3780 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3781 8
3782 );
3783 }
3784
3785 #[tokio::test]
3786 async fn read_frame_returns_single_frame_per_group() {
3787 let mut producer = track_producer("test", None);
3788 let mut consumer = producer.subscribe(None);
3789
3790 producer.write_frame(Timestamp::ZERO, b"hello".as_slice()).unwrap();
3791 producer.write_frame(Timestamp::ZERO, b"world".as_slice()).unwrap();
3792
3793 let frame = consumer
3794 .read_frame()
3795 .now_or_never()
3796 .expect("should not block")
3797 .expect("would have errored")
3798 .expect("track should not be closed");
3799 assert_eq!(&frame.payload[..], b"hello");
3800
3801 let frame = consumer
3802 .read_frame()
3803 .now_or_never()
3804 .expect("should not block")
3805 .expect("would have errored")
3806 .expect("track should not be closed");
3807 assert_eq!(&frame.payload[..], b"world");
3808 }
3809
3810 #[tokio::test]
3811 async fn read_frame_preserves_timestamp() {
3812 let mut producer = track_producer("test", None);
3813 let mut consumer = producer.subscribe(None);
3814
3815 producer
3816 .write_frame(Timestamp::from_micros(20_000).unwrap(), b"hello".as_slice())
3817 .unwrap();
3818
3819 let frame = consumer
3820 .read_frame()
3821 .now_or_never()
3822 .expect("should not block")
3823 .expect("would have errored")
3824 .expect("track should not be closed");
3825 assert_eq!(frame.timestamp.as_micros(), 20_000);
3826 assert_eq!(&frame.payload[..], b"hello");
3827 }
3828
3829 #[tokio::test]
3830 async fn read_frame_skips_stalled_group_for_newer_ready_frame() {
3831 let mut producer = track_producer("test", None);
3832 let mut consumer = producer.subscribe(None);
3833
3834 let _stalled = producer.create_group(group::Info { sequence: 3 }).unwrap();
3836 let mut g5 = producer.create_group(group::Info { sequence: 5 }).unwrap();
3838 g5.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"later"))
3839 .unwrap();
3840 g5.finish().unwrap();
3841
3842 let frame = consumer
3844 .read_frame()
3845 .now_or_never()
3846 .expect("should not block on stalled earlier group")
3847 .expect("would have errored")
3848 .expect("track should not be closed");
3849 assert_eq!(&frame.payload[..], b"later");
3850 }
3851
3852 #[tokio::test]
3853 async fn read_frame_discards_rest_of_multi_frame_group() {
3854 let mut producer = track_producer("test", None);
3855 let mut consumer = producer.subscribe(None);
3856
3857 let mut g0 = producer.create_group(group::Info { sequence: 0 }).unwrap();
3859 g0.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"one"))
3860 .unwrap();
3861 g0.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"two"))
3862 .unwrap();
3863 g0.finish().unwrap();
3864
3865 producer.write_frame(Timestamp::ZERO, b"next".as_slice()).unwrap();
3867
3868 let frame = consumer
3869 .read_frame()
3870 .now_or_never()
3871 .expect("should not block")
3872 .expect("would have errored")
3873 .expect("track should not be closed");
3874 assert_eq!(&frame.payload[..], b"one");
3875
3876 let frame = consumer
3878 .read_frame()
3879 .now_or_never()
3880 .expect("should not block")
3881 .expect("would have errored")
3882 .expect("track should not be closed");
3883 assert_eq!(&frame.payload[..], b"next");
3884 }
3885
3886 #[tokio::test]
3887 async fn read_frame_waits_for_pending_group_after_finish() {
3888 let mut producer = track_producer("test", None);
3891 let mut consumer = producer.subscribe(None);
3892
3893 let mut g0 = producer.create_group(group::Info { sequence: 0 }).unwrap();
3894 producer.finish().unwrap();
3895
3896 assert!(
3898 consumer.read_frame().now_or_never().is_none(),
3899 "read_frame must block on a pending group even after finish()"
3900 );
3901
3902 g0.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"late"))
3904 .unwrap();
3905 let frame = consumer
3906 .read_frame()
3907 .now_or_never()
3908 .expect("should not block once a frame is written")
3909 .expect("would have errored")
3910 .expect("track should not be closed");
3911 assert_eq!(&frame.payload[..], b"late");
3912 }
3913
3914 #[tokio::test]
3915 async fn read_frame_respects_start_at() {
3916 let mut producer = track_producer("test", None);
3919 let mut consumer = producer.subscribe(None);
3920 consumer.start_at(5);
3921
3922 let mut g3 = producer.create_group(group::Info { sequence: 3 }).unwrap();
3924 g3.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"skip-me"))
3925 .unwrap();
3926 g3.finish().unwrap();
3927
3928 let mut g5 = producer.create_group(group::Info { sequence: 5 }).unwrap();
3929 g5.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"keep"))
3930 .unwrap();
3931 g5.finish().unwrap();
3932
3933 let frame = consumer
3934 .read_frame()
3935 .now_or_never()
3936 .expect("should not block")
3937 .expect("would have errored")
3938 .expect("track should not be closed");
3939 assert_eq!(&frame.payload[..], b"keep");
3940 }
3941
3942 #[tokio::test]
3943 async fn read_frame_returns_none_when_finished() {
3944 let mut producer = track_producer("test", None);
3945 let mut consumer = producer.subscribe(None);
3946
3947 producer.write_frame(Timestamp::ZERO, b"only".as_slice()).unwrap();
3948 producer.finish().unwrap();
3949
3950 let frame = consumer
3951 .read_frame()
3952 .now_or_never()
3953 .expect("should not block")
3954 .expect("would have errored")
3955 .expect("track should not be closed");
3956 assert_eq!(&frame.payload[..], b"only");
3957
3958 let done = consumer
3959 .read_frame()
3960 .now_or_never()
3961 .expect("should not block")
3962 .expect("would have errored");
3963 assert!(done.is_none());
3964 }
3965
3966 #[test]
3967 fn append_group_returns_bounds_exceeded_on_sequence_overflow() {
3968 let mut producer = track_producer("test", None);
3969 {
3970 let mut state = producer.state.write().ok().unwrap();
3971 state.max_sequence = Some(u64::MAX);
3972 }
3973
3974 assert!(matches!(producer.append_group(), Err(Error::BoundsExceeded(_))));
3975 }
3976
3977 #[tokio::test]
3978 async fn fetch_cache_hit() {
3979 let mut producer = track_producer("test", None);
3980
3981 let mut group = producer.append_group().unwrap(); group
3984 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hello"))
3985 .unwrap();
3986 group.finish().unwrap();
3987
3988 let dynamic = producer.dynamic();
3991 let consumer = producer.consume();
3992 assert!(consumer.peek_group(0).is_some());
3993 let mut g = consumer.fetch_group(0, None).await.unwrap();
3994 assert_eq!(g.sequence, 0);
3995 assert_eq!(&g.read_frame().await.unwrap().unwrap().payload[..], b"hello");
3996
3997 assert!(dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending());
3999 }
4000
4001 #[tokio::test]
4002 async fn fetch_miss_signals_dynamic() {
4003 let producer = track_producer("test", None);
4004 let dynamic = producer.dynamic();
4005 let consumer = producer.consume();
4006
4007 assert!(consumer.peek_group(5).is_none());
4011 let pending = consumer.fetch_group(5, group::Fetch::default().with_priority(7));
4012 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
4013
4014 let req = dynamic
4015 .requested_group()
4016 .now_or_never()
4017 .expect("should not block")
4018 .unwrap();
4019 assert_eq!(req.sequence(), 5);
4020 assert_eq!(req.priority(), 7);
4021
4022 let mut group = req.accept(None).unwrap();
4024 group
4025 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hi"))
4026 .unwrap();
4027 group.finish().unwrap();
4028
4029 let mut g = pending.await.unwrap();
4030 assert_eq!(g.sequence, 5);
4031 assert_eq!(&g.read_frame().await.unwrap().unwrap().payload[..], b"hi");
4032 }
4033
4034 #[tokio::test]
4035 async fn fetch_miss_rejects() {
4036 let producer = track_producer("test", None);
4037 let dynamic = producer.dynamic();
4038 let consumer = producer.consume();
4039
4040 let pending = consumer.fetch_group(5, None);
4041 let req = dynamic
4042 .requested_group()
4043 .now_or_never()
4044 .expect("should not block")
4045 .unwrap();
4046
4047 req.reject(Error::Cancel);
4048 assert!(matches!(pending.await, Err(Error::Cancel)));
4049 let fetch = producer.state.read().fetch.clone();
4050 assert!(fetch.read().is_empty());
4051 }
4052
4053 #[tokio::test]
4054 async fn fetch_miss_drop_rejects() {
4055 let producer = track_producer("test", None);
4056 let dynamic = producer.dynamic();
4057 let consumer = producer.consume();
4058
4059 let pending = consumer.fetch_group(5, None);
4060 let req = dynamic
4061 .requested_group()
4062 .now_or_never()
4063 .expect("should not block")
4064 .unwrap();
4065
4066 drop(req);
4067 assert!(matches!(pending.await, Err(Error::Dropped)));
4068 }
4069
4070 #[tokio::test]
4071 async fn fetch_reject_does_not_poison_retry() {
4072 let producer = track_producer("test", None);
4073 let dynamic = producer.dynamic();
4074 let consumer = producer.consume();
4075
4076 let pending = consumer.fetch_group(5, None);
4077 let req = dynamic
4078 .requested_group()
4079 .now_or_never()
4080 .expect("should not block")
4081 .unwrap();
4082 req.reject(Error::Cancel);
4083 assert!(matches!(pending.await, Err(Error::Cancel)));
4084
4085 let retry = consumer.fetch_group(5, None);
4086 let req = dynamic
4087 .requested_group()
4088 .now_or_never()
4089 .expect("should not block")
4090 .unwrap();
4091 let mut group = req.accept(None).unwrap();
4092 group
4093 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"retry"))
4094 .unwrap();
4095 group.finish().unwrap();
4096
4097 let mut group = retry.await.unwrap();
4098 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"retry");
4099 }
4100
4101 #[tokio::test]
4102 async fn fetch_coalesces_concurrent() {
4103 let producer = track_producer("test", None);
4104 let dynamic = producer.dynamic();
4105 let consumer = producer.consume();
4106
4107 let first = consumer.fetch_group(5, group::Fetch::default().with_priority(1));
4110 let second = consumer.fetch_group(5, group::Fetch::default().with_priority(7));
4111 assert!(kio::Pollable::poll(&*first, &kio::Waiter::noop()).is_pending());
4112
4113 let req = dynamic
4114 .requested_group()
4115 .now_or_never()
4116 .expect("should not block")
4117 .unwrap();
4118 assert_eq!(req.sequence(), 5);
4119 assert_eq!(req.priority(), 7);
4120 assert!(
4121 dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending(),
4122 "the second fetch queued a duplicate request"
4123 );
4124
4125 let third = consumer.fetch_group(5, None);
4127
4128 let mut group = req.accept(None).unwrap();
4130 group
4131 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hi"))
4132 .unwrap();
4133 group.finish().unwrap();
4134
4135 assert_eq!(first.await.unwrap().sequence, 5);
4136 assert_eq!(second.await.unwrap().sequence, 5);
4137 assert_eq!(third.await.unwrap().sequence, 5);
4138 }
4139
4140 #[tokio::test]
4141 async fn fetch_coalesced_reject_fails_all() {
4142 let producer = track_producer("test", None);
4143 let dynamic = producer.dynamic();
4144 let consumer = producer.consume();
4145
4146 let first = consumer.fetch_group(5, None);
4147 let second = consumer.fetch_group(5, None);
4148 let req = dynamic
4149 .requested_group()
4150 .now_or_never()
4151 .expect("should not block")
4152 .unwrap();
4153 req.reject(Error::Cancel);
4154
4155 assert!(matches!(first.await, Err(Error::Cancel)));
4156 assert!(matches!(second.await, Err(Error::Cancel)));
4157
4158 let retry = consumer.fetch_group(5, None);
4160 assert!(kio::Pollable::poll(&*retry, &kio::Waiter::noop()).is_pending());
4161 let req = dynamic
4162 .requested_group()
4163 .now_or_never()
4164 .expect("should not block")
4165 .unwrap();
4166 assert_eq!(req.sequence(), 5);
4167 }
4168
4169 #[tokio::test]
4170 async fn fetch_queued_fails_when_handlers_leave() {
4171 let producer = track_producer("test", None);
4172 let dynamic = producer.dynamic();
4173 let consumer = producer.consume();
4174
4175 let pending = consumer.fetch_group(5, None);
4177 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
4178 drop(dynamic);
4179 assert!(matches!(pending.await, Err(Error::NotFound)));
4180
4181 let fetch = producer.state.read().fetch.clone();
4183 assert!(fetch.read().is_empty());
4184 }
4185
4186 #[tokio::test]
4187 async fn fetch_miss_no_dynamic_not_found() {
4188 let mut producer = track_producer("test", None);
4191 producer.append_group().unwrap(); let consumer = producer.consume();
4193 assert!(matches!(consumer.fetch_group(5, None).await, Err(Error::NotFound)));
4194 }
4195
4196 #[tokio::test]
4197 async fn fetch_past_final_not_found() {
4198 let mut producer = track_producer("test", None);
4199 producer.append_group().unwrap(); producer.finish().unwrap(); let dynamic = producer.dynamic();
4205 let consumer = producer.consume();
4206 assert!(matches!(consumer.fetch_group(5, None).await, Err(Error::NotFound)));
4207
4208 assert!(dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending());
4210 }
4211
4212 fn pooled_producer(capacity: u64) -> (Producer, cache::Pool) {
4214 let pool = cache::Pool::new(capacity);
4215 let broadcast = broadcast::Info {
4216 origin: crate::origin::Info::default().with_pool(pool.clone()),
4217 ..Default::default()
4218 };
4219 let producer = Producer::new(Arc::new(broadcast), "test", None);
4220 (producer, pool)
4221 }
4222
4223 fn finished_group(producer: &mut Producer, size: usize) -> u64 {
4224 let mut group = producer.append_group().unwrap();
4225 group
4226 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; size]))
4227 .unwrap();
4228 group.finish().unwrap();
4229 group.sequence
4230 }
4231
4232 #[tokio::test]
4235 async fn debt_evicts_oldest_group() {
4236 tokio::time::pause();
4237
4238 let (mut producer, pool) = pooled_producer(10_000);
4240
4241 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let consumer = producer.consume();
4246 assert!(consumer.peek_group(0).is_none(), "oldest group is evicted");
4247 assert!(consumer.peek_group(2).is_some(), "latest group survives");
4248 assert!(pool.used() <= 21_000, "usage hovers near capacity: {}", pool.used());
4251
4252 let mut subscriber = producer.subscribe(None);
4254 assert!(subscriber.assert_group().sequence > 0, "evicted group is not delivered");
4255 }
4256
4257 #[tokio::test]
4259 async fn latest_group_never_evicted() {
4260 tokio::time::pause();
4261
4262 let (mut producer, pool) = pooled_producer(100);
4264 finished_group(&mut producer, 1000); assert!(pool.used() > 100, "the latest may exceed the budget");
4266
4267 finished_group(&mut producer, 1000); finished_group(&mut producer, 1000); let consumer = producer.consume();
4272 assert!(consumer.peek_group(0).is_none());
4273 let mut group = consumer.peek_group(2).expect("latest survives");
4274 assert_eq!(group.read_frame().await.unwrap().unwrap().payload.len(), 1000);
4275 }
4276
4277 #[tokio::test]
4281 async fn fetch_refresh_survives_eviction() {
4282 tokio::time::pause();
4283
4284 let (mut producer, _pool) = pooled_producer(10_000);
4285 let consumer = producer.consume();
4286
4287 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_secs(1)).await;
4289 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_secs(1)).await;
4291 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_millis(500)).await;
4293
4294 let mut fetched = consumer.fetch_group(0, None).await.unwrap();
4296 assert_eq!(fetched.read_frame().await.unwrap().unwrap().payload.len(), 3_000);
4297 tokio::time::advance(Duration::from_millis(500)).await;
4298
4299 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_secs(1)).await;
4303 finished_group(&mut producer, 3_000); assert!(consumer.peek_group(0).is_some(), "refreshed group survives");
4306 assert!(consumer.peek_group(1).is_none(), "unread group is evicted instead");
4307 }
4308
4309 #[tokio::test]
4312 async fn eviction_aborts_readers() {
4313 tokio::time::pause();
4314
4315 let (mut producer, _pool) = pooled_producer(10_000);
4316 let mut subscriber = producer.subscribe(None);
4317
4318 finished_group(&mut producer, 10_000); let mut group0 = subscriber.assert_group();
4320
4321 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let read = group0.read_frame().await;
4325 assert!(matches!(read, Err(Error::Evicted)), "expected Evicted, got {read:?}");
4326 }
4327
4328 #[tokio::test]
4332 async fn small_writes_carry_debt() {
4333 tokio::time::pause();
4334
4335 let (mut producer, pool) = pooled_producer(22_000);
4336 let consumer = producer.consume();
4337
4338 finished_group(&mut producer, 20_000); for _ in 0..3 {
4343 finished_group(&mut producer, 1_000);
4344 }
4345 assert!(consumer.peek_group(0).is_some(), "debt smaller than the victim carries");
4346
4347 for _ in 0..20 {
4349 finished_group(&mut producer, 1_000);
4350 }
4351 assert!(
4352 consumer.peek_group(0).is_none(),
4353 "accumulated debt evicts the large group"
4354 );
4355 assert!(pool.used() <= 24_000, "usage hovers near capacity: {}", pool.used());
4358 }
4359
4360 #[tokio::test]
4364 async fn payment_capped_per_write() {
4365 tokio::time::pause();
4366
4367 let (mut producer, pool) = pooled_producer(1 << 40);
4368 for _ in 0..10 {
4369 finished_group(&mut producer, 1_000);
4370 }
4371
4372 pool.resize(100);
4374 let before = pool.used();
4375
4376 finished_group(&mut producer, 1_000);
4378
4379 let consumer = producer.consume();
4380 assert!(consumer.peek_group(0).is_none(), "the oldest groups are evicted");
4381 assert!(consumer.peek_group(1).is_none());
4382 assert!(consumer.peek_group(2).is_some(), "the backlog drains gradually");
4383 assert!(pool.used() > before - 4_000, "one write must not dump the backlog");
4384 }
4385
4386 #[tokio::test]
4390 async fn accept_preserves_write_accounting() {
4391 tokio::time::pause();
4392
4393 let pool = cache::Pool::new(12_000);
4394 let broadcast = broadcast::Info {
4395 origin: crate::origin::Info::default().with_pool(pool.clone()),
4396 ..Default::default()
4397 };
4398 let request = Request::new(Arc::new(broadcast), "test");
4399 let dynamic = request.dynamic();
4400 let consumer = request.consume();
4401
4402 let pending = consumer.fetch_group(0, None);
4404 let req = dynamic
4405 .requested_group()
4406 .now_or_never()
4407 .expect("should not block")
4408 .unwrap();
4409 let mut backfill = req.accept(None).unwrap();
4410 pending.await.unwrap();
4411 backfill
4412 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 30_000]))
4413 .unwrap();
4414
4415 let mut producer = request.accept(None);
4418 producer.append_group().unwrap().finish().unwrap();
4419 producer.append_group().unwrap().finish().unwrap();
4420
4421 assert!(
4422 producer.consume().peek_group(0).is_none(),
4423 "pre-accept backfill growth is reclaimed after accept"
4424 );
4425 assert!(pool.used() <= 13_000, "usage converges: {}", pool.used());
4426 }
4427
4428 #[tokio::test]
4431 async fn recreated_sequence_bounds_eviction_hints() {
4432 let (mut producer, _pool) = pooled_producer(1 << 40);
4433 producer.create_group(5u64.into()).unwrap().finish().unwrap();
4434
4435 for _ in 0..200 {
4436 let group = producer.create_group(1u64.into()).unwrap();
4437 group.abort(Error::Cancel).unwrap();
4438 }
4439
4440 let state = producer.state.read();
4441 assert!(
4442 state.evict.len() <= 2 * state.lookup.len() + EVICT_SLACK,
4443 "stale hints are compacted: {} entries for {} slots",
4444 state.evict.len(),
4445 state.lookup.len()
4446 );
4447 }
4448
4449 #[tokio::test]
4452 async fn same_tick_write_outranks_inserted() {
4453 tokio::time::pause();
4454
4455 let (mut producer, _pool) = pooled_producer(10_000);
4457
4458 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();
4465 assert!(consumer.peek_group(0).is_none(), "insert-only content pays first");
4466 assert!(consumer.peek_group(1).is_some(), "same-tick written content survives");
4467 }
4468
4469 #[tokio::test]
4472 async fn frame_only_writer_pays() {
4473 tokio::time::pause();
4474
4475 let (mut producer, pool) = pooled_producer(2_000);
4476 let mut demoted = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); demoted
4482 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 300_000]))
4483 .unwrap();
4484
4485 assert!(
4486 pool.used() <= 5_000,
4487 "the frame write settled the debt: {}",
4488 pool.used()
4489 );
4490 assert!(matches!(demoted.finish(), Err(Error::Evicted)));
4491 }
4492
4493 #[tokio::test]
4496 async fn each_track_owns_its_account() {
4497 let broadcast = Arc::new(broadcast::Info::default());
4498 let info = Info::default();
4499 let a = Producer::new(broadcast.clone(), "a", info.clone());
4500 let b = Producer::new(broadcast, "b", info);
4501
4502 let a = a.state.read().cache.clone();
4503 let b = b.state.read().cache.clone();
4504 assert!(!Arc::ptr_eq(&a, &b), "each track owns its account");
4505 }
4506
4507 #[tokio::test]
4510 async fn a_dynamic_defers_teardown() {
4511 let (mut producer, pool) = pooled_producer(1 << 40);
4512 let dynamic = producer.dynamic();
4513 finished_group(&mut producer, 100);
4514
4515 drop(producer);
4516 assert!(pool.used() > 0, "the handler still serves the cache");
4517
4518 drop(dynamic);
4519 assert_eq!(pool.used(), 0, "the last handle tears it down");
4520 }
4521
4522 #[tokio::test]
4528 async fn finished_track_frees_its_cache() {
4529 let (mut producer, pool) = pooled_producer(1 << 40);
4530 finished_group(&mut producer, 100);
4531 producer.finish().unwrap();
4532
4533 let state = producer.state.downgrade();
4534 drop(producer);
4535
4536 assert!(state.upgrade().is_none(), "the track state is freed");
4537 assert_eq!(pool.used(), 0, "so are its cached bytes");
4538 }
4539
4540 #[tokio::test]
4544 async fn teardown_ignores_a_settling_group() {
4545 let (mut producer, pool) = pooled_producer(1 << 40);
4546 finished_group(&mut producer, 100);
4547
4548 let settling = producer.state.downgrade().upgrade().expect("open");
4550 drop(producer);
4551
4552 assert_eq!(pool.used(), 0, "the abrupt teardown still released the cache");
4553 drop(settling);
4554 }
4555
4556 #[tokio::test]
4559 async fn cached_group_outlives_its_track() {
4560 let (mut producer, pool) = pooled_producer(1 << 40);
4561 let sequence = finished_group(&mut producer, 100);
4562 let group = producer.consume().peek_group(sequence).expect("cached");
4563 producer.finish().unwrap();
4564
4565 let state = producer.state.downgrade();
4566 drop(producer);
4567 assert!(state.upgrade().is_none(), "the track state is freed");
4568 assert!(pool.used() > 0, "the retained group keeps its own bytes");
4569
4570 drop(group);
4571 assert_eq!(pool.used(), 0, "which it releases when dropped");
4572 }
4573
4574 #[tokio::test]
4578 async fn pre_accept_backfill_settles_late_writes() {
4579 tokio::time::pause();
4580
4581 let pool = cache::Pool::new(2_000);
4582 let broadcast = broadcast::Info {
4583 origin: crate::origin::Info::default().with_pool(pool.clone()),
4584 ..Default::default()
4585 };
4586 let request = Request::new(Arc::new(broadcast), "test");
4587 let dynamic = request.dynamic();
4588 let consumer = request.consume();
4589
4590 let pending = consumer.fetch_group(0, None);
4592 let req = dynamic
4593 .requested_group()
4594 .now_or_never()
4595 .expect("should not block")
4596 .unwrap();
4597 let mut backfill = req.accept(None).unwrap();
4598 pending.await.unwrap();
4599
4600 let mut producer = request.accept(None);
4602 producer.append_group().unwrap().finish().unwrap();
4603
4604 backfill
4607 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 300_000]))
4608 .unwrap();
4609
4610 assert!(
4611 pool.used() <= 5_000,
4612 "the frame write settled the debt: {}",
4613 pool.used()
4614 );
4615 }
4616
4617 #[tokio::test]
4621 async fn write_restarts_retention_clock() {
4622 tokio::time::pause();
4623
4624 let (mut producer, _pool) = pooled_producer(1 << 40);
4625 let mut straggler = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
4630 straggler
4631 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 100]))
4632 .unwrap();
4633 producer.append_group().unwrap().finish().unwrap(); let consumer = producer.consume();
4636 assert!(consumer.peek_group(0).is_some(), "the write restarted the clock");
4637
4638 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
4640 producer.append_group().unwrap().finish().unwrap(); assert!(consumer.peek_group(0).is_none(), "idle content still expires");
4642 }
4643
4644 #[tokio::test]
4647 async fn refreshed_front_does_not_starve_expiry() {
4648 tokio::time::pause();
4649
4650 let (mut producer, _pool) = pooled_producer(1 << 40);
4651 let dynamic = producer.dynamic();
4652 let consumer = producer.consume();
4653
4654 producer.create_group(10u64.into()).unwrap().finish().unwrap();
4655 for sequence in 1..=5u64 {
4656 let pending = consumer.fetch_group(sequence, 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(vec![0u8; 100]))
4665 .unwrap();
4666 group.finish().unwrap();
4667 pending.await.unwrap();
4668 }
4669
4670 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
4673 for sequence in 1..=4u64 {
4674 consumer.fetch_group(sequence, None).await.unwrap();
4675 }
4676
4677 for _ in 0..3 {
4679 producer.append_group().unwrap().finish().unwrap();
4680 }
4681 assert!(consumer.peek_group(5).is_none(), "expired backfill is reclaimed");
4682 assert!(consumer.peek_group(1).is_some(), "refreshed backfill survives");
4683 }
4684
4685 #[tokio::test]
4688 async fn recreated_sequence_delivered_once() {
4689 let (mut producer, _pool) = pooled_producer(1 << 40);
4690
4691 producer.create_group(0u64.into()).unwrap().finish().unwrap();
4692 let aborted = producer.create_group(1u64.into()).unwrap();
4693 aborted.abort(Error::Cancel).unwrap();
4694 producer.create_group(2u64.into()).unwrap().finish().unwrap();
4695 producer.create_group(1u64.into()).unwrap().finish().unwrap();
4696
4697 let mut subscriber = producer.subscribe(None);
4698 assert_eq!(subscriber.assert_group().sequence, 0);
4699 assert_eq!(subscriber.assert_group().sequence, 2);
4700 assert_eq!(
4701 subscriber.assert_group().sequence,
4702 1,
4703 "replacement arrives at its own position"
4704 );
4705 subscriber.assert_no_group();
4706 }
4707
4708 #[tokio::test]
4712 async fn datagrams_do_not_block_eviction() {
4713 tokio::time::pause();
4714
4715 let (mut producer, pool) = pooled_producer(1_000);
4716 for _ in 0..10 {
4717 finished_group(&mut producer, 1_000);
4718 producer.append_datagram(Timestamp::ZERO, &b"beat"[..]).unwrap();
4719 }
4720
4721 let consumer = producer.consume();
4722 assert!(consumer.peek_group(0).is_none(), "old groups still evict");
4723 assert!(
4724 pool.used() < 4 * 1_256,
4725 "interleaved datagrams must not bypass the budget: {}",
4726 pool.used()
4727 );
4728 }
4729
4730 #[tokio::test]
4734 async fn aborted_group_leaves_no_ghost_sample() {
4735 tokio::time::pause();
4736
4737 let (mut producer, pool) = pooled_producer(1 << 40);
4738 let group0 = producer.append_group().unwrap();
4739 producer.append_group().unwrap(); assert!(pool.average().is_some(), "demoted group is sampled");
4742 group0.abort(Error::Cancel).unwrap();
4743 assert_eq!(pool.average(), None, "the abort must remove the sample");
4744 }
4745
4746 #[tokio::test]
4749 async fn empty_groups_repay_overhead() {
4750 tokio::time::pause();
4751
4752 let (mut producer, pool) = pooled_producer(1_000);
4753 for _ in 0..100 {
4754 let mut group = producer.append_group().unwrap();
4755 group.finish().unwrap();
4756 }
4757
4758 assert!(
4759 pool.used() <= 3_000,
4760 "empty-group overhead must stay near the budget: {}",
4761 pool.used()
4762 );
4763 }
4764
4765 #[tokio::test]
4768 async fn growth_on_demoted_group_is_billed() {
4769 tokio::time::pause();
4770
4771 let (mut producer, pool) = pooled_producer(2_000);
4772 let mut straggler = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); straggler
4777 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 10_000]))
4778 .unwrap();
4779
4780 producer.append_group().unwrap().finish().unwrap(); let consumer = producer.consume();
4784 assert!(consumer.peek_group(0).is_none(), "the ballooned group is evicted");
4785 assert!(pool.used() <= 3_000, "growth is reclaimed: {}", pool.used());
4786 }
4787
4788 #[tokio::test]
4791 async fn refilled_sequence_stays_out_of_subscriptions() {
4792 let (mut producer, _pool) = pooled_producer(1 << 40);
4793 let dynamic = producer.dynamic();
4794 let consumer = producer.consume();
4795
4796 producer.create_group(0u64.into()).unwrap().finish().unwrap();
4797 let aborted = producer.create_group(1u64.into()).unwrap();
4798 aborted.abort(Error::Cancel).unwrap();
4799 producer.create_group(2u64.into()).unwrap().finish().unwrap();
4800
4801 let pending = consumer.fetch_group(1, None);
4804 let req = dynamic
4805 .requested_group()
4806 .now_or_never()
4807 .expect("should not block")
4808 .unwrap();
4809 let mut group = req.accept(None).unwrap();
4810 group
4811 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"backfill"))
4812 .unwrap();
4813 group.finish().unwrap();
4814 pending.await.unwrap();
4815
4816 assert!(consumer.peek_group(1).is_some());
4818 let mut subscriber = producer.subscribe(None);
4819 assert_eq!(subscriber.assert_group().sequence, 0);
4820 assert_eq!(subscriber.assert_group().sequence, 2);
4821 subscriber.assert_no_group();
4822 }
4823
4824 #[tokio::test]
4827 async fn expired_backfill_behind_refreshed_reclaimed() {
4828 tokio::time::pause();
4829
4830 let (mut producer, _pool) = pooled_producer(1 << 40);
4831 let dynamic = producer.dynamic();
4832 let consumer = producer.consume();
4833
4834 producer.create_group(5u64.into()).unwrap().finish().unwrap();
4835 for sequence in [2u64, 3u64] {
4836 let pending = consumer.fetch_group(sequence, None);
4837 let req = dynamic
4838 .requested_group()
4839 .now_or_never()
4840 .expect("should not block")
4841 .unwrap();
4842 let mut group = req.accept(None).unwrap();
4843 group
4844 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 100]))
4845 .unwrap();
4846 group.finish().unwrap();
4847 pending.await.unwrap();
4848 }
4849
4850 tokio::time::advance(Duration::from_secs(4)).await;
4852 consumer.fetch_group(2, None).await.unwrap();
4853 tokio::time::advance(DEFAULT_LATENCY_MAX - Duration::from_secs(2)).await;
4854 producer.create_group(6u64.into()).unwrap().finish().unwrap();
4855
4856 let consumer = producer.consume();
4857 assert!(consumer.peek_group(2).is_some(), "refreshed backfill survives");
4858 assert!(consumer.peek_group(3).is_none(), "expired backfill is reclaimed");
4859 }
4860
4861 #[tokio::test]
4864 async fn same_tick_fetch_protects() {
4865 tokio::time::pause();
4866
4867 let (mut producer, _pool) = pooled_producer(10_000);
4869 let consumer = producer.consume();
4870
4871 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();
4876
4877 finished_group(&mut producer, 3_000); finished_group(&mut producer, 3_000); assert!(consumer.peek_group(0).is_some(), "same-tick refresh protects");
4881 assert!(consumer.peek_group(1).is_none(), "the unread group dies instead");
4882 }
4883
4884 #[tokio::test]
4888 async fn refetched_latest_stays_protected() {
4889 tokio::time::pause();
4890
4891 let (mut producer, _pool) = pooled_producer(10_000);
4892 let dynamic = producer.dynamic();
4893 let consumer = producer.consume();
4894
4895 let straggler = producer.append_group().unwrap(); let latest = producer.append_group().unwrap(); latest.abort(Error::Cancel).unwrap();
4900
4901 let pending = consumer.fetch_group(1, None);
4903 let req = dynamic
4904 .requested_group()
4905 .now_or_never()
4906 .expect("should not block")
4907 .unwrap();
4908 let mut group = req.accept(None).unwrap();
4909 group
4910 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 1000]))
4911 .unwrap();
4912 group.finish().unwrap();
4913 pending.await.unwrap();
4914
4915 {
4918 let state = producer.state.read();
4919 assert!(state.lookup.contains_key(&1), "refetched group is cached");
4920 assert!(
4921 state.evict.iter().all(|(sequence, _)| *sequence != 1),
4922 "the live edge must not be an eviction candidate"
4923 );
4924 }
4925 drop(straggler);
4926 }
4927
4928 #[tokio::test]
4931 async fn eviction_allows_refetch() {
4932 tokio::time::pause();
4933
4934 let (mut producer, _pool) = pooled_producer(10_000);
4935 let dynamic = producer.dynamic();
4936
4937 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let consumer = producer.consume();
4942 assert!(consumer.peek_group(0).is_none());
4943 let pending = consumer.fetch_group(0, None);
4944
4945 let req = dynamic
4946 .requested_group()
4947 .now_or_never()
4948 .expect("should not block")
4949 .unwrap();
4950 assert_eq!(req.sequence(), 0);
4951
4952 let mut group = req.accept(None).unwrap();
4953 group
4954 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"refetched"))
4955 .unwrap();
4956 group.finish().unwrap();
4957
4958 let mut group = pending.await.unwrap();
4959 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"refetched");
4960 }
4961
4962 #[tokio::test]
4965 async fn fetched_backfill_not_subscribed() {
4966 let (mut producer, _pool) = pooled_producer(1 << 40);
4967 let dynamic = producer.dynamic();
4968 let consumer = producer.consume();
4969
4970 producer.create_group(5u64.into()).unwrap().finish().unwrap();
4972 producer.create_group(6u64.into()).unwrap().finish().unwrap();
4973
4974 let pending = consumer.fetch_group(2, None);
4976 let req = dynamic
4977 .requested_group()
4978 .now_or_never()
4979 .expect("should not block")
4980 .unwrap();
4981 let mut group = req.accept(None).unwrap();
4982 group
4983 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"backfill"))
4984 .unwrap();
4985 group.finish().unwrap();
4986 let mut fetched = pending.await.unwrap();
4987 assert_eq!(&fetched.read_frame().await.unwrap().unwrap().payload[..], b"backfill");
4988 assert!(consumer.peek_group(2).is_some(), "backfill is cached for later fetches");
4989
4990 let mut subscriber = producer.subscribe(None);
4992 assert_eq!(subscriber.assert_group().sequence, 5);
4993 assert_eq!(subscriber.assert_group().sequence, 6);
4994 subscriber.assert_no_group();
4995 }
4996
4997 #[tokio::test]
5000 async fn expired_backfill_reclaimed() {
5001 tokio::time::pause();
5002
5003 let (mut producer, pool) = pooled_producer(1 << 40);
5004 let dynamic = producer.dynamic();
5005 let consumer = producer.consume();
5006
5007 producer.create_group(5u64.into()).unwrap().finish().unwrap();
5008
5009 let pending = consumer.fetch_group(2, None);
5011 let req = dynamic
5012 .requested_group()
5013 .now_or_never()
5014 .expect("should not block")
5015 .unwrap();
5016 let mut group = req.accept(None).unwrap();
5017 group
5018 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 1000]))
5019 .unwrap();
5020 group.finish().unwrap();
5021 pending.await.unwrap();
5022 let used = pool.used();
5023
5024 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
5026 producer.create_group(6u64.into()).unwrap().finish().unwrap();
5027
5028 assert!(consumer.peek_group(2).is_none(), "expired backfill is reclaimed");
5029 assert!(pool.used() < used, "its bytes are released");
5030 }
5031
5032 #[tokio::test]
5033 async fn fetch_aborts_with_track() {
5034 let producer = track_producer("test", None);
5035 let dynamic = producer.dynamic();
5036 let consumer = producer.consume();
5037
5038 let pending = consumer.fetch_group(3, None);
5039 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
5040
5041 producer.abort(Error::Cancel).unwrap();
5042 assert!(pending.await.is_err());
5043 drop(dynamic);
5044 }
5045}