1use crate::{Error, Result, Timescale, Timestamp, coding};
17use crate::{broadcast, cache, frame, group, stats};
18
19use super::{Datagram, Requests};
20
21pub use super::subscription::Subscription;
22
23use std::{
24 collections::{BTreeMap, 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 parked: BTreeMap::new(),
1200 }),
1201 stats: stats::Scope::default(),
1203 _stats_sub: stats::Subscription::default(),
1204 }
1205 }
1206
1207 pub async fn subscription_changed(&mut self) -> Result<Option<Subscription>> {
1213 kio::wait(|waiter| self.poll_subscription_changed(waiter)).await
1214 }
1215
1216 pub fn subscription(&self) -> Option<Subscription> {
1224 let state = self.state.read();
1225 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
1226 drop(state);
1227 snapshot_subscription(&subs, bound)
1228 }
1229
1230 pub fn poll_subscription_changed(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Subscription>>> {
1234 if self.state.poll_closed(waiter).is_ready() {
1237 let abort = self.state.read().abort.clone();
1238 return Poll::Ready(Err(abort.unwrap_or(Error::Dropped)));
1239 }
1240
1241 let state = self.state.read();
1243 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
1244 drop(state);
1245
1246 let prev = &self.prev_subscription;
1247 let mut combined = None;
1248 let mut guard = match subs.poll(waiter, |subs| {
1249 let next = combined_subscription(subs, bound, waiter);
1250 if &next == prev {
1251 Poll::Pending
1252 } else {
1253 combined = next;
1254 Poll::Ready(())
1255 }
1256 }) {
1257 Poll::Ready(guard) => guard,
1258 Poll::Pending => return Poll::Pending,
1259 };
1260 guard.retain(|sub| !sub.is_closed());
1262 drop(guard);
1263 self.prev_subscription = combined.clone();
1264 Poll::Ready(Ok(combined))
1265 }
1266
1267 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
1269 self.state.poll_unused(waiter).map(|_| ())
1270 }
1271
1272 pub fn dynamic(&self) -> Dynamic {
1276 Dynamic::new(self.name.clone(), self.state.clone(), self.alive.clone())
1277 }
1278
1279 fn modify(&self) -> Result<kio::Mut<'_, TrackState>> {
1280 TrackState::modify(&self.state)
1281 }
1282}
1283
1284fn poll_requested_group(
1288 state: &kio::Producer<TrackState>,
1289 fetch: &kio::Shared<FetchState>,
1290 waiter: &kio::Waiter,
1291) -> Poll<Result<GroupRequest>> {
1292 if let Poll::Ready(mut guard) = fetch.poll(waiter, |fetch| {
1294 if fetch.has_queued() {
1295 Poll::Ready(())
1296 } else {
1297 Poll::Pending
1298 }
1299 }) {
1300 let sequence = guard.pop().expect("predicate guaranteed a request");
1301 let pending = guard.get(&sequence).expect("popped key must be pending");
1305 let priority = pending.priority;
1306 let result = pending.result.clone();
1307 drop(guard);
1308 return Poll::Ready(Ok(GroupRequest {
1309 state: state.clone(),
1310 fetch: fetch.clone(),
1311 sequence,
1312 priority,
1313 result,
1314 done: false,
1315 }));
1316 }
1317
1318 match state.poll_ref(waiter, |state| match &state.abort {
1320 Some(err) => Poll::Ready(err.clone()),
1321 None => Poll::Pending,
1322 }) {
1323 Poll::Ready(Ok(err)) => Poll::Ready(Err(err)),
1324 Poll::Ready(Err(closed)) => Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped))),
1325 Poll::Pending => Poll::Pending,
1326 }
1327}
1328
1329pub struct Dynamic {
1339 name: Arc<str>,
1340 state: kio::Producer<TrackState>,
1342 fetch: kio::Shared<FetchState>,
1344 alive: Arc<Alive>,
1347}
1348
1349impl Dynamic {
1350 fn new(name: Arc<str>, state: kio::Producer<TrackState>, alive: Arc<Alive>) -> Self {
1351 let fetch = state.read().fetch.clone();
1352 fetch.lock().add_handler();
1353 Self {
1354 name,
1355 state,
1356 fetch,
1357 alive,
1358 }
1359 }
1360
1361 pub fn name(&self) -> &str {
1363 &self.name
1364 }
1365
1366 pub async fn requested_group(&self) -> Result<GroupRequest> {
1372 kio::wait(|waiter| self.poll_requested_group(waiter)).await
1373 }
1374
1375 pub fn poll_requested_group(&self, waiter: &kio::Waiter) -> Poll<Result<GroupRequest>> {
1377 poll_requested_group(&self.state, &self.fetch, waiter)
1378 }
1379
1380 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
1382 self.state.poll_unused(waiter).map(|_| ())
1383 }
1384}
1385
1386impl Clone for Dynamic {
1387 fn clone(&self) -> Self {
1388 self.fetch.lock().add_handler();
1390 Self {
1391 name: self.name.clone(),
1392 state: self.state.clone(),
1393 fetch: self.fetch.clone(),
1394 alive: self.alive.clone(),
1395 }
1396 }
1397}
1398
1399impl Drop for Dynamic {
1400 fn drop(&mut self) {
1401 let mut fetch = self.fetch.lock();
1407 if fetch.remove_handler() {
1408 fetch.drain_queued();
1409 }
1410 }
1411}
1412
1413struct Alive {
1422 name: Arc<str>,
1423 state: kio::Producer<TrackState>,
1424
1425 published: AtomicBool,
1428
1429 stats: OnceLock<stats::Subscription>,
1432}
1433
1434impl Alive {
1435 fn new(name: Arc<str>, state: kio::Producer<TrackState>) -> Arc<Self> {
1436 Arc::new(Self {
1437 name,
1438 state,
1439 published: Default::default(),
1440 stats: Default::default(),
1441 })
1442 }
1443
1444 fn publish(&self, stats: Option<&stats::Scope>) {
1448 self.published.store(true, Ordering::Relaxed);
1449 if let Some(scope) = stats {
1450 let _ = self.stats.set(scope.subscribe());
1453 }
1454 }
1455}
1456
1457impl Drop for Alive {
1458 fn drop(&mut self) {
1459 if !self.published.load(Ordering::Relaxed) {
1461 return;
1462 }
1463 match self.state.write() {
1471 Ok(mut state) => {
1472 if state.final_sequence.is_some() || state.abort.is_some() {
1473 return;
1474 }
1475 tracing::warn!(
1476 track = %self.name,
1477 "track::Producer dropped without finish() or abort()"
1478 );
1479 state.clear_cache();
1480 state.datagrams.clear();
1481 }
1482 Err(state) => {
1483 if state.final_sequence.is_some() || state.abort.is_some() {
1484 return;
1485 }
1486 tracing::warn!(
1487 track = %self.name,
1488 "track::Producer dropped without finish() or abort()"
1489 );
1490 }
1491 }
1492 }
1493}
1494
1495fn combined_subscription(subs: &Subscriptions, bound: Option<Duration>, waiter: &kio::Waiter) -> Option<Subscription> {
1501 let mut combined = None;
1502 for sub in subs.iter() {
1503 if sub.is_closed() {
1508 continue;
1509 }
1510 let _ = sub.poll_closed(waiter);
1515 if let Poll::Ready(Ok(sub)) = sub.poll(waiter, |sub| sub.poll_combined(&combined)) {
1516 combined = Some(sub);
1517 }
1518 }
1519 clamp_combined(combined, bound)
1520}
1521
1522fn snapshot_subscription(subs: &kio::Shared<Subscriptions>, bound: Option<Duration>) -> Option<Subscription> {
1524 let mut combined: Option<Subscription> = None;
1525 for sub in subs.read().iter() {
1526 if sub.is_closed() {
1528 continue;
1529 }
1530 if let Poll::Ready(merged) = sub.read().poll_combined(&combined) {
1531 combined = Some(merged);
1532 }
1533 }
1534 clamp_combined(combined, bound)
1535}
1536
1537fn clamp_combined(combined: Option<Subscription>, bound: Option<Duration>) -> Option<Subscription> {
1545 let mut combined = combined?;
1546 if let Some(bound) = bound {
1547 combined.latency_max = combined.latency_max.min(bound);
1548 }
1549 Some(combined)
1550}
1551
1552fn register_subscription(state: kio::Ref<'_, TrackState>, subscription: &kio::Producer<Subscription>) {
1556 if state.is_closed() {
1557 return;
1558 }
1559 let subs = state.subscriptions.clone();
1560 drop(state);
1561 subs.lock().push(subscription.consume());
1562}
1563
1564#[derive(Clone)]
1566pub(crate) struct TrackWeak {
1567 name: Arc<str>,
1568 state: kio::ProducerWeak<TrackState>,
1569}
1570
1571impl TrackWeak {
1572 pub fn consume(&self) -> Consumer {
1573 Consumer::plain(self.name.clone(), self.state.consume())
1574 }
1575
1576 pub(crate) fn name(&self) -> &Arc<str> {
1579 &self.name
1580 }
1581
1582 pub(crate) fn is_used(&self) -> bool {
1585 !self.state.is_closed() && self.state.is_used()
1586 }
1587
1588 pub(crate) fn poll_used(&self, waiter: &kio::Waiter) {
1591 let _ = self.state.poll_used(waiter);
1592 }
1593
1594 pub(crate) fn poll_unused(&self, waiter: &kio::Waiter) {
1597 let _ = self.state.poll_unused(waiter);
1598 }
1599}
1600
1601impl super::WeakEntry for TrackWeak {
1602 fn is_closed(&self) -> bool {
1603 self.state.is_closed()
1604 }
1605
1606 fn same_channel(&self, other: &Self) -> bool {
1607 self.state.same_channel(&other.state)
1608 }
1609}
1610
1611#[derive(Clone)]
1620pub struct Demand {
1621 name: Arc<str>,
1622 state: kio::ProducerWeak<TrackState>,
1623}
1624
1625impl Demand {
1626 pub fn name(&self) -> &str {
1628 &self.name
1629 }
1630
1631 pub async fn used(&self) -> Result<()> {
1633 self.state.used().await.map_err(|_| self.abort_reason())
1634 }
1635
1636 pub async fn unused(&self) -> Result<()> {
1638 self.state.unused().await.map_err(|_| self.abort_reason())
1639 }
1640
1641 pub async fn closed(&self) -> Error {
1643 self.state.closed().await;
1644 self.abort_reason()
1645 }
1646
1647 fn abort_reason(&self) -> Error {
1649 self.state.read().abort.clone().unwrap_or(Error::Dropped)
1650 }
1651}
1652
1653#[derive(Clone)]
1664pub struct Consumer {
1665 name: Arc<str>,
1666 inner: ConsumerKind,
1667 stats: stats::Scope,
1670}
1671
1672#[derive(Clone)]
1673enum ConsumerKind {
1674 Plain(kio::Consumer<TrackState>),
1675 Spliced(super::resume::Consumer),
1676}
1677
1678impl Consumer {
1679 fn plain(name: Arc<str>, state: kio::Consumer<TrackState>) -> Self {
1680 Self {
1681 name,
1682 inner: ConsumerKind::Plain(state),
1683 stats: stats::Scope::default(),
1684 }
1685 }
1686
1687 pub(crate) fn spliced(name: Arc<str>, resume: super::resume::Consumer) -> Self {
1689 Self {
1690 name,
1691 inner: ConsumerKind::Spliced(resume),
1692 stats: stats::Scope::default(),
1693 }
1694 }
1695
1696 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
1699 self.stats = scope;
1700 self
1701 }
1702
1703 pub fn name(&self) -> &str {
1705 &self.name
1706 }
1707
1708 pub fn subscribe(&self, subscription: impl Into<Option<Subscription>>) -> kio::Pending<Subscribing> {
1714 let subscription = kio::Producer::new(subscription.into().unwrap_or_default());
1715
1716 let inner = match &self.inner {
1717 ConsumerKind::Plain(state) => {
1718 register_subscription(state.read(), &subscription);
1721 SubscribingKind::Plain(state.clone())
1722 }
1723 ConsumerKind::Spliced(resume) => SubscribingKind::Spliced(resume.clone()),
1725 };
1726
1727 kio::Pending::new(Subscribing {
1728 name: self.name.clone(),
1729 inner,
1730 subscription,
1731 stats: self.stats.clone(),
1732 })
1733 }
1734
1735 #[cfg(test)]
1739 pub(crate) fn peek_group(&self, sequence: u64) -> Option<group::Consumer> {
1740 match &self.inner {
1741 ConsumerKind::Plain(state) => state.read().cached_group(sequence),
1742 ConsumerKind::Spliced(_) => None,
1745 }
1746 }
1747
1748 pub fn fetch_group(&self, sequence: u64, options: impl Into<Option<group::Fetch>>) -> kio::Pending<Fetching> {
1760 let options = options.into().unwrap_or_default();
1761
1762 self.stats.fetch();
1766
1767 let state = match &self.inner {
1768 ConsumerKind::Plain(state) => state,
1769 ConsumerKind::Spliced(resume) => {
1772 return kio::Pending::new(Fetching {
1773 inner: FetchingKind::Spliced(resume.fetch_group(sequence, options)),
1774 stats: self.stats.clone(),
1775 });
1776 }
1777 };
1778
1779 let mut result = None;
1780
1781 let (fetch, unresolved) = {
1785 let state = state.read();
1786 (state.fetch.clone(), state.poll_fetch_cached(sequence).is_pending())
1787 };
1788
1789 if unresolved {
1790 let mut fetch = fetch.lock();
1791 if let Some(pending) = fetch.join(&sequence) {
1792 pending.priority = pending.priority.max(options.priority);
1795 result = Some(pending.result.consume());
1796 } else {
1797 let producer = kio::Producer::<FetchOutcome>::default();
1801 let consumer = producer.consume();
1802 let attempt = PendingFetch {
1803 priority: options.priority,
1804 result: producer,
1805 };
1806 if fetch.insert(sequence, attempt).is_ok() {
1807 result = Some(consumer);
1808 }
1809 }
1810 }
1811
1812 kio::Pending::new(Fetching {
1813 inner: FetchingKind::Plain {
1814 state: state.clone(),
1815 fetch,
1816 sequence,
1817 result,
1818 },
1819 stats: self.stats.clone(),
1820 })
1821 }
1822
1823 pub fn info(&self) -> kio::Pending<Querying> {
1830 kio::Pending::new(Querying {
1831 inner: match &self.inner {
1832 ConsumerKind::Plain(state) => QueryingKind::Plain(state.clone()),
1833 ConsumerKind::Spliced(resume) => QueryingKind::Spliced(resume.clone()),
1834 },
1835 })
1836 }
1837
1838 pub fn latest(&self) -> Option<u64> {
1840 match &self.inner {
1841 ConsumerKind::Plain(state) => state.read().max_sequence,
1842 ConsumerKind::Spliced(resume) => resume.latest(),
1843 }
1844 }
1845
1846 pub(crate) fn poll_complete(&self, waiter: &kio::Waiter) -> Poll<Result<()>> {
1851 let ConsumerKind::Plain(state) = &self.inner else {
1852 return Poll::Pending;
1854 };
1855 match ready!(state.poll(waiter, |state| {
1856 if state.is_complete() {
1857 Poll::Ready(())
1858 } else {
1859 Poll::Pending
1860 }
1861 })) {
1862 Ok(_) => Poll::Ready(Ok(())),
1863 Err(closed) => Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped))),
1866 }
1867 }
1868}
1869
1870pub struct Subscribing {
1873 name: Arc<str>,
1874 inner: SubscribingKind,
1875 subscription: kio::Producer<Subscription>,
1876 stats: stats::Scope,
1877}
1878
1879enum SubscribingKind {
1880 Plain(kio::Consumer<TrackState>),
1881 Spliced(super::resume::Consumer),
1882}
1883
1884impl Subscribing {
1885 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<Subscriber>> {
1888 match &self.inner {
1889 SubscribingKind::Plain(state) => {
1890 let info = ready!(state.poll(waiter, |state| state.poll_info()))
1892 .map_err(|e| e.abort.clone().unwrap_or(Error::Dropped))??;
1893
1894 Poll::Ready(Ok(Subscriber {
1895 name: self.name.clone(),
1896 info,
1897 inner: SubscriberKind::Plain(PlainSubscriber {
1898 state: state.clone(),
1899 subscription: self.subscription.clone(),
1900 index: 0,
1901 datagram_index: 0,
1902 min_sequence: 0,
1903 next_sequence: 0,
1904 end_sequence: None,
1905 parked: BTreeMap::new(),
1906 }),
1907 stats: self.stats.clone(),
1908 _stats_sub: self.stats.subscribe(),
1909 }))
1910 }
1911 SubscribingKind::Spliced(resume) => {
1912 let info = ready!(resume.poll_info(waiter))?;
1915
1916 Poll::Ready(Ok(Subscriber {
1917 name: self.name.clone(),
1918 info,
1919 inner: SubscriberKind::Spliced(Box::new(resume.subscribe_shared(self.subscription.clone()))),
1920 stats: self.stats.clone(),
1921 _stats_sub: self.stats.subscribe(),
1922 }))
1923 }
1924 }
1925 }
1926
1927 pub fn update(&mut self, subscription: Subscription) -> Result<()> {
1932 let mut state = self.subscription.write().map_err(|_| Error::Closed)?;
1933 *state = subscription;
1934 Ok(())
1935 }
1936}
1937
1938impl kio::Pollable for Subscribing {
1939 type Output = Result<Subscriber>;
1940
1941 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
1942 self.poll_ok(waiter)
1943 }
1944}
1945
1946pub struct Querying {
1949 inner: QueryingKind,
1950}
1951
1952enum QueryingKind {
1953 Plain(kio::Consumer<TrackState>),
1954 Spliced(super::resume::Consumer),
1955}
1956
1957impl Querying {
1958 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<Info>> {
1960 match &self.inner {
1961 QueryingKind::Plain(state) => {
1962 let info = ready!(state.poll(waiter, |state| state.poll_info()))
1964 .map_err(|e| e.abort.clone().unwrap_or(Error::Dropped))??;
1965 Poll::Ready(Ok(info))
1966 }
1967 QueryingKind::Spliced(resume) => resume.poll_info(waiter),
1968 }
1969 }
1970}
1971
1972impl kio::Pollable for Querying {
1973 type Output = Result<Info>;
1974
1975 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
1976 self.poll_ok(waiter)
1977 }
1978}
1979
1980pub struct GroupRequest {
1989 state: kio::Producer<TrackState>,
1990 fetch: kio::Shared<FetchState>,
1992 sequence: u64,
1993 priority: u8,
1994 result: kio::Producer<FetchOutcome>,
1996 done: bool,
1997}
1998
1999impl GroupRequest {
2000 pub fn sequence(&self) -> u64 {
2002 self.sequence
2003 }
2004
2005 pub fn priority(&self) -> u8 {
2007 self.priority
2008 }
2009
2010 pub fn accept(mut self, info: impl Into<Option<Info>>) -> Result<group::Producer> {
2018 self.done = true;
2019 let res = TrackState::modify(&self.state)
2023 .and_then(|mut state| state.insert_group_request(self.sequence, info.into()));
2024 self.remove();
2025 res
2026 }
2027
2028 pub fn reject(mut self, err: Error) {
2030 self.done = true;
2031 self.remove();
2034 if let Ok(mut outcome) = self.result.write() {
2035 outcome.rejected = Some(err);
2036 }
2037 }
2038
2039 fn remove(&self) {
2042 self.fetch
2043 .lock()
2044 .remove_if(&self.sequence, |pending| pending.result.same_channel(&self.result));
2045 }
2046}
2047
2048impl Drop for GroupRequest {
2049 fn drop(&mut self) {
2050 if self.done {
2051 return;
2052 }
2053 self.remove();
2054 if let Ok(mut outcome) = self.result.write() {
2055 outcome.rejected = Some(Error::Dropped);
2056 }
2057 }
2058}
2059
2060pub struct Fetching {
2066 inner: FetchingKind,
2067 stats: stats::Scope,
2070}
2071
2072enum FetchingKind {
2073 Plain {
2074 state: kio::Consumer<TrackState>,
2075 fetch: kio::Shared<FetchState>,
2076 sequence: u64,
2077 result: Option<kio::Consumer<FetchOutcome>>,
2079 },
2080 Spliced(kio::Pending<super::resume::Fetching>),
2082}
2083
2084impl kio::Pollable for Fetching {
2085 type Output = Result<group::Consumer>;
2086
2087 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
2088 let (state, fetch, sequence, result) = match &self.inner {
2089 FetchingKind::Plain {
2090 state,
2091 fetch,
2092 sequence,
2093 result,
2094 } => (state, fetch, *sequence, result.as_ref()),
2095 FetchingKind::Spliced(spliced) => {
2096 return kio::Pollable::poll(&**spliced, waiter)
2099 .map(|res| res.map(|group| group.with_meter(self.stats.meter())));
2100 }
2101 };
2102
2103 match state.poll(waiter, |state| state.poll_fetch_cached(sequence)) {
2106 Poll::Ready(Ok(res)) => return Poll::Ready(res.map(|group| group.with_meter(self.stats.meter()))),
2107 Poll::Ready(Err(closed)) => {
2108 return Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped)));
2109 }
2110 Poll::Pending => {}
2111 }
2112
2113 let Some(result) = result else {
2115 return match fetch.poll(waiter, |fetch| match fetch.has_handlers() {
2118 false => Poll::Ready(()),
2119 true => Poll::Pending,
2120 }) {
2121 Poll::Ready(_guard) => Poll::Ready(Err(Error::NotFound)),
2122 Poll::Pending => Poll::Pending,
2123 };
2124 };
2125
2126 match result.poll(waiter, |outcome| match &outcome.rejected {
2129 Some(err) => Poll::Ready(err.clone()),
2130 None => Poll::Pending,
2131 }) {
2132 Poll::Ready(Ok(err)) => Poll::Ready(Err(err)),
2133 Poll::Ready(Err(_closed)) => Poll::Ready(Err(Error::NotFound)),
2134 Poll::Pending => Poll::Pending,
2135 }
2136 }
2137}
2138
2139pub struct Subscriber {
2162 name: Arc<str>,
2163 info: Info,
2164 inner: SubscriberKind,
2165 stats: stats::Scope,
2168 _stats_sub: stats::Subscription,
2171}
2172
2173enum SubscriberKind {
2174 Plain(PlainSubscriber),
2175 Spliced(Box<super::resume::Subscriber>),
2177}
2178
2179struct PlainSubscriber {
2181 state: kio::Consumer<TrackState>,
2182
2183 subscription: kio::Producer<Subscription>,
2184 index: usize,
2186 datagram_index: usize,
2188 min_sequence: u64,
2190 next_sequence: u64,
2193 end_sequence: Option<u64>,
2198 parked: BTreeMap<u64, group::Consumer>,
2203}
2204
2205impl PlainSubscriber {
2206 fn poll<F, R>(&self, waiter: &kio::Waiter, f: F) -> Poll<Result<R>>
2208 where
2209 F: Fn(&kio::Ref<'_, TrackState>) -> Poll<Result<R>>,
2210 {
2211 Poll::Ready(match ready!(self.state.poll(waiter, f)) {
2212 Ok(res) => res,
2213 Err(state) => Err(state.abort.clone().unwrap_or(Error::Dropped)),
2215 })
2216 }
2217
2218 fn poll_recv_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2219 self.parked
2224 .retain(|sequence, group| *sequence >= self.min_sequence && !group.is_aborted());
2225
2226 if let Some(&sequence) = self.parked.keys().next()
2228 && self.end_sequence.is_none_or(|end| sequence <= end)
2229 {
2230 return Poll::Ready(Ok(self.parked.remove(&sequence)));
2231 }
2232
2233 loop {
2234 let Some((consumer, found_index)) =
2235 ready!(self.poll(waiter, |state| state.poll_recv_group(self.index, self.min_sequence))?)
2236 else {
2237 if self.parked.is_empty() {
2240 return Poll::Ready(Ok(None));
2241 }
2242 return Poll::Pending;
2243 };
2244 self.index = found_index + 1;
2245
2246 if self.end_sequence.is_some_and(|end| consumer.sequence > end) {
2249 self.parked.insert(consumer.sequence, consumer);
2250 continue;
2251 }
2252 return Poll::Ready(Ok(Some(consumer)));
2253 }
2254 }
2255
2256 fn poll_recv_datagram(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Datagram>>> {
2257 let Some((datagram, found_index)) =
2258 ready!(self.poll(waiter, |state| state.poll_recv_datagram(self.datagram_index))?)
2259 else {
2260 return Poll::Ready(Ok(None));
2261 };
2262
2263 self.datagram_index = found_index + 1;
2264 self.next_sequence = self.next_sequence.max(datagram.sequence.saturating_add(1));
2265 Poll::Ready(Ok(Some(datagram)))
2266 }
2267
2268 fn poll_next_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2269 let floor = self.next_sequence.max(self.min_sequence);
2270 let Some(group) = ready!(self.poll(waiter, |state| state.poll_next_in_range(floor, self.end_sequence))?) else {
2271 return Poll::Ready(Ok(None));
2272 };
2273 self.next_sequence = group.sequence.saturating_add(1);
2274 Poll::Ready(Ok(Some(group)))
2275 }
2276
2277 fn poll_read_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Frame>>> {
2278 let lower = self.min_sequence.max(self.next_sequence);
2279 let Some((frame, found_index, sequence)) =
2280 ready!(self.poll(waiter, |state| { state.poll_read_frame(self.index, lower, waiter) })?)
2281 else {
2282 return Poll::Ready(Ok(None));
2283 };
2284
2285 self.index = found_index + 1;
2286 self.next_sequence = sequence.saturating_add(1);
2287 Poll::Ready(Ok(Some(frame)))
2288 }
2289}
2290
2291#[derive(Clone)]
2297pub struct SubscriberControl {
2298 subscription: kio::Producer<Subscription>,
2299}
2300
2301impl SubscriberControl {
2302 pub fn subscription(&self) -> Subscription {
2304 self.subscription.read().clone()
2305 }
2306
2307 pub fn update(&self, subscription: Subscription) -> Result<()> {
2312 let mut state = self.subscription.write().map_err(|_| Error::Closed)?;
2313 *state = subscription;
2314 Ok(())
2315 }
2316}
2317
2318impl Subscriber {
2319 pub fn info(&self) -> &Info {
2324 &self.info
2325 }
2326
2327 pub fn name(&self) -> &str {
2329 &self.name
2330 }
2331
2332 pub fn control(&self) -> SubscriberControl {
2334 SubscriberControl {
2335 subscription: match &self.inner {
2336 SubscriberKind::Plain(plain) => plain.subscription.clone(),
2337 SubscriberKind::Spliced(spliced) => spliced.prefs(),
2338 },
2339 }
2340 }
2341
2342 pub fn poll_recv_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2359 let meter = self.stats.meter();
2360 let res = match &mut self.inner {
2361 SubscriberKind::Plain(plain) => plain.poll_recv_group(waiter),
2362 SubscriberKind::Spliced(spliced) => spliced.poll_recv_group(waiter),
2363 };
2364 res.map(|res| res.map(|group| group.map(|group| group.with_meter(meter))))
2365 }
2366
2367 pub async fn recv_group(&mut self) -> Result<Option<group::Consumer>> {
2374 kio::wait(|waiter| self.poll_recv_group(waiter)).await
2375 }
2376
2377 pub fn poll_recv_datagram(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Datagram>>> {
2388 let meter = self.stats.meter();
2389 let res = match &mut self.inner {
2390 SubscriberKind::Plain(plain) => plain.poll_recv_datagram(waiter),
2391 SubscriberKind::Spliced(spliced) => spliced.poll_recv_datagram(waiter),
2392 };
2393 if let Poll::Ready(Ok(Some(datagram))) = &res {
2396 meter.datagram(datagram.payload.len() as u64);
2397 }
2398 res
2399 }
2400
2401 pub async fn recv_datagram(&mut self) -> Result<Option<Datagram>> {
2408 kio::wait(|waiter| self.poll_recv_datagram(waiter)).await
2409 }
2410
2411 pub fn poll_next_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2420 let meter = self.stats.meter();
2421 let res = match &mut self.inner {
2422 SubscriberKind::Plain(plain) => plain.poll_next_group(waiter),
2423 SubscriberKind::Spliced(spliced) => spliced.poll_next_group(waiter),
2424 };
2425 res.map(|res| res.map(|group| group.map(|group| group.with_meter(meter))))
2426 }
2427
2428 pub async fn next_group(&mut self) -> Result<Option<group::Consumer>> {
2434 kio::wait(|waiter| self.poll_next_group(waiter)).await
2435 }
2436
2437 pub fn poll_read_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Frame>>> {
2441 let meter = self.stats.meter();
2442 let res = match &mut self.inner {
2443 SubscriberKind::Plain(plain) => plain.poll_read_frame(waiter),
2444 SubscriberKind::Spliced(spliced) => spliced.poll_read_frame(waiter),
2445 };
2446 if let Poll::Ready(Ok(Some(frame))) = &res {
2449 meter.group();
2450 meter.frames(1);
2451 meter.bytes(frame.payload.len() as u64);
2452 }
2453 res
2454 }
2455
2456 pub async fn read_frame(&mut self) -> Result<Option<frame::Frame>> {
2461 kio::wait(|waiter| self.poll_read_frame(waiter)).await
2462 }
2463
2464 pub fn is_clone(&self, other: &Self) -> bool {
2466 match (&self.inner, &other.inner) {
2467 (SubscriberKind::Plain(a), SubscriberKind::Plain(b)) => a.state.same_channel(&b.state),
2468 (SubscriberKind::Spliced(a), SubscriberKind::Spliced(b)) => a.is_clone(b),
2469 _ => false,
2470 }
2471 }
2472
2473 pub fn poll_finished(&mut self, waiter: &kio::Waiter) -> Poll<Result<u64>> {
2475 match &mut self.inner {
2476 SubscriberKind::Plain(plain) => plain.poll(waiter, |state| state.poll_finished()),
2477 SubscriberKind::Spliced(spliced) => spliced.poll_finished(waiter),
2478 }
2479 }
2480
2481 pub async fn finished(&mut self) -> Result<u64> {
2489 kio::wait(|waiter| self.poll_finished(waiter)).await
2490 }
2491
2492 pub fn start_at(&mut self, sequence: u64) {
2499 match &mut self.inner {
2500 SubscriberKind::Plain(plain) => plain.min_sequence = sequence,
2501 SubscriberKind::Spliced(spliced) => spliced.start_at(sequence),
2502 }
2503 }
2504
2505 pub fn end_at(&mut self, sequence: impl Into<Option<u64>>) {
2518 match &mut self.inner {
2519 SubscriberKind::Plain(plain) => plain.end_sequence = sequence.into(),
2520 SubscriberKind::Spliced(spliced) => spliced.end_at(sequence),
2521 }
2522 }
2523
2524 pub fn subscription(&self) -> Subscription {
2526 self.control().subscription()
2527 }
2528
2529 pub fn update(&mut self, subscription: Subscription) -> Result<()> {
2535 match &mut self.inner {
2536 SubscriberKind::Plain(plain) => {
2537 let mut state = plain.subscription.write().map_err(|_| Error::Closed)?;
2538 *state = subscription;
2539 }
2540 SubscriberKind::Spliced(spliced) => spliced.update(subscription),
2541 }
2542 Ok(())
2543 }
2544
2545 pub fn latest(&self) -> Option<u64> {
2547 match &self.inner {
2548 SubscriberKind::Plain(plain) => plain.state.read().max_sequence,
2549 SubscriberKind::Spliced(spliced) => spliced.latest(),
2550 }
2551 }
2552}
2553
2554pub struct Request {
2566 name: Arc<str>,
2567 broadcast: Arc<broadcast::Info>,
2569 state: kio::Producer<TrackState>,
2570
2571 prev_subscription: Option<Subscription>,
2573
2574 alive: Arc<Alive>,
2577
2578 _dynamic: Dynamic,
2583
2584 stats: stats::Scope,
2587}
2588
2589impl Request {
2590 pub(crate) fn new(broadcast: Arc<broadcast::Info>, name: impl Into<Arc<str>>) -> Self {
2591 let name = name.into();
2592 let state = TrackState::spawn(broadcast.clone());
2593 let alive = Alive::new(name.clone(), state.clone());
2594 let dynamic = Dynamic::new(name.clone(), state.clone(), alive.clone());
2595 Self {
2596 name,
2597 broadcast,
2598 state,
2599 prev_subscription: None,
2600 alive,
2601 _dynamic: dynamic,
2602 stats: stats::Scope::default(),
2603 }
2604 }
2605
2606 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
2609 self.stats = scope;
2610 self
2611 }
2612
2613 pub fn name(&self) -> &str {
2615 &self.name
2616 }
2617
2618 pub fn consume(&self) -> Consumer {
2620 Consumer::plain(self.name.clone(), self.state.consume())
2621 }
2622
2623 pub fn dynamic(&self) -> Dynamic {
2627 Dynamic::new(self.name.clone(), self.state.clone(), self.alive.clone())
2628 }
2629
2630 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
2633 self.state.poll_unused(waiter).map(|_| ())
2634 }
2635
2636 pub fn accept(self, info: impl Into<Option<Info>>) -> Producer {
2643 if let Ok(mut state) = self.state.write() {
2646 state.install(info.into().unwrap_or_default());
2647 }
2648 self.alive.publish(Some(&self.stats));
2651 Producer {
2652 name: self.name,
2653 broadcast: self.broadcast,
2654 state: self.state,
2655 prev_subscription: None,
2656 alive: self.alive,
2657 stats: self.stats,
2658 }
2659 }
2660
2661 pub fn reject(self, err: Error) {
2663 if let Ok(mut state) = self.state.write() {
2664 state.abort = Some(err);
2665 }
2666 }
2667
2668 pub fn subscription(&self) -> Option<Subscription> {
2671 let state = self.state.read();
2672 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
2673 drop(state);
2674 snapshot_subscription(&subs, bound)
2675 }
2676
2677 pub async fn subscription_changed(&mut self) -> Option<Subscription> {
2680 kio::wait(|waiter| self.poll_subscription_changed(waiter)).await
2681 }
2682
2683 pub fn poll_subscription_changed(&mut self, waiter: &kio::Waiter) -> Poll<Option<Subscription>> {
2685 let state = self.state.read();
2686 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
2687 drop(state);
2688
2689 let prev = &self.prev_subscription;
2690 let mut combined = None;
2691 let mut guard = ready!(subs.poll(waiter, |subs| {
2692 let next = combined_subscription(subs, bound, waiter);
2693 if &next == prev {
2694 Poll::Pending
2695 } else {
2696 combined = next;
2697 Poll::Ready(())
2698 }
2699 }));
2700 guard.retain(|sub| !sub.is_closed());
2702 drop(guard);
2703 self.prev_subscription = combined.clone();
2704 Poll::Ready(combined)
2705 }
2706
2707 pub(super) fn weak(&self) -> TrackWeak {
2708 TrackWeak {
2709 name: self.name.clone(),
2710 state: self.state.weak(),
2711 }
2712 }
2713}
2714
2715#[cfg(test)]
2716use futures::FutureExt;
2717
2718#[cfg(test)]
2719#[allow(missing_docs)] impl Subscriber {
2721 pub fn assert_group(&mut self) -> group::Consumer {
2722 self.recv_group()
2723 .now_or_never()
2724 .expect("group would have blocked")
2725 .expect("would have errored")
2726 .expect("track was closed")
2727 }
2728
2729 pub fn assert_no_group(&mut self) {
2730 assert!(
2731 self.recv_group().now_or_never().is_none(),
2732 "recv_group would not have blocked"
2733 );
2734 }
2735
2736 pub fn assert_not_closed(&mut self) {
2737 assert!(self.finished().now_or_never().is_none(), "should not be closed");
2738 }
2739
2740 pub fn assert_closed(&mut self) {
2741 assert!(self.finished().now_or_never().is_some(), "should be closed");
2742 }
2743
2744 pub fn assert_error(&mut self) {
2746 assert!(
2747 self.finished().now_or_never().expect("should not block").is_err(),
2748 "should be error"
2749 );
2750 }
2751
2752 pub fn assert_is_clone(&self, other: &Self) {
2753 assert!(self.is_clone(other), "should be clone");
2754 }
2755
2756 pub fn assert_not_clone(&self, other: &Self) {
2757 assert!(!self.is_clone(other), "should not be clone");
2758 }
2759}
2760
2761#[cfg(test)]
2762mod test {
2763 use super::*;
2764 use crate::model::test_tracing::count_drop_warnings;
2765
2766 fn track_producer(name: impl Into<Arc<str>>, info: impl Into<Option<Info>>) -> Producer {
2769 Producer::new(Arc::new(broadcast::Info::default()), name, info)
2770 }
2771
2772 fn live_groups(state: &TrackState) -> usize {
2774 state.lookup.len()
2775 }
2776
2777 fn first_live_sequence(state: &TrackState) -> u64 {
2779 state
2780 .arrival
2781 .iter()
2782 .find(|(sequence, stamp)| state.lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp))
2783 .map(|(sequence, _)| *sequence)
2784 .unwrap()
2785 }
2786
2787 fn recv_datagram(dg: &mut Subscriber) -> Datagram {
2789 dg.recv_datagram()
2790 .now_or_never()
2791 .expect("datagram would have blocked")
2792 .expect("would have errored")
2793 .expect("track was closed")
2794 }
2795
2796 #[tokio::test]
2797 async fn append_datagram_shares_group_sequence() {
2798 let mut producer = track_producer("test", None);
2799 let ts = Timestamp::from_millis(10).unwrap();
2800
2801 assert_eq!(producer.append_group().unwrap().sequence, 0);
2803 assert_eq!(producer.append_datagram(ts, &b"a"[..]).unwrap(), 1);
2804 assert_eq!(producer.append_group().unwrap().sequence, 2);
2805 assert_eq!(producer.append_datagram(ts, &b"b"[..]).unwrap(), 3);
2806 assert_eq!(producer.latest(), Some(3));
2807 }
2808
2809 #[tokio::test]
2810 async fn append_datagram_roundtrip() {
2811 let mut producer = track_producer("test", None);
2812 let mut dg = producer.subscribe(None);
2813
2814 let ts = Timestamp::from_millis(42).unwrap();
2815 let seq = producer.append_datagram(ts, &b"hello"[..]).unwrap();
2816
2817 let got = recv_datagram(&mut dg);
2818 assert_eq!(got.sequence, seq);
2819 assert_eq!(got.timestamp, ts);
2820 assert_eq!(&got.payload[..], b"hello");
2821 }
2822
2823 #[tokio::test]
2824 async fn write_datagram_preserves_sequence() {
2825 let mut producer = track_producer("test", None);
2826 let mut dg = producer.subscribe(None);
2827
2828 let ts = Timestamp::from_millis(5).unwrap();
2829 producer
2831 .write_datagram(Datagram {
2832 sequence: 100,
2833 timestamp: ts,
2834 payload: bytes::Bytes::from_static(b"x"),
2835 })
2836 .unwrap();
2837
2838 assert_eq!(recv_datagram(&mut dg).sequence, 100);
2839 assert_eq!(producer.append_group().unwrap().sequence, 101);
2841 }
2842
2843 #[tokio::test]
2844 async fn recv_datagram_advances_ordered_group_cursor() {
2845 let mut producer = track_producer("test", None);
2846 let mut subscriber = producer.subscribe(None);
2847 let ts = Timestamp::from_millis(5).unwrap();
2848
2849 producer
2850 .write_datagram(Datagram {
2851 sequence: 5,
2852 timestamp: ts,
2853 payload: bytes::Bytes::from_static(b"x"),
2854 })
2855 .unwrap();
2856 assert_eq!(recv_datagram(&mut subscriber).sequence, 5);
2857
2858 producer.create_group(group::Info { sequence: 3 }).unwrap();
2859 producer.create_group(group::Info { sequence: 6 }).unwrap();
2860
2861 let group = subscriber
2862 .next_group()
2863 .now_or_never()
2864 .expect("group would have blocked")
2865 .expect("would have errored")
2866 .expect("track was closed");
2867 assert_eq!(group.sequence, 6);
2868 }
2869
2870 #[tokio::test]
2871 async fn datagram_normalized_to_track_timescale() {
2872 let info = Info::default().with_timescale(Timescale::MICRO);
2873 let mut producer = track_producer("test", info);
2874 let mut dg = producer.subscribe(None);
2875
2876 producer
2878 .append_datagram(Timestamp::from_millis(2).unwrap(), &b"z"[..])
2879 .unwrap();
2880 let got = recv_datagram(&mut dg);
2881 assert_eq!(got.timestamp.scale(), Timescale::MICRO);
2882 assert_eq!(got.timestamp.value(), 2_000);
2883 }
2884
2885 #[tokio::test]
2886 async fn datagram_rejects_oversized() {
2887 let mut producer = track_producer("test", None);
2888 let big = bytes::Bytes::from(vec![0u8; crate::model::datagram::MAX_DATAGRAM_PAYLOAD + 1]);
2889 let ts = Timestamp::from_millis(0).unwrap();
2890 assert!(matches!(
2891 producer.append_datagram(ts, big.clone()),
2892 Err(Error::FrameTooLarge)
2893 ));
2894 assert!(matches!(
2895 producer.write_datagram(Datagram {
2896 sequence: 0,
2897 timestamp: ts,
2898 payload: big,
2899 }),
2900 Err(Error::FrameTooLarge)
2901 ));
2902 }
2903
2904 #[tokio::test]
2905 async fn datagram_fanout_to_subscribers() {
2906 let mut producer = track_producer("test", None);
2907 let mut a = producer.subscribe(None);
2909 let mut b = producer.subscribe(None);
2910 let ts = Timestamp::from_millis(1).unwrap();
2911
2912 producer.append_datagram(ts, &b"first"[..]).unwrap();
2913 producer.append_datagram(ts, &b"second"[..]).unwrap();
2914
2915 assert_eq!(&recv_datagram(&mut a).payload[..], b"first");
2917 assert_eq!(&recv_datagram(&mut a).payload[..], b"second");
2918 assert_eq!(&recv_datagram(&mut b).payload[..], b"first");
2919 assert_eq!(&recv_datagram(&mut b).payload[..], b"second");
2920 }
2921
2922 #[tokio::test]
2923 async fn datagram_evicts_stale() {
2924 tokio::time::pause();
2925
2926 let mut producer = track_producer("test", None);
2927 let mut dg = producer.subscribe(None);
2928 let ts = Timestamp::from_millis(0).unwrap();
2929
2930 producer.append_datagram(ts, &b"old"[..]).unwrap(); tokio::time::advance(MAX_DATAGRAM_AGE + Duration::from_millis(10)).await;
2934 producer.append_datagram(ts, &b"new"[..]).unwrap(); let got = recv_datagram(&mut dg);
2938 assert_eq!(got.sequence, 1);
2939 assert_eq!(&got.payload[..], b"new");
2940 }
2941
2942 #[tokio::test]
2943 async fn datagram_recv_pends_until_written() {
2944 let mut producer = track_producer("test", None);
2945 let mut dg = producer.subscribe(None);
2946
2947 assert!(
2948 dg.recv_datagram().now_or_never().is_none(),
2949 "should block with no datagrams"
2950 );
2951
2952 producer
2953 .append_datagram(Timestamp::from_millis(0).unwrap(), &b"go"[..])
2954 .unwrap();
2955 assert_eq!(&recv_datagram(&mut dg).payload[..], b"go");
2956 }
2957
2958 #[tokio::test]
2962 async fn datagram_wire_roundtrip_between_tracks() {
2963 use crate::coding::{Decode, Encode};
2964 use crate::lite;
2965
2966 let version = lite::Version::Lite05;
2967
2968 let mut origin = track_producer("test", None);
2970 let mut origin_dg = origin.subscribe(None);
2971 let ts = Timestamp::from_millis(7).unwrap();
2972 let seq = origin.append_datagram(ts, &b"payload"[..]).unwrap();
2973
2974 let d = recv_datagram(&mut origin_dg);
2975 let body = lite::Datagram {
2976 subscribe: 5,
2977 sequence: d.sequence,
2978 timestamp: d.timestamp.value(),
2979 payload: d.payload.clone(),
2980 }
2981 .encode_bytes(version)
2982 .unwrap();
2983
2984 let mut slice = &body[..];
2986 let wire = lite::Datagram::decode(&mut slice, version).unwrap();
2987 let mut downstream = track_producer("test", None);
2988 let mut downstream_dg = downstream.subscribe(None);
2989 downstream
2990 .write_datagram(Datagram {
2991 sequence: wire.sequence,
2992 timestamp: Timestamp::new(wire.timestamp, Timescale::MILLI).unwrap(),
2993 payload: wire.payload,
2994 })
2995 .unwrap();
2996
2997 let got = recv_datagram(&mut downstream_dg);
2998 assert_eq!(got.sequence, seq);
2999 assert_eq!(got.timestamp, ts);
3000 assert_eq!(&got.payload[..], b"payload");
3001 }
3002
3003 #[tokio::test]
3004 async fn evict_expired_groups() {
3005 tokio::time::pause();
3006
3007 let mut producer = track_producer("test", None);
3008
3009 producer.append_group().unwrap(); producer.append_group().unwrap(); producer.append_group().unwrap(); {
3015 let state = producer.state.read();
3016 assert_eq!(live_groups(&state), 3);
3017 assert_eq!(state.offset, 0);
3018 }
3019
3020 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3022
3023 producer.append_group().unwrap(); {
3030 let state = producer.state.read();
3031 assert_eq!(live_groups(&state), 1);
3032 assert_eq!(first_live_sequence(&state), 3);
3033 assert_eq!(state.offset, 3);
3034 assert!(!state.lookup.contains_key(&0));
3035 assert!(!state.lookup.contains_key(&1));
3036 assert!(!state.lookup.contains_key(&2));
3037 assert!(state.lookup.contains_key(&3));
3038 }
3039 }
3040
3041 #[tokio::test]
3045 async fn aging_out_a_finished_group_keeps_the_clean_end() {
3046 tokio::time::pause();
3047
3048 let mut producer = track_producer("test", None);
3049 let mut group = producer.create_group(group::Info { sequence: 0 }).unwrap();
3050 let mut consumer = group.consume();
3051
3052 group
3053 .write_frame(Timestamp::from_millis(0).unwrap(), b"hello".as_slice())
3054 .unwrap();
3055 assert_eq!(consumer.next_frame().await.unwrap().unwrap().size, 5);
3056
3057 tokio::time::advance(DEFAULT_LATENCY_MAX * 12).await;
3059 group.finish().unwrap();
3060 let _next = producer.create_group(group::Info { sequence: 1 }).unwrap();
3061
3062 assert!(consumer.next_frame().await.unwrap().is_none());
3063 }
3064
3065 #[tokio::test]
3066 async fn evict_keeps_max_sequence() {
3067 tokio::time::pause();
3068
3069 let mut producer = track_producer("test", None);
3070 producer.append_group().unwrap(); tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3074
3075 producer.append_group().unwrap(); {
3079 let state = producer.state.read();
3080 assert_eq!(live_groups(&state), 1);
3081 assert_eq!(first_live_sequence(&state), 1);
3082 assert_eq!(state.offset, 1);
3083 }
3084 }
3085
3086 #[tokio::test]
3087 async fn no_eviction_when_fresh() {
3088 tokio::time::pause();
3089
3090 let mut producer = track_producer("test", None);
3091 producer.append_group().unwrap(); producer.append_group().unwrap(); producer.append_group().unwrap(); {
3096 let state = producer.state.read();
3097 assert_eq!(live_groups(&state), 3);
3098 assert_eq!(state.offset, 0);
3099 }
3100 }
3101
3102 #[tokio::test]
3103 async fn consumer_skips_evicted_groups() {
3104 tokio::time::pause();
3105
3106 let mut producer = track_producer("test", None);
3107 producer.append_group().unwrap(); let mut consumer = producer.subscribe(None);
3110
3111 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3112 producer.append_group().unwrap(); let group = consumer.assert_group();
3116 assert_eq!(group.sequence, 1);
3117 }
3118
3119 #[tokio::test]
3120 async fn cache_age_controls_eviction() {
3121 tokio::time::pause();
3122
3123 let mut producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(1)));
3125 producer.append_group().unwrap(); tokio::time::advance(Duration::from_secs(2)).await;
3129 producer.append_group().unwrap(); let state = producer.state.read();
3133 assert_eq!(live_groups(&state), 1);
3134 assert_eq!(first_live_sequence(&state), 1);
3135 }
3136
3137 #[test]
3138 fn latency_max_clamped_to_cache() {
3139 let producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(2)));
3140
3141 let mut subscriber = producer.subscribe(Subscription::default().with_latency_max(Duration::from_secs(10)));
3145 assert_eq!(subscriber.subscription().latency_max, Duration::from_secs(10));
3146 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3147
3148 subscriber
3150 .update(Subscription::default().with_latency_max(Duration::from_millis(500)))
3151 .unwrap();
3152 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_millis(500));
3153
3154 subscriber
3155 .update(Subscription::default().with_latency_max(Duration::ZERO))
3156 .unwrap();
3157 assert_eq!(producer.subscription().unwrap().latency_max, Duration::ZERO);
3158 }
3159
3160 fn track_producer_capped(name: impl Into<Arc<str>>, info: Info, cap: Duration) -> Producer {
3163 let origin = crate::origin::Info::default().with_cache_duration(cap);
3164 Producer::new(Arc::new(broadcast::Info { origin }), name, info)
3165 }
3166
3167 #[test]
3168 fn origin_cache_duration_clamps_latency_max() {
3169 let capped = track_producer_capped(
3172 "test",
3173 Info::default().with_latency_max(Duration::from_secs(60)),
3174 Duration::from_secs(1),
3175 );
3176 assert_eq!(capped.state.read().latency_bound(), Some(Duration::from_secs(1)));
3177
3178 let under = track_producer_capped(
3179 "test",
3180 Info::default().with_latency_max(Duration::from_millis(500)),
3181 Duration::from_secs(1),
3182 );
3183 assert_eq!(under.state.read().latency_bound(), Some(Duration::from_millis(500)));
3184 }
3185
3186 #[tokio::test]
3187 async fn origin_cache_duration_caps_eviction() {
3188 tokio::time::pause();
3189
3190 let mut producer = track_producer_capped(
3192 "test",
3193 Info::default().with_latency_max(Duration::from_secs(60)),
3194 Duration::from_secs(1),
3195 );
3196 producer.append_group().unwrap(); tokio::time::advance(Duration::from_secs(2)).await;
3200 producer.append_group().unwrap(); let state = producer.state.read();
3204 assert_eq!(live_groups(&state), 1);
3205 assert_eq!(first_live_sequence(&state), 1);
3206 }
3207
3208 #[test]
3209 fn latency_max_clamped_via_every_update_path() {
3210 let producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(2)));
3211 let over = Subscription::default().with_latency_max(Duration::from_secs(10));
3212
3213 let mut subscriber = producer.subscribe(over.clone());
3216 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3217
3218 subscriber.control().update(over.clone()).unwrap();
3219 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3220
3221 subscriber.update(over).unwrap();
3222 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3223 }
3224
3225 #[test]
3226 fn latency_max_aggregate_clamps_the_max_across_subscribers() {
3227 let producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(2)));
3228
3229 let _a = producer.subscribe(Subscription::default().with_latency_max(Duration::from_millis(500)));
3232 let _b = producer.subscribe(Subscription::default().with_latency_max(Duration::from_secs(10)));
3233
3234 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3235 }
3236
3237 #[test]
3238 fn subscriber_control_updates_while_read_future_is_pending() {
3239 let producer = track_producer("test", None);
3240 let mut subscriber = producer.subscribe(None);
3241 let control = subscriber.control();
3242
3243 let mut recv = Box::pin(subscriber.recv_group());
3244 assert!(recv.as_mut().now_or_never().is_none());
3245
3246 control
3247 .update(Subscription::default().with_priority(7).with_ordered(false))
3248 .unwrap();
3249
3250 let aggregate = producer.subscription().expect("expected an active subscription");
3251 assert_eq!(aggregate.priority, 7);
3252 assert!(!aggregate.ordered);
3253 }
3254
3255 #[test]
3256 fn dropped_subscriber_leaves_no_ghost_in_aggregate() {
3257 let mut producer = track_producer("test", None);
3262 let a = producer.subscribe(Subscription::default().with_priority(5));
3263
3264 let waiter = kio::Waiter::noop();
3266 assert!(
3267 matches!(producer.poll_subscription_changed(&waiter), Poll::Ready(Ok(Some(_)))),
3268 "one live subscriber should aggregate to Some",
3269 );
3270
3271 drop(a);
3273
3274 assert!(
3276 matches!(producer.poll_subscription_changed(&waiter), Poll::Ready(Ok(None))),
3277 "a dropped subscriber must not linger in the aggregate",
3278 );
3279
3280 assert!(
3282 producer.subscription().is_none(),
3283 "snapshot must exclude a dropped subscriber",
3284 );
3285 }
3286
3287 #[test]
3288 fn dropped_subscriber_wakes_the_aggregate() {
3289 use std::sync::atomic::{AtomicBool, Ordering};
3296
3297 let mut producer = track_producer("test", None);
3298 let a = producer.subscribe(Subscription::default().with_priority(5));
3299
3300 let woken = Arc::new(AtomicBool::new(false));
3301 let waiter = kio::Waiter::new(futures::task::waker(Arc::new(FlagWake(woken.clone()))));
3302
3303 assert!(matches!(
3305 producer.poll_subscription_changed(&waiter),
3306 Poll::Ready(Ok(Some(_)))
3307 ));
3308 assert!(
3309 producer.poll_subscription_changed(&waiter).is_pending(),
3310 "the aggregate is unchanged, so this poll must park",
3311 );
3312 assert!(!woken.load(Ordering::SeqCst), "nothing happened yet");
3313
3314 drop(a);
3315 assert!(
3316 woken.load(Ordering::SeqCst),
3317 "the last subscriber leaving must wake the aggregate watcher",
3318 );
3319 }
3320
3321 struct FlagWake(Arc<std::sync::atomic::AtomicBool>);
3323
3324 impl futures::task::ArcWake for FlagWake {
3325 fn wake_by_ref(arc_self: &Arc<Self>) {
3326 arc_self.0.store(true, std::sync::atomic::Ordering::SeqCst);
3327 }
3328 }
3329
3330 #[tokio::test]
3331 async fn out_of_order_max_sequence_at_front() {
3332 tokio::time::pause();
3333
3334 let mut producer = track_producer("test", None);
3335
3336 producer.create_group(group::Info { sequence: 5 }).unwrap();
3338 producer.create_group(group::Info { sequence: 3 }).unwrap();
3339 producer.create_group(group::Info { sequence: 4 }).unwrap();
3340
3341 {
3343 let state = producer.state.read();
3344 assert_eq!(state.max_sequence, Some(5));
3345 }
3346
3347 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3349
3350 producer.append_group().unwrap(); {
3356 let state = producer.state.read();
3357 assert_eq!(live_groups(&state), 1);
3358 assert_eq!(first_live_sequence(&state), 6);
3359 assert!(!state.lookup.contains_key(&3));
3360 assert!(!state.lookup.contains_key(&4));
3361 assert!(!state.lookup.contains_key(&5));
3362 assert!(state.lookup.contains_key(&6));
3363 }
3364 }
3365
3366 #[tokio::test]
3367 async fn max_sequence_at_front_blocks_trim() {
3368 tokio::time::pause();
3369
3370 let mut producer = track_producer("test", None);
3371
3372 producer.create_group(group::Info { sequence: 5 }).unwrap();
3374
3375 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3376
3377 producer.create_group(group::Info { sequence: 3 }).unwrap();
3379
3380 {
3383 let state = producer.state.read();
3384 assert_eq!(live_groups(&state), 2);
3385 assert_eq!(state.offset, 0);
3386 }
3387
3388 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3390
3391 producer.create_group(group::Info { sequence: 2 }).unwrap();
3393
3394 {
3399 let state = producer.state.read();
3400 assert_eq!(live_groups(&state), 2);
3401 assert_eq!(state.offset, 0);
3402 assert!(state.lookup.contains_key(&5));
3403 assert!(!state.lookup.contains_key(&3));
3404 assert!(state.lookup.contains_key(&2));
3405 }
3406
3407 let mut consumer = producer.subscribe(None);
3409 let group = consumer.assert_group();
3410 assert_eq!(group.sequence, 5);
3412 }
3413
3414 #[tokio::test]
3415 async fn abort_clears_cached_groups() {
3416 let mut producer = track_producer("test", None);
3417 producer.append_group().unwrap();
3418 producer.append_group().unwrap();
3419
3420 let mut consumer = producer.subscribe(None);
3422 assert_eq!(live_groups(&producer.state.read()), 2);
3423
3424 producer.clone().abort(Error::Cancel).unwrap();
3425
3426 {
3427 let state = producer.state.read();
3428 assert!(state.lookup.is_empty(), "cached groups should be dropped on abort");
3429 assert!(state.arrival.is_empty());
3430 assert!(state.evict.is_empty());
3431 }
3432
3433 let result = consumer.recv_group().now_or_never().expect("should not block");
3435 assert!(matches!(result, Err(Error::Cancel)));
3436 }
3437
3438 #[tokio::test]
3439 async fn drop_unfinished_clears_cached_groups() {
3440 let producer = track_producer("test", None);
3441 let mut writer = producer.clone();
3442 writer.append_group().unwrap();
3443
3444 let mut consumer = producer.subscribe(None);
3446 assert_eq!(live_groups(&producer.state.read()), 1);
3447
3448 drop(writer);
3450 drop(producer);
3451
3452 let result = consumer.recv_group().now_or_never().expect("should not block");
3453 assert!(matches!(result, Err(Error::Dropped)));
3454 }
3455
3456 #[tokio::test]
3457 async fn drop_after_abort_does_not_warn() {
3458 let warns = count_drop_warnings("track::Producer dropped without finish", || {
3461 let producer = track_producer("test", None);
3462 let keep = producer.clone();
3463 let mut writer = producer.clone();
3464 let mut group = writer.append_group().unwrap();
3465 group.finish().unwrap();
3466 let _consumer = producer.subscribe(None);
3467 writer.abort(Error::Cancel).unwrap();
3468 drop(keep);
3469 });
3470 assert_eq!(warns, 0, "abort-then-drop must not emit unfinished-producer WARN");
3471 }
3472
3473 #[tokio::test]
3474 async fn drop_unfinished_warns() {
3475 let warns = count_drop_warnings("track::Producer dropped without finish", || {
3476 let producer = track_producer("test", None);
3477 let mut writer = producer.clone();
3478 writer.append_group().unwrap();
3479 let _consumer = producer.subscribe(None);
3480 drop(writer);
3481 drop(producer);
3482 });
3483 assert!(warns >= 1, "unfinished drop must emit unfinished-producer WARN");
3484 }
3485
3486 #[tokio::test]
3487 async fn drop_finished_keeps_cached_groups() {
3488 let mut producer = track_producer("test", None);
3489 producer.append_group().unwrap();
3490 producer.finish().unwrap();
3491
3492 let mut consumer = producer.subscribe(None);
3493 drop(producer);
3494
3495 assert_eq!(consumer.assert_group().sequence, 0);
3497 let done = consumer.recv_group().now_or_never().expect("should not block").unwrap();
3498 assert!(done.is_none(), "consumer should drain then see clean finish");
3499 }
3500
3501 #[test]
3502 fn append_finish_cannot_be_rewritten() {
3503 let mut producer = track_producer("test", None);
3504
3505 assert!(producer.finish().is_ok());
3507 assert!(producer.finish().is_err());
3508 assert!(producer.append_group().is_err());
3509 }
3510
3511 #[test]
3512 fn finish_after_groups() {
3513 let mut producer = track_producer("test", None);
3514
3515 producer.append_group().unwrap();
3516 assert!(producer.finish().is_ok());
3517 assert!(producer.finish().is_err());
3518 assert!(producer.append_group().is_err());
3519 }
3520
3521 #[test]
3522 fn finish_at_rejects_a_boundary_at_or_below_the_live_edge() {
3523 let mut producer = track_producer("test", None);
3524 producer.create_group(group::Info { sequence: 5 }).unwrap();
3525
3526 assert!(producer.finish_at(4).is_err());
3529 assert!(producer.finish_at(5).is_err());
3530 assert!(producer.finish_at(6).is_ok());
3531
3532 {
3533 let state = producer.state.read();
3534 assert_eq!(state.final_sequence, Some(6));
3535 }
3536
3537 assert!(producer.finish_at(6).is_err());
3539 assert!(producer.create_group(group::Info { sequence: 4 }).is_ok());
3540 assert!(producer.create_group(group::Info { sequence: 6 }).is_err());
3541 }
3542
3543 #[test]
3544 fn final_sequence_reports_the_declared_boundary() {
3545 let mut producer = track_producer("test", None);
3546 assert_eq!(producer.final_sequence(), None);
3547
3548 producer.create_group(group::Info { sequence: 5 }).unwrap();
3549 assert_eq!(producer.final_sequence(), None, "a group does not declare a boundary");
3550
3551 producer.finish_at(9).unwrap();
3552 assert_eq!(producer.final_sequence(), Some(9));
3553
3554 assert!(producer.finish().is_err());
3556 }
3557
3558 #[test]
3559 fn final_sequence_reports_the_live_edge_after_finish() {
3560 let mut producer = track_producer("test", None);
3561 producer.create_group(group::Info { sequence: 5 }).unwrap();
3562 producer.finish().unwrap();
3563 assert_eq!(producer.final_sequence(), Some(6));
3564 }
3565
3566 #[tokio::test]
3567 async fn finish_at_declares_a_future_boundary() {
3568 let mut producer = track_producer("test", None);
3569 producer.create_group(group::Info { sequence: 5 }).unwrap();
3570
3571 producer.finish_at(7).unwrap();
3573
3574 let mut consumer = producer.subscribe(None);
3575 assert_eq!(consumer.assert_group().sequence, 5);
3576
3577 let boundary = consumer
3580 .finished()
3581 .now_or_never()
3582 .expect("boundary is known immediately")
3583 .expect("would have errored");
3584 assert_eq!(boundary, 7);
3585 assert!(
3586 consumer.recv_group().now_or_never().is_none(),
3587 "should wait for the outstanding group"
3588 );
3589
3590 producer.create_group(group::Info { sequence: 6 }).unwrap();
3592 assert_eq!(consumer.assert_group().sequence, 6);
3593 let done = consumer
3594 .recv_group()
3595 .now_or_never()
3596 .expect("should not block")
3597 .expect("would have errored");
3598 assert!(done.is_none(), "track completes once the boundary is reached");
3599 }
3600
3601 #[tokio::test]
3602 async fn recv_group_finishes_without_waiting_for_gaps() {
3603 let mut producer = track_producer("test", None);
3604 producer.create_group(group::Info { sequence: 1 }).unwrap();
3605 producer.finish().unwrap();
3606
3607 let mut consumer = producer.subscribe(None);
3608 assert_eq!(consumer.assert_group().sequence, 1);
3609
3610 let done = consumer
3611 .recv_group()
3612 .now_or_never()
3613 .expect("should not block")
3614 .expect("would have errored");
3615 assert!(done.is_none(), "track should finish without waiting for gaps");
3616 }
3617
3618 #[tokio::test]
3619 async fn next_group_skips_late_arrivals() {
3620 let mut producer = track_producer("test", None);
3621 let mut consumer = producer.subscribe(None);
3622
3623 producer.create_group(group::Info { sequence: 5 }).unwrap();
3625 let group = consumer
3626 .next_group()
3627 .now_or_never()
3628 .expect("should not block")
3629 .expect("would have errored")
3630 .expect("track should not be closed");
3631 assert_eq!(group.sequence, 5);
3632
3633 producer.create_group(group::Info { sequence: 3 }).unwrap();
3635 producer.create_group(group::Info { sequence: 4 }).unwrap();
3637 producer.create_group(group::Info { sequence: 7 }).unwrap();
3639
3640 let group = consumer
3641 .next_group()
3642 .now_or_never()
3643 .expect("should not block")
3644 .expect("would have errored")
3645 .expect("track should not be closed");
3646 assert_eq!(group.sequence, 7);
3647
3648 assert!(
3650 consumer.next_group().now_or_never().is_none(),
3651 "should block waiting for a higher sequence"
3652 );
3653 }
3654
3655 #[tokio::test]
3656 async fn next_group_returns_arrivals_in_order() {
3657 let mut producer = track_producer("test", None);
3658 let mut consumer = producer.subscribe(None);
3659
3660 producer.create_group(group::Info { sequence: 3 }).unwrap();
3662 producer.create_group(group::Info { sequence: 5 }).unwrap();
3663
3664 let group = consumer
3665 .next_group()
3666 .now_or_never()
3667 .expect("should not block")
3668 .expect("would have errored")
3669 .expect("track should not be closed");
3670 assert_eq!(group.sequence, 3);
3671
3672 let group = consumer
3673 .next_group()
3674 .now_or_never()
3675 .expect("should not block")
3676 .expect("would have errored")
3677 .expect("track should not be closed");
3678 assert_eq!(group.sequence, 5);
3679 }
3680
3681 #[tokio::test]
3682 async fn next_group_and_recv_group_use_independent_cursors() {
3683 let mut producer = track_producer("test", None);
3684 let mut consumer = producer.subscribe(None);
3685
3686 producer.create_group(group::Info { sequence: 5 }).unwrap();
3688 producer.create_group(group::Info { sequence: 3 }).unwrap();
3689
3690 let group = consumer
3693 .next_group()
3694 .now_or_never()
3695 .expect("should not block")
3696 .expect("would have errored")
3697 .expect("track should not be closed");
3698 assert_eq!(group.sequence, 3);
3699
3700 assert_eq!(consumer.assert_group().sequence, 5);
3703 }
3704
3705 #[tokio::test]
3706 async fn end_at_caps_next_group() {
3707 let mut producer = track_producer("test", None);
3708 let mut consumer = producer.subscribe(None);
3709
3710 for s in 0..6 {
3711 producer.create_group(group::Info { sequence: s }).unwrap();
3712 }
3713
3714 consumer.end_at(2);
3715
3716 assert_eq!(
3718 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3719 0
3720 );
3721 assert_eq!(
3722 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3723 1
3724 );
3725 assert_eq!(
3726 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3727 2
3728 );
3729
3730 assert!(
3732 consumer.next_group().now_or_never().is_none(),
3733 "capped consumer must block instead of returning out-of-range groups"
3734 );
3735 }
3736
3737 #[tokio::test]
3738 async fn end_at_release_drains_cached_groups() {
3739 let mut producer = track_producer("test", None);
3740 let mut consumer = producer.subscribe(None);
3741
3742 for s in 0..6 {
3743 producer.create_group(group::Info { sequence: s }).unwrap();
3744 }
3745
3746 consumer.end_at(1);
3747 assert_eq!(
3748 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3749 0
3750 );
3751 assert_eq!(
3752 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3753 1
3754 );
3755 assert!(consumer.next_group().now_or_never().is_none(), "capped at 1");
3756
3757 consumer.end_at(4);
3759 assert_eq!(
3760 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3761 2
3762 );
3763 assert_eq!(
3764 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3765 3
3766 );
3767 assert_eq!(
3768 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3769 4
3770 );
3771 assert!(consumer.next_group().now_or_never().is_none(), "capped at 4");
3772
3773 consumer.end_at(None);
3775 assert_eq!(
3776 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3777 5
3778 );
3779 assert!(consumer.next_group().now_or_never().is_none(), "no more groups");
3780 }
3781
3782 #[tokio::test]
3783 async fn end_at_lower_than_cursor_parks_consumer() {
3784 let mut producer = track_producer("test", None);
3785 let mut consumer = producer.subscribe(None);
3786
3787 for s in 0..3 {
3788 producer.create_group(group::Info { sequence: s }).unwrap();
3789 }
3790
3791 assert_eq!(
3793 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3794 0
3795 );
3796 assert_eq!(
3797 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3798 1
3799 );
3800 assert_eq!(
3801 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3802 2
3803 );
3804
3805 consumer.end_at(1);
3807 producer.create_group(group::Info { sequence: 3 }).unwrap();
3808 producer.create_group(group::Info { sequence: 4 }).unwrap();
3809 assert!(
3810 consumer.next_group().now_or_never().is_none(),
3811 "cap is below cursor; nothing returnable until cap rises"
3812 );
3813
3814 consumer.end_at(None);
3816 assert_eq!(
3817 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3818 3
3819 );
3820 assert_eq!(
3821 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3822 4
3823 );
3824 }
3825
3826 #[tokio::test]
3827 async fn end_at_toggling_around_late_arrivals() {
3828 let mut producer = track_producer("test", None);
3829 let mut consumer = producer.subscribe(None);
3830
3831 consumer.end_at(5);
3832
3833 producer.create_group(group::Info { sequence: 2 }).unwrap();
3835 producer.create_group(group::Info { sequence: 5 }).unwrap();
3836 producer.create_group(group::Info { sequence: 3 }).unwrap();
3837 producer.create_group(group::Info { sequence: 8 }).unwrap();
3839 producer.create_group(group::Info { sequence: 4 }).unwrap();
3840
3841 assert_eq!(
3843 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3844 2
3845 );
3846 assert_eq!(
3847 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3848 3
3849 );
3850 assert_eq!(
3851 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3852 4
3853 );
3854 assert_eq!(
3855 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3856 5
3857 );
3858 assert!(consumer.next_group().now_or_never().is_none());
3860
3861 consumer.end_at(10);
3863 assert_eq!(
3864 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3865 8
3866 );
3867 }
3868
3869 #[tokio::test]
3873 async fn end_at_parks_recv_group() {
3874 let mut producer = track_producer("test", None);
3875 let mut consumer = producer.subscribe(None);
3876
3877 for s in 0..3 {
3878 producer.create_group(group::Info { sequence: s }).unwrap();
3879 }
3880
3881 consumer.end_at(1);
3882 assert_eq!(
3883 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3884 0
3885 );
3886 assert_eq!(
3887 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3888 1
3889 );
3890 assert!(consumer.recv_group().now_or_never().is_none(), "capped at 1");
3891
3892 producer.finish().unwrap();
3894 assert!(
3895 consumer.recv_group().now_or_never().is_none(),
3896 "still parked after finish"
3897 );
3898
3899 consumer.end_at(None);
3900 assert_eq!(
3901 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3902 2
3903 );
3904 assert!(
3905 matches!(consumer.recv_group().now_or_never(), Some(Ok(None))),
3906 "finished once the parked group drains"
3907 );
3908 }
3909
3910 #[tokio::test]
3913 async fn recv_group_serves_arrivals_behind_the_cap() {
3914 let mut producer = track_producer("test", None);
3915 let mut consumer = producer.subscribe(None);
3916
3917 consumer.end_at(1);
3918
3919 producer.create_group(group::Info { sequence: 2 }).unwrap();
3921 producer.create_group(group::Info { sequence: 0 }).unwrap();
3922 producer.create_group(group::Info { sequence: 1 }).unwrap();
3923
3924 assert_eq!(
3925 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3926 0
3927 );
3928 assert_eq!(
3929 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3930 1
3931 );
3932 assert!(consumer.recv_group().now_or_never().is_none(), "capped at 1");
3933
3934 consumer.end_at(2);
3935 assert_eq!(
3936 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3937 2
3938 );
3939 }
3940
3941 #[tokio::test]
3944 async fn start_at_drops_parked_recv_groups() {
3945 let mut producer = track_producer("test", None);
3946 let mut consumer = producer.subscribe(None);
3947
3948 consumer.end_at(0);
3949 producer.create_group(group::Info { sequence: 1 }).unwrap();
3950 assert!(
3951 consumer.recv_group().now_or_never().is_none(),
3952 "group 1 parked at the cap"
3953 );
3954
3955 consumer.start_at(2);
3956 consumer.end_at(None);
3957 producer.create_group(group::Info { sequence: 2 }).unwrap();
3958 assert_eq!(
3959 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3960 2,
3961 "the overtaken parked group is dropped, not re-offered"
3962 );
3963 }
3964
3965 #[tokio::test]
3969 async fn evicted_parked_recv_groups_are_dropped() {
3970 let mut producer = track_producer("test", None);
3971 let mut consumer = producer.subscribe(None);
3972
3973 producer.create_group(group::Info { sequence: 0 }).unwrap();
3974 assert_eq!(
3975 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3976 0
3977 );
3978
3979 consumer.end_at(0);
3980 let straggler = producer.create_group(group::Info { sequence: 1 }).unwrap();
3981 assert!(
3982 consumer.recv_group().now_or_never().is_none(),
3983 "group 1 parked at the cap"
3984 );
3985
3986 straggler.abort(Error::Old).unwrap();
3988 producer.finish().unwrap();
3989
3990 consumer.end_at(None);
3991 assert!(
3992 matches!(consumer.recv_group().now_or_never(), Some(Ok(None))),
3993 "a dead parked group must not be delivered or hold the stream open"
3994 );
3995 }
3996
3997 #[tokio::test]
3998 async fn read_frame_returns_single_frame_per_group() {
3999 let mut producer = track_producer("test", None);
4000 let mut consumer = producer.subscribe(None);
4001
4002 producer.write_frame(Timestamp::ZERO, b"hello".as_slice()).unwrap();
4003 producer.write_frame(Timestamp::ZERO, b"world".as_slice()).unwrap();
4004
4005 let frame = consumer
4006 .read_frame()
4007 .now_or_never()
4008 .expect("should not block")
4009 .expect("would have errored")
4010 .expect("track should not be closed");
4011 assert_eq!(&frame.payload[..], b"hello");
4012
4013 let frame = consumer
4014 .read_frame()
4015 .now_or_never()
4016 .expect("should not block")
4017 .expect("would have errored")
4018 .expect("track should not be closed");
4019 assert_eq!(&frame.payload[..], b"world");
4020 }
4021
4022 #[tokio::test]
4023 async fn read_frame_preserves_timestamp() {
4024 let mut producer = track_producer("test", None);
4025 let mut consumer = producer.subscribe(None);
4026
4027 producer
4028 .write_frame(Timestamp::from_micros(20_000).unwrap(), b"hello".as_slice())
4029 .unwrap();
4030
4031 let frame = consumer
4032 .read_frame()
4033 .now_or_never()
4034 .expect("should not block")
4035 .expect("would have errored")
4036 .expect("track should not be closed");
4037 assert_eq!(frame.timestamp.as_micros(), 20_000);
4038 assert_eq!(&frame.payload[..], b"hello");
4039 }
4040
4041 #[tokio::test]
4042 async fn read_frame_skips_stalled_group_for_newer_ready_frame() {
4043 let mut producer = track_producer("test", None);
4044 let mut consumer = producer.subscribe(None);
4045
4046 let _stalled = producer.create_group(group::Info { sequence: 3 }).unwrap();
4048 let mut g5 = producer.create_group(group::Info { sequence: 5 }).unwrap();
4050 g5.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"later"))
4051 .unwrap();
4052 g5.finish().unwrap();
4053
4054 let frame = consumer
4056 .read_frame()
4057 .now_or_never()
4058 .expect("should not block on stalled earlier group")
4059 .expect("would have errored")
4060 .expect("track should not be closed");
4061 assert_eq!(&frame.payload[..], b"later");
4062 }
4063
4064 #[tokio::test]
4065 async fn read_frame_discards_rest_of_multi_frame_group() {
4066 let mut producer = track_producer("test", None);
4067 let mut consumer = producer.subscribe(None);
4068
4069 let mut g0 = producer.create_group(group::Info { sequence: 0 }).unwrap();
4071 g0.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"one"))
4072 .unwrap();
4073 g0.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"two"))
4074 .unwrap();
4075 g0.finish().unwrap();
4076
4077 producer.write_frame(Timestamp::ZERO, b"next".as_slice()).unwrap();
4079
4080 let frame = consumer
4081 .read_frame()
4082 .now_or_never()
4083 .expect("should not block")
4084 .expect("would have errored")
4085 .expect("track should not be closed");
4086 assert_eq!(&frame.payload[..], b"one");
4087
4088 let frame = consumer
4090 .read_frame()
4091 .now_or_never()
4092 .expect("should not block")
4093 .expect("would have errored")
4094 .expect("track should not be closed");
4095 assert_eq!(&frame.payload[..], b"next");
4096 }
4097
4098 #[tokio::test]
4099 async fn read_frame_waits_for_pending_group_after_finish() {
4100 let mut producer = track_producer("test", None);
4103 let mut consumer = producer.subscribe(None);
4104
4105 let mut g0 = producer.create_group(group::Info { sequence: 0 }).unwrap();
4106 producer.finish().unwrap();
4107
4108 assert!(
4110 consumer.read_frame().now_or_never().is_none(),
4111 "read_frame must block on a pending group even after finish()"
4112 );
4113
4114 g0.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"late"))
4116 .unwrap();
4117 let frame = consumer
4118 .read_frame()
4119 .now_or_never()
4120 .expect("should not block once a frame is written")
4121 .expect("would have errored")
4122 .expect("track should not be closed");
4123 assert_eq!(&frame.payload[..], b"late");
4124 }
4125
4126 #[tokio::test]
4127 async fn read_frame_respects_start_at() {
4128 let mut producer = track_producer("test", None);
4131 let mut consumer = producer.subscribe(None);
4132 consumer.start_at(5);
4133
4134 let mut g3 = producer.create_group(group::Info { sequence: 3 }).unwrap();
4136 g3.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"skip-me"))
4137 .unwrap();
4138 g3.finish().unwrap();
4139
4140 let mut g5 = producer.create_group(group::Info { sequence: 5 }).unwrap();
4141 g5.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"keep"))
4142 .unwrap();
4143 g5.finish().unwrap();
4144
4145 let frame = consumer
4146 .read_frame()
4147 .now_or_never()
4148 .expect("should not block")
4149 .expect("would have errored")
4150 .expect("track should not be closed");
4151 assert_eq!(&frame.payload[..], b"keep");
4152 }
4153
4154 #[tokio::test]
4155 async fn read_frame_returns_none_when_finished() {
4156 let mut producer = track_producer("test", None);
4157 let mut consumer = producer.subscribe(None);
4158
4159 producer.write_frame(Timestamp::ZERO, b"only".as_slice()).unwrap();
4160 producer.finish().unwrap();
4161
4162 let frame = consumer
4163 .read_frame()
4164 .now_or_never()
4165 .expect("should not block")
4166 .expect("would have errored")
4167 .expect("track should not be closed");
4168 assert_eq!(&frame.payload[..], b"only");
4169
4170 let done = consumer
4171 .read_frame()
4172 .now_or_never()
4173 .expect("should not block")
4174 .expect("would have errored");
4175 assert!(done.is_none());
4176 }
4177
4178 #[test]
4179 fn append_group_returns_bounds_exceeded_on_sequence_overflow() {
4180 let mut producer = track_producer("test", None);
4181 {
4182 let mut state = producer.state.write().ok().unwrap();
4183 state.max_sequence = Some(u64::MAX);
4184 }
4185
4186 assert!(matches!(producer.append_group(), Err(Error::BoundsExceeded(_))));
4187 }
4188
4189 #[tokio::test]
4190 async fn fetch_cache_hit() {
4191 let mut producer = track_producer("test", None);
4192
4193 let mut group = producer.append_group().unwrap(); group
4196 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hello"))
4197 .unwrap();
4198 group.finish().unwrap();
4199
4200 let dynamic = producer.dynamic();
4203 let consumer = producer.consume();
4204 assert!(consumer.peek_group(0).is_some());
4205 let mut g = consumer.fetch_group(0, None).await.unwrap();
4206 assert_eq!(g.sequence, 0);
4207 assert_eq!(&g.read_frame().await.unwrap().unwrap().payload[..], b"hello");
4208
4209 assert!(dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending());
4211 }
4212
4213 #[tokio::test]
4214 async fn fetch_miss_signals_dynamic() {
4215 let producer = track_producer("test", None);
4216 let dynamic = producer.dynamic();
4217 let consumer = producer.consume();
4218
4219 assert!(consumer.peek_group(5).is_none());
4223 let pending = consumer.fetch_group(5, group::Fetch::default().with_priority(7));
4224 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
4225
4226 let req = dynamic
4227 .requested_group()
4228 .now_or_never()
4229 .expect("should not block")
4230 .unwrap();
4231 assert_eq!(req.sequence(), 5);
4232 assert_eq!(req.priority(), 7);
4233
4234 let mut group = req.accept(None).unwrap();
4236 group
4237 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hi"))
4238 .unwrap();
4239 group.finish().unwrap();
4240
4241 let mut g = pending.await.unwrap();
4242 assert_eq!(g.sequence, 5);
4243 assert_eq!(&g.read_frame().await.unwrap().unwrap().payload[..], b"hi");
4244 }
4245
4246 #[tokio::test]
4247 async fn fetch_miss_rejects() {
4248 let producer = track_producer("test", None);
4249 let dynamic = producer.dynamic();
4250 let consumer = producer.consume();
4251
4252 let pending = consumer.fetch_group(5, None);
4253 let req = dynamic
4254 .requested_group()
4255 .now_or_never()
4256 .expect("should not block")
4257 .unwrap();
4258
4259 req.reject(Error::Cancel);
4260 assert!(matches!(pending.await, Err(Error::Cancel)));
4261 let fetch = producer.state.read().fetch.clone();
4262 assert!(fetch.read().is_empty());
4263 }
4264
4265 #[tokio::test]
4266 async fn fetch_miss_drop_rejects() {
4267 let producer = track_producer("test", None);
4268 let dynamic = producer.dynamic();
4269 let consumer = producer.consume();
4270
4271 let pending = consumer.fetch_group(5, None);
4272 let req = dynamic
4273 .requested_group()
4274 .now_or_never()
4275 .expect("should not block")
4276 .unwrap();
4277
4278 drop(req);
4279 assert!(matches!(pending.await, Err(Error::Dropped)));
4280 }
4281
4282 #[tokio::test]
4283 async fn fetch_reject_does_not_poison_retry() {
4284 let producer = track_producer("test", None);
4285 let dynamic = producer.dynamic();
4286 let consumer = producer.consume();
4287
4288 let pending = consumer.fetch_group(5, None);
4289 let req = dynamic
4290 .requested_group()
4291 .now_or_never()
4292 .expect("should not block")
4293 .unwrap();
4294 req.reject(Error::Cancel);
4295 assert!(matches!(pending.await, Err(Error::Cancel)));
4296
4297 let retry = consumer.fetch_group(5, None);
4298 let req = dynamic
4299 .requested_group()
4300 .now_or_never()
4301 .expect("should not block")
4302 .unwrap();
4303 let mut group = req.accept(None).unwrap();
4304 group
4305 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"retry"))
4306 .unwrap();
4307 group.finish().unwrap();
4308
4309 let mut group = retry.await.unwrap();
4310 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"retry");
4311 }
4312
4313 #[tokio::test]
4314 async fn fetch_coalesces_concurrent() {
4315 let producer = track_producer("test", None);
4316 let dynamic = producer.dynamic();
4317 let consumer = producer.consume();
4318
4319 let first = consumer.fetch_group(5, group::Fetch::default().with_priority(1));
4322 let second = consumer.fetch_group(5, group::Fetch::default().with_priority(7));
4323 assert!(kio::Pollable::poll(&*first, &kio::Waiter::noop()).is_pending());
4324
4325 let req = dynamic
4326 .requested_group()
4327 .now_or_never()
4328 .expect("should not block")
4329 .unwrap();
4330 assert_eq!(req.sequence(), 5);
4331 assert_eq!(req.priority(), 7);
4332 assert!(
4333 dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending(),
4334 "the second fetch queued a duplicate request"
4335 );
4336
4337 let third = consumer.fetch_group(5, None);
4339
4340 let mut group = req.accept(None).unwrap();
4342 group
4343 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hi"))
4344 .unwrap();
4345 group.finish().unwrap();
4346
4347 assert_eq!(first.await.unwrap().sequence, 5);
4348 assert_eq!(second.await.unwrap().sequence, 5);
4349 assert_eq!(third.await.unwrap().sequence, 5);
4350 }
4351
4352 #[tokio::test]
4353 async fn fetch_coalesced_reject_fails_all() {
4354 let producer = track_producer("test", None);
4355 let dynamic = producer.dynamic();
4356 let consumer = producer.consume();
4357
4358 let first = consumer.fetch_group(5, None);
4359 let second = consumer.fetch_group(5, None);
4360 let req = dynamic
4361 .requested_group()
4362 .now_or_never()
4363 .expect("should not block")
4364 .unwrap();
4365 req.reject(Error::Cancel);
4366
4367 assert!(matches!(first.await, Err(Error::Cancel)));
4368 assert!(matches!(second.await, Err(Error::Cancel)));
4369
4370 let retry = consumer.fetch_group(5, None);
4372 assert!(kio::Pollable::poll(&*retry, &kio::Waiter::noop()).is_pending());
4373 let req = dynamic
4374 .requested_group()
4375 .now_or_never()
4376 .expect("should not block")
4377 .unwrap();
4378 assert_eq!(req.sequence(), 5);
4379 }
4380
4381 #[tokio::test]
4382 async fn fetch_queued_fails_when_handlers_leave() {
4383 let producer = track_producer("test", None);
4384 let dynamic = producer.dynamic();
4385 let consumer = producer.consume();
4386
4387 let pending = consumer.fetch_group(5, None);
4389 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
4390 drop(dynamic);
4391 assert!(matches!(pending.await, Err(Error::NotFound)));
4392
4393 let fetch = producer.state.read().fetch.clone();
4395 assert!(fetch.read().is_empty());
4396 }
4397
4398 #[tokio::test]
4399 async fn fetch_miss_no_dynamic_not_found() {
4400 let mut producer = track_producer("test", None);
4403 producer.append_group().unwrap(); let consumer = producer.consume();
4405 assert!(matches!(consumer.fetch_group(5, None).await, Err(Error::NotFound)));
4406 }
4407
4408 #[tokio::test]
4409 async fn fetch_past_final_not_found() {
4410 let mut producer = track_producer("test", None);
4411 producer.append_group().unwrap(); producer.finish().unwrap(); let dynamic = producer.dynamic();
4417 let consumer = producer.consume();
4418 assert!(matches!(consumer.fetch_group(5, None).await, Err(Error::NotFound)));
4419
4420 assert!(dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending());
4422 }
4423
4424 fn pooled_producer(capacity: u64) -> (Producer, cache::Pool) {
4426 let pool = cache::Pool::new(capacity);
4427 let broadcast = broadcast::Info {
4428 origin: crate::origin::Info::default().with_pool(pool.clone()),
4429 ..Default::default()
4430 };
4431 let producer = Producer::new(Arc::new(broadcast), "test", None);
4432 (producer, pool)
4433 }
4434
4435 fn finished_group(producer: &mut Producer, size: usize) -> u64 {
4436 let mut group = producer.append_group().unwrap();
4437 group
4438 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; size]))
4439 .unwrap();
4440 group.finish().unwrap();
4441 group.sequence
4442 }
4443
4444 #[tokio::test]
4447 async fn debt_evicts_oldest_group() {
4448 tokio::time::pause();
4449
4450 let (mut producer, pool) = pooled_producer(10_000);
4452
4453 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let consumer = producer.consume();
4458 assert!(consumer.peek_group(0).is_none(), "oldest group is evicted");
4459 assert!(consumer.peek_group(2).is_some(), "latest group survives");
4460 assert!(pool.used() <= 21_000, "usage hovers near capacity: {}", pool.used());
4463
4464 let mut subscriber = producer.subscribe(None);
4466 assert!(subscriber.assert_group().sequence > 0, "evicted group is not delivered");
4467 }
4468
4469 #[tokio::test]
4471 async fn latest_group_never_evicted() {
4472 tokio::time::pause();
4473
4474 let (mut producer, pool) = pooled_producer(100);
4476 finished_group(&mut producer, 1000); assert!(pool.used() > 100, "the latest may exceed the budget");
4478
4479 finished_group(&mut producer, 1000); finished_group(&mut producer, 1000); let consumer = producer.consume();
4484 assert!(consumer.peek_group(0).is_none());
4485 let mut group = consumer.peek_group(2).expect("latest survives");
4486 assert_eq!(group.read_frame().await.unwrap().unwrap().payload.len(), 1000);
4487 }
4488
4489 #[tokio::test]
4493 async fn fetch_refresh_survives_eviction() {
4494 tokio::time::pause();
4495
4496 let (mut producer, _pool) = pooled_producer(10_000);
4497 let consumer = producer.consume();
4498
4499 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_secs(1)).await;
4501 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_secs(1)).await;
4503 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_millis(500)).await;
4505
4506 let mut fetched = consumer.fetch_group(0, None).await.unwrap();
4508 assert_eq!(fetched.read_frame().await.unwrap().unwrap().payload.len(), 3_000);
4509 tokio::time::advance(Duration::from_millis(500)).await;
4510
4511 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_secs(1)).await;
4515 finished_group(&mut producer, 3_000); assert!(consumer.peek_group(0).is_some(), "refreshed group survives");
4518 assert!(consumer.peek_group(1).is_none(), "unread group is evicted instead");
4519 }
4520
4521 #[tokio::test]
4524 async fn eviction_aborts_readers() {
4525 tokio::time::pause();
4526
4527 let (mut producer, _pool) = pooled_producer(10_000);
4528 let mut subscriber = producer.subscribe(None);
4529
4530 finished_group(&mut producer, 10_000); let mut group0 = subscriber.assert_group();
4532
4533 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let read = group0.read_frame().await;
4537 assert!(matches!(read, Err(Error::Evicted)), "expected Evicted, got {read:?}");
4538 }
4539
4540 #[tokio::test]
4544 async fn small_writes_carry_debt() {
4545 tokio::time::pause();
4546
4547 let (mut producer, pool) = pooled_producer(22_000);
4548 let consumer = producer.consume();
4549
4550 finished_group(&mut producer, 20_000); for _ in 0..3 {
4555 finished_group(&mut producer, 1_000);
4556 }
4557 assert!(consumer.peek_group(0).is_some(), "debt smaller than the victim carries");
4558
4559 for _ in 0..20 {
4561 finished_group(&mut producer, 1_000);
4562 }
4563 assert!(
4564 consumer.peek_group(0).is_none(),
4565 "accumulated debt evicts the large group"
4566 );
4567 assert!(pool.used() <= 24_000, "usage hovers near capacity: {}", pool.used());
4570 }
4571
4572 #[tokio::test]
4576 async fn payment_capped_per_write() {
4577 tokio::time::pause();
4578
4579 let (mut producer, pool) = pooled_producer(1 << 40);
4580 for _ in 0..10 {
4581 finished_group(&mut producer, 1_000);
4582 }
4583
4584 pool.resize(100);
4586 let before = pool.used();
4587
4588 finished_group(&mut producer, 1_000);
4590
4591 let consumer = producer.consume();
4592 assert!(consumer.peek_group(0).is_none(), "the oldest groups are evicted");
4593 assert!(consumer.peek_group(1).is_none());
4594 assert!(consumer.peek_group(2).is_some(), "the backlog drains gradually");
4595 assert!(pool.used() > before - 4_000, "one write must not dump the backlog");
4596 }
4597
4598 #[tokio::test]
4602 async fn accept_preserves_write_accounting() {
4603 tokio::time::pause();
4604
4605 let pool = cache::Pool::new(12_000);
4606 let broadcast = broadcast::Info {
4607 origin: crate::origin::Info::default().with_pool(pool.clone()),
4608 ..Default::default()
4609 };
4610 let request = Request::new(Arc::new(broadcast), "test");
4611 let dynamic = request.dynamic();
4612 let consumer = request.consume();
4613
4614 let pending = consumer.fetch_group(0, None);
4616 let req = dynamic
4617 .requested_group()
4618 .now_or_never()
4619 .expect("should not block")
4620 .unwrap();
4621 let mut backfill = req.accept(None).unwrap();
4622 pending.await.unwrap();
4623 backfill
4624 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 30_000]))
4625 .unwrap();
4626
4627 let mut producer = request.accept(None);
4630 producer.append_group().unwrap().finish().unwrap();
4631 producer.append_group().unwrap().finish().unwrap();
4632
4633 assert!(
4634 producer.consume().peek_group(0).is_none(),
4635 "pre-accept backfill growth is reclaimed after accept"
4636 );
4637 assert!(pool.used() <= 13_000, "usage converges: {}", pool.used());
4638 }
4639
4640 #[tokio::test]
4643 async fn recreated_sequence_bounds_eviction_hints() {
4644 let (mut producer, _pool) = pooled_producer(1 << 40);
4645 producer.create_group(5u64.into()).unwrap().finish().unwrap();
4646
4647 for _ in 0..200 {
4648 let group = producer.create_group(1u64.into()).unwrap();
4649 group.abort(Error::Cancel).unwrap();
4650 }
4651
4652 let state = producer.state.read();
4653 assert!(
4654 state.evict.len() <= 2 * state.lookup.len() + EVICT_SLACK,
4655 "stale hints are compacted: {} entries for {} slots",
4656 state.evict.len(),
4657 state.lookup.len()
4658 );
4659 }
4660
4661 #[tokio::test]
4664 async fn same_tick_write_outranks_inserted() {
4665 tokio::time::pause();
4666
4667 let (mut producer, _pool) = pooled_producer(10_000);
4669
4670 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();
4677 assert!(consumer.peek_group(0).is_none(), "insert-only content pays first");
4678 assert!(consumer.peek_group(1).is_some(), "same-tick written content survives");
4679 }
4680
4681 #[tokio::test]
4684 async fn frame_only_writer_pays() {
4685 tokio::time::pause();
4686
4687 let (mut producer, pool) = pooled_producer(2_000);
4688 let mut demoted = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); demoted
4694 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 300_000]))
4695 .unwrap();
4696
4697 assert!(
4698 pool.used() <= 5_000,
4699 "the frame write settled the debt: {}",
4700 pool.used()
4701 );
4702 assert!(matches!(demoted.finish(), Err(Error::Evicted)));
4703 }
4704
4705 #[tokio::test]
4708 async fn each_track_owns_its_account() {
4709 let broadcast = Arc::new(broadcast::Info::default());
4710 let info = Info::default();
4711 let a = Producer::new(broadcast.clone(), "a", info.clone());
4712 let b = Producer::new(broadcast, "b", info);
4713
4714 let a = a.state.read().cache.clone();
4715 let b = b.state.read().cache.clone();
4716 assert!(!Arc::ptr_eq(&a, &b), "each track owns its account");
4717 }
4718
4719 #[tokio::test]
4722 async fn a_dynamic_defers_teardown() {
4723 let (mut producer, pool) = pooled_producer(1 << 40);
4724 let dynamic = producer.dynamic();
4725 finished_group(&mut producer, 100);
4726
4727 drop(producer);
4728 assert!(pool.used() > 0, "the handler still serves the cache");
4729
4730 drop(dynamic);
4731 assert_eq!(pool.used(), 0, "the last handle tears it down");
4732 }
4733
4734 #[tokio::test]
4740 async fn finished_track_frees_its_cache() {
4741 let (mut producer, pool) = pooled_producer(1 << 40);
4742 finished_group(&mut producer, 100);
4743 producer.finish().unwrap();
4744
4745 let state = producer.state.downgrade();
4746 drop(producer);
4747
4748 assert!(state.upgrade().is_none(), "the track state is freed");
4749 assert_eq!(pool.used(), 0, "so are its cached bytes");
4750 }
4751
4752 #[tokio::test]
4756 async fn teardown_ignores_a_settling_group() {
4757 let (mut producer, pool) = pooled_producer(1 << 40);
4758 finished_group(&mut producer, 100);
4759
4760 let settling = producer.state.downgrade().upgrade().expect("open");
4762 drop(producer);
4763
4764 assert_eq!(pool.used(), 0, "the abrupt teardown still released the cache");
4765 drop(settling);
4766 }
4767
4768 #[tokio::test]
4771 async fn cached_group_outlives_its_track() {
4772 let (mut producer, pool) = pooled_producer(1 << 40);
4773 let sequence = finished_group(&mut producer, 100);
4774 let group = producer.consume().peek_group(sequence).expect("cached");
4775 producer.finish().unwrap();
4776
4777 let state = producer.state.downgrade();
4778 drop(producer);
4779 assert!(state.upgrade().is_none(), "the track state is freed");
4780 assert!(pool.used() > 0, "the retained group keeps its own bytes");
4781
4782 drop(group);
4783 assert_eq!(pool.used(), 0, "which it releases when dropped");
4784 }
4785
4786 #[tokio::test]
4790 async fn pre_accept_backfill_settles_late_writes() {
4791 tokio::time::pause();
4792
4793 let pool = cache::Pool::new(2_000);
4794 let broadcast = broadcast::Info {
4795 origin: crate::origin::Info::default().with_pool(pool.clone()),
4796 ..Default::default()
4797 };
4798 let request = Request::new(Arc::new(broadcast), "test");
4799 let dynamic = request.dynamic();
4800 let consumer = request.consume();
4801
4802 let pending = consumer.fetch_group(0, None);
4804 let req = dynamic
4805 .requested_group()
4806 .now_or_never()
4807 .expect("should not block")
4808 .unwrap();
4809 let mut backfill = req.accept(None).unwrap();
4810 pending.await.unwrap();
4811
4812 let mut producer = request.accept(None);
4814 producer.append_group().unwrap().finish().unwrap();
4815
4816 backfill
4819 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 300_000]))
4820 .unwrap();
4821
4822 assert!(
4823 pool.used() <= 5_000,
4824 "the frame write settled the debt: {}",
4825 pool.used()
4826 );
4827 }
4828
4829 #[tokio::test]
4833 async fn write_restarts_retention_clock() {
4834 tokio::time::pause();
4835
4836 let (mut producer, _pool) = pooled_producer(1 << 40);
4837 let mut straggler = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
4842 straggler
4843 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 100]))
4844 .unwrap();
4845 producer.append_group().unwrap().finish().unwrap(); let consumer = producer.consume();
4848 assert!(consumer.peek_group(0).is_some(), "the write restarted the clock");
4849
4850 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
4852 producer.append_group().unwrap().finish().unwrap(); assert!(consumer.peek_group(0).is_none(), "idle content still expires");
4854 }
4855
4856 #[tokio::test]
4859 async fn refreshed_front_does_not_starve_expiry() {
4860 tokio::time::pause();
4861
4862 let (mut producer, _pool) = pooled_producer(1 << 40);
4863 let dynamic = producer.dynamic();
4864 let consumer = producer.consume();
4865
4866 producer.create_group(10u64.into()).unwrap().finish().unwrap();
4867 for sequence in 1..=5u64 {
4868 let pending = consumer.fetch_group(sequence, None);
4869 let req = dynamic
4870 .requested_group()
4871 .now_or_never()
4872 .expect("should not block")
4873 .unwrap();
4874 let mut group = req.accept(None).unwrap();
4875 group
4876 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 100]))
4877 .unwrap();
4878 group.finish().unwrap();
4879 pending.await.unwrap();
4880 }
4881
4882 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
4885 for sequence in 1..=4u64 {
4886 consumer.fetch_group(sequence, None).await.unwrap();
4887 }
4888
4889 for _ in 0..3 {
4891 producer.append_group().unwrap().finish().unwrap();
4892 }
4893 assert!(consumer.peek_group(5).is_none(), "expired backfill is reclaimed");
4894 assert!(consumer.peek_group(1).is_some(), "refreshed backfill survives");
4895 }
4896
4897 #[tokio::test]
4900 async fn recreated_sequence_delivered_once() {
4901 let (mut producer, _pool) = pooled_producer(1 << 40);
4902
4903 producer.create_group(0u64.into()).unwrap().finish().unwrap();
4904 let aborted = producer.create_group(1u64.into()).unwrap();
4905 aborted.abort(Error::Cancel).unwrap();
4906 producer.create_group(2u64.into()).unwrap().finish().unwrap();
4907 producer.create_group(1u64.into()).unwrap().finish().unwrap();
4908
4909 let mut subscriber = producer.subscribe(None);
4910 assert_eq!(subscriber.assert_group().sequence, 0);
4911 assert_eq!(subscriber.assert_group().sequence, 2);
4912 assert_eq!(
4913 subscriber.assert_group().sequence,
4914 1,
4915 "replacement arrives at its own position"
4916 );
4917 subscriber.assert_no_group();
4918 }
4919
4920 #[tokio::test]
4924 async fn datagrams_do_not_block_eviction() {
4925 tokio::time::pause();
4926
4927 let (mut producer, pool) = pooled_producer(1_000);
4928 for _ in 0..10 {
4929 finished_group(&mut producer, 1_000);
4930 producer.append_datagram(Timestamp::ZERO, &b"beat"[..]).unwrap();
4931 }
4932
4933 let consumer = producer.consume();
4934 assert!(consumer.peek_group(0).is_none(), "old groups still evict");
4935 assert!(
4936 pool.used() < 4 * 1_256,
4937 "interleaved datagrams must not bypass the budget: {}",
4938 pool.used()
4939 );
4940 }
4941
4942 #[tokio::test]
4946 async fn aborted_group_leaves_no_ghost_sample() {
4947 tokio::time::pause();
4948
4949 let (mut producer, pool) = pooled_producer(1 << 40);
4950 let group0 = producer.append_group().unwrap();
4951 producer.append_group().unwrap(); assert!(pool.average().is_some(), "demoted group is sampled");
4954 group0.abort(Error::Cancel).unwrap();
4955 assert_eq!(pool.average(), None, "the abort must remove the sample");
4956 }
4957
4958 #[tokio::test]
4961 async fn empty_groups_repay_overhead() {
4962 tokio::time::pause();
4963
4964 let (mut producer, pool) = pooled_producer(1_000);
4965 for _ in 0..100 {
4966 let mut group = producer.append_group().unwrap();
4967 group.finish().unwrap();
4968 }
4969
4970 assert!(
4971 pool.used() <= 3_000,
4972 "empty-group overhead must stay near the budget: {}",
4973 pool.used()
4974 );
4975 }
4976
4977 #[tokio::test]
4980 async fn growth_on_demoted_group_is_billed() {
4981 tokio::time::pause();
4982
4983 let (mut producer, pool) = pooled_producer(2_000);
4984 let mut straggler = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); straggler
4989 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 10_000]))
4990 .unwrap();
4991
4992 producer.append_group().unwrap().finish().unwrap(); let consumer = producer.consume();
4996 assert!(consumer.peek_group(0).is_none(), "the ballooned group is evicted");
4997 assert!(pool.used() <= 3_000, "growth is reclaimed: {}", pool.used());
4998 }
4999
5000 #[tokio::test]
5003 async fn refilled_sequence_stays_out_of_subscriptions() {
5004 let (mut producer, _pool) = pooled_producer(1 << 40);
5005 let dynamic = producer.dynamic();
5006 let consumer = producer.consume();
5007
5008 producer.create_group(0u64.into()).unwrap().finish().unwrap();
5009 let aborted = producer.create_group(1u64.into()).unwrap();
5010 aborted.abort(Error::Cancel).unwrap();
5011 producer.create_group(2u64.into()).unwrap().finish().unwrap();
5012
5013 let pending = consumer.fetch_group(1, None);
5016 let req = dynamic
5017 .requested_group()
5018 .now_or_never()
5019 .expect("should not block")
5020 .unwrap();
5021 let mut group = req.accept(None).unwrap();
5022 group
5023 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"backfill"))
5024 .unwrap();
5025 group.finish().unwrap();
5026 pending.await.unwrap();
5027
5028 assert!(consumer.peek_group(1).is_some());
5030 let mut subscriber = producer.subscribe(None);
5031 assert_eq!(subscriber.assert_group().sequence, 0);
5032 assert_eq!(subscriber.assert_group().sequence, 2);
5033 subscriber.assert_no_group();
5034 }
5035
5036 #[tokio::test]
5039 async fn expired_backfill_behind_refreshed_reclaimed() {
5040 tokio::time::pause();
5041
5042 let (mut producer, _pool) = pooled_producer(1 << 40);
5043 let dynamic = producer.dynamic();
5044 let consumer = producer.consume();
5045
5046 producer.create_group(5u64.into()).unwrap().finish().unwrap();
5047 for sequence in [2u64, 3u64] {
5048 let pending = consumer.fetch_group(sequence, None);
5049 let req = dynamic
5050 .requested_group()
5051 .now_or_never()
5052 .expect("should not block")
5053 .unwrap();
5054 let mut group = req.accept(None).unwrap();
5055 group
5056 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 100]))
5057 .unwrap();
5058 group.finish().unwrap();
5059 pending.await.unwrap();
5060 }
5061
5062 tokio::time::advance(Duration::from_secs(4)).await;
5064 consumer.fetch_group(2, None).await.unwrap();
5065 tokio::time::advance(DEFAULT_LATENCY_MAX - Duration::from_secs(2)).await;
5066 producer.create_group(6u64.into()).unwrap().finish().unwrap();
5067
5068 let consumer = producer.consume();
5069 assert!(consumer.peek_group(2).is_some(), "refreshed backfill survives");
5070 assert!(consumer.peek_group(3).is_none(), "expired backfill is reclaimed");
5071 }
5072
5073 #[tokio::test]
5076 async fn same_tick_fetch_protects() {
5077 tokio::time::pause();
5078
5079 let (mut producer, _pool) = pooled_producer(10_000);
5081 let consumer = producer.consume();
5082
5083 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();
5088
5089 finished_group(&mut producer, 3_000); finished_group(&mut producer, 3_000); assert!(consumer.peek_group(0).is_some(), "same-tick refresh protects");
5093 assert!(consumer.peek_group(1).is_none(), "the unread group dies instead");
5094 }
5095
5096 #[tokio::test]
5100 async fn refetched_latest_stays_protected() {
5101 tokio::time::pause();
5102
5103 let (mut producer, _pool) = pooled_producer(10_000);
5104 let dynamic = producer.dynamic();
5105 let consumer = producer.consume();
5106
5107 let straggler = producer.append_group().unwrap(); let latest = producer.append_group().unwrap(); latest.abort(Error::Cancel).unwrap();
5112
5113 let pending = consumer.fetch_group(1, None);
5115 let req = dynamic
5116 .requested_group()
5117 .now_or_never()
5118 .expect("should not block")
5119 .unwrap();
5120 let mut group = req.accept(None).unwrap();
5121 group
5122 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 1000]))
5123 .unwrap();
5124 group.finish().unwrap();
5125 pending.await.unwrap();
5126
5127 {
5130 let state = producer.state.read();
5131 assert!(state.lookup.contains_key(&1), "refetched group is cached");
5132 assert!(
5133 state.evict.iter().all(|(sequence, _)| *sequence != 1),
5134 "the live edge must not be an eviction candidate"
5135 );
5136 }
5137 drop(straggler);
5138 }
5139
5140 #[tokio::test]
5143 async fn eviction_allows_refetch() {
5144 tokio::time::pause();
5145
5146 let (mut producer, _pool) = pooled_producer(10_000);
5147 let dynamic = producer.dynamic();
5148
5149 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let consumer = producer.consume();
5154 assert!(consumer.peek_group(0).is_none());
5155 let pending = consumer.fetch_group(0, None);
5156
5157 let req = dynamic
5158 .requested_group()
5159 .now_or_never()
5160 .expect("should not block")
5161 .unwrap();
5162 assert_eq!(req.sequence(), 0);
5163
5164 let mut group = req.accept(None).unwrap();
5165 group
5166 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"refetched"))
5167 .unwrap();
5168 group.finish().unwrap();
5169
5170 let mut group = pending.await.unwrap();
5171 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"refetched");
5172 }
5173
5174 #[tokio::test]
5177 async fn fetched_backfill_not_subscribed() {
5178 let (mut producer, _pool) = pooled_producer(1 << 40);
5179 let dynamic = producer.dynamic();
5180 let consumer = producer.consume();
5181
5182 producer.create_group(5u64.into()).unwrap().finish().unwrap();
5184 producer.create_group(6u64.into()).unwrap().finish().unwrap();
5185
5186 let pending = consumer.fetch_group(2, None);
5188 let req = dynamic
5189 .requested_group()
5190 .now_or_never()
5191 .expect("should not block")
5192 .unwrap();
5193 let mut group = req.accept(None).unwrap();
5194 group
5195 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"backfill"))
5196 .unwrap();
5197 group.finish().unwrap();
5198 let mut fetched = pending.await.unwrap();
5199 assert_eq!(&fetched.read_frame().await.unwrap().unwrap().payload[..], b"backfill");
5200 assert!(consumer.peek_group(2).is_some(), "backfill is cached for later fetches");
5201
5202 let mut subscriber = producer.subscribe(None);
5204 assert_eq!(subscriber.assert_group().sequence, 5);
5205 assert_eq!(subscriber.assert_group().sequence, 6);
5206 subscriber.assert_no_group();
5207 }
5208
5209 #[tokio::test]
5212 async fn expired_backfill_reclaimed() {
5213 tokio::time::pause();
5214
5215 let (mut producer, pool) = pooled_producer(1 << 40);
5216 let dynamic = producer.dynamic();
5217 let consumer = producer.consume();
5218
5219 producer.create_group(5u64.into()).unwrap().finish().unwrap();
5220
5221 let pending = consumer.fetch_group(2, None);
5223 let req = dynamic
5224 .requested_group()
5225 .now_or_never()
5226 .expect("should not block")
5227 .unwrap();
5228 let mut group = req.accept(None).unwrap();
5229 group
5230 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 1000]))
5231 .unwrap();
5232 group.finish().unwrap();
5233 pending.await.unwrap();
5234 let used = pool.used();
5235
5236 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
5238 producer.create_group(6u64.into()).unwrap().finish().unwrap();
5239
5240 assert!(consumer.peek_group(2).is_none(), "expired backfill is reclaimed");
5241 assert!(pool.used() < used, "its bytes are released");
5242 }
5243
5244 #[tokio::test]
5245 async fn fetch_aborts_with_track() {
5246 let producer = track_producer("test", None);
5247 let dynamic = producer.dynamic();
5248 let consumer = producer.consume();
5249
5250 let pending = consumer.fetch_group(3, None);
5251 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
5252
5253 producer.abort(Error::Cancel).unwrap();
5254 assert!(pending.await.is_err());
5255 drop(dynamic);
5256 }
5257}