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 frame = crate::IntoBytes::into_bytes(frame);
1041 if frame.len() as u64 > group::MAX_GROUP_CACHE {
1042 return Err(Error::FrameTooLarge);
1043 }
1044 let mut group = self.append_group()?;
1045 group.write_frame(timestamp, frame)?;
1046 group.finish()?;
1047 Ok(())
1048 }
1049
1050 pub fn finish(&mut self) -> Result<()> {
1056 let mut state = self.modify()?;
1057 let final_sequence = match state.max_sequence {
1058 Some(max) => max.checked_add(1).ok_or(coding::BoundsExceeded)?,
1059 None => 0,
1060 };
1061 state.set_final(final_sequence)
1062 }
1063
1064 pub fn finish_at(&mut self, final_sequence: u64) -> Result<()> {
1077 self.modify()?.set_final(final_sequence)
1078 }
1079
1080 pub fn final_sequence(&self) -> Option<u64> {
1085 self.state.read().final_sequence
1086 }
1087
1088 pub fn abort(self, err: Error) -> Result<()> {
1099 let mut guard = self.modify()?;
1100 guard.abort = Some(err);
1101 guard.clear_cache();
1102 guard.datagrams.clear();
1103 guard.close();
1104 Ok(())
1105 }
1106
1107 pub async fn unused(&self) -> Result<()> {
1109 self.state.unused().await.map_err(|_| self.abort_reason())
1110 }
1111
1112 pub async fn used(&self) -> Result<()> {
1114 self.state.used().await.map_err(|_| self.abort_reason())
1115 }
1116
1117 pub async fn closed(&self) -> Error {
1119 kio::wait(|waiter| self.poll_closed(waiter)).await
1120 }
1121
1122 pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<Error> {
1124 self.state.poll_closed(waiter).map(|()| self.abort_reason())
1125 }
1126
1127 fn abort_reason(&self) -> Error {
1129 self.state.read().abort.clone().unwrap_or(Error::Dropped)
1130 }
1131
1132 pub fn is_closed(&self) -> bool {
1134 self.state.read().is_closed()
1135 }
1136
1137 pub fn latest(&self) -> Option<u64> {
1139 self.state.read().max_sequence
1140 }
1141
1142 pub fn is_clone(&self, other: &Self) -> bool {
1144 self.state.same_channel(&other.state)
1145 }
1146
1147 pub(crate) fn weak(&self) -> TrackWeak {
1149 TrackWeak {
1150 name: self.name.clone(),
1151 state: self.state.weak(),
1152 }
1153 }
1154
1155 pub fn demand(&self) -> Demand {
1163 Demand {
1164 name: self.name.clone(),
1165 state: self.state.weak(),
1166 }
1167 }
1168
1169 pub fn consume(&self) -> Consumer {
1174 Consumer::plain(self.name.clone(), self.state.consume())
1175 }
1176
1177 pub fn subscribe(&self, subscription: impl Into<Option<Subscription>>) -> Subscriber {
1182 let preferences = subscription.into().unwrap_or_default();
1183
1184 let info = self.state.read().info.clone().expect("producer always has info");
1189 let subscription = kio::Producer::new(preferences);
1190 register_subscription(self.state.read(), &subscription);
1191
1192 Subscriber {
1193 name: self.name.clone(),
1194 info,
1195 inner: SubscriberKind::Plain(PlainSubscriber {
1196 state: self.state.consume(),
1197 subscription,
1198 index: 0,
1199 datagram_index: 0,
1200 min_sequence: 0,
1201 next_sequence: 0,
1202 end_sequence: None,
1203 parked: BTreeMap::new(),
1204 }),
1205 stats: stats::Scope::default(),
1207 _stats_sub: stats::Subscription::default(),
1208 }
1209 }
1210
1211 pub async fn subscription_changed(&mut self) -> Result<Option<Subscription>> {
1217 kio::wait(|waiter| self.poll_subscription_changed(waiter)).await
1218 }
1219
1220 pub fn subscription(&self) -> Option<Subscription> {
1228 let state = self.state.read();
1229 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
1230 drop(state);
1231 snapshot_subscription(&subs, bound)
1232 }
1233
1234 pub fn poll_subscription_changed(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Subscription>>> {
1238 if self.state.poll_closed(waiter).is_ready() {
1241 let abort = self.state.read().abort.clone();
1242 return Poll::Ready(Err(abort.unwrap_or(Error::Dropped)));
1243 }
1244
1245 let state = self.state.read();
1247 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
1248 drop(state);
1249
1250 let prev = &self.prev_subscription;
1251 let mut combined = None;
1252 let mut guard = match subs.poll(waiter, |subs| {
1253 let next = combined_subscription(subs, bound, waiter);
1254 if &next == prev {
1255 Poll::Pending
1256 } else {
1257 combined = next;
1258 Poll::Ready(())
1259 }
1260 }) {
1261 Poll::Ready(guard) => guard,
1262 Poll::Pending => return Poll::Pending,
1263 };
1264 guard.retain(|sub| !sub.is_closed());
1266 drop(guard);
1267 self.prev_subscription = combined.clone();
1268 Poll::Ready(Ok(combined))
1269 }
1270
1271 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
1273 self.state.poll_unused(waiter).map(|_| ())
1274 }
1275
1276 pub fn dynamic(&self) -> Dynamic {
1280 Dynamic::new(self.name.clone(), self.state.clone(), self.alive.clone())
1281 }
1282
1283 fn modify(&self) -> Result<kio::Mut<'_, TrackState>> {
1284 TrackState::modify(&self.state)
1285 }
1286}
1287
1288fn poll_requested_group(
1292 state: &kio::Producer<TrackState>,
1293 fetch: &kio::Shared<FetchState>,
1294 waiter: &kio::Waiter,
1295) -> Poll<Result<GroupRequest>> {
1296 if let Poll::Ready(mut guard) = fetch.poll(waiter, |fetch| {
1298 if fetch.has_queued() {
1299 Poll::Ready(())
1300 } else {
1301 Poll::Pending
1302 }
1303 }) {
1304 let sequence = guard.pop().expect("predicate guaranteed a request");
1305 let pending = guard.get(&sequence).expect("popped key must be pending");
1309 let priority = pending.priority;
1310 let result = pending.result.clone();
1311 drop(guard);
1312 return Poll::Ready(Ok(GroupRequest {
1313 state: state.clone(),
1314 fetch: fetch.clone(),
1315 sequence,
1316 priority,
1317 result,
1318 done: false,
1319 }));
1320 }
1321
1322 match state.poll_ref(waiter, |state| match &state.abort {
1324 Some(err) => Poll::Ready(err.clone()),
1325 None => Poll::Pending,
1326 }) {
1327 Poll::Ready(Ok(err)) => Poll::Ready(Err(err)),
1328 Poll::Ready(Err(closed)) => Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped))),
1329 Poll::Pending => Poll::Pending,
1330 }
1331}
1332
1333pub struct Dynamic {
1343 name: Arc<str>,
1344 state: kio::Producer<TrackState>,
1346 fetch: kio::Shared<FetchState>,
1348 alive: Arc<Alive>,
1351}
1352
1353impl Dynamic {
1354 fn new(name: Arc<str>, state: kio::Producer<TrackState>, alive: Arc<Alive>) -> Self {
1355 let fetch = state.read().fetch.clone();
1356 fetch.lock().add_handler();
1357 Self {
1358 name,
1359 state,
1360 fetch,
1361 alive,
1362 }
1363 }
1364
1365 pub fn name(&self) -> &str {
1367 &self.name
1368 }
1369
1370 pub async fn requested_group(&self) -> Result<GroupRequest> {
1376 kio::wait(|waiter| self.poll_requested_group(waiter)).await
1377 }
1378
1379 pub fn poll_requested_group(&self, waiter: &kio::Waiter) -> Poll<Result<GroupRequest>> {
1381 poll_requested_group(&self.state, &self.fetch, waiter)
1382 }
1383
1384 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
1386 self.state.poll_unused(waiter).map(|_| ())
1387 }
1388}
1389
1390impl Clone for Dynamic {
1391 fn clone(&self) -> Self {
1392 self.fetch.lock().add_handler();
1394 Self {
1395 name: self.name.clone(),
1396 state: self.state.clone(),
1397 fetch: self.fetch.clone(),
1398 alive: self.alive.clone(),
1399 }
1400 }
1401}
1402
1403impl Drop for Dynamic {
1404 fn drop(&mut self) {
1405 let mut fetch = self.fetch.lock();
1411 if fetch.remove_handler() {
1412 fetch.drain_queued();
1413 }
1414 }
1415}
1416
1417struct Alive {
1426 name: Arc<str>,
1427 state: kio::Producer<TrackState>,
1428
1429 published: AtomicBool,
1432
1433 stats: OnceLock<stats::Subscription>,
1436}
1437
1438impl Alive {
1439 fn new(name: Arc<str>, state: kio::Producer<TrackState>) -> Arc<Self> {
1440 Arc::new(Self {
1441 name,
1442 state,
1443 published: Default::default(),
1444 stats: Default::default(),
1445 })
1446 }
1447
1448 fn publish(&self, stats: Option<&stats::Scope>) {
1452 self.published.store(true, Ordering::Relaxed);
1453 if let Some(scope) = stats {
1454 let _ = self.stats.set(scope.subscribe());
1457 }
1458 }
1459}
1460
1461impl Drop for Alive {
1462 fn drop(&mut self) {
1463 if !self.published.load(Ordering::Relaxed) {
1465 return;
1466 }
1467 match self.state.write() {
1475 Ok(mut state) => {
1476 if state.final_sequence.is_some() || state.abort.is_some() {
1477 return;
1478 }
1479 tracing::warn!(
1480 track = %self.name,
1481 "track::Producer dropped without finish() or abort()"
1482 );
1483 state.clear_cache();
1484 state.datagrams.clear();
1485 }
1486 Err(state) => {
1487 if state.final_sequence.is_some() || state.abort.is_some() {
1488 return;
1489 }
1490 tracing::warn!(
1491 track = %self.name,
1492 "track::Producer dropped without finish() or abort()"
1493 );
1494 }
1495 }
1496 }
1497}
1498
1499fn combined_subscription(subs: &Subscriptions, bound: Option<Duration>, waiter: &kio::Waiter) -> Option<Subscription> {
1505 let mut combined = None;
1506 for sub in subs.iter() {
1507 if sub.is_closed() {
1512 continue;
1513 }
1514 let _ = sub.poll_closed(waiter);
1519 if let Poll::Ready(Ok(sub)) = sub.poll(waiter, |sub| sub.poll_combined(&combined)) {
1520 combined = Some(sub);
1521 }
1522 }
1523 clamp_combined(combined, bound)
1524}
1525
1526fn snapshot_subscription(subs: &kio::Shared<Subscriptions>, bound: Option<Duration>) -> Option<Subscription> {
1528 let mut combined: Option<Subscription> = None;
1529 for sub in subs.read().iter() {
1530 if sub.is_closed() {
1532 continue;
1533 }
1534 if let Poll::Ready(merged) = sub.read().poll_combined(&combined) {
1535 combined = Some(merged);
1536 }
1537 }
1538 clamp_combined(combined, bound)
1539}
1540
1541fn clamp_combined(combined: Option<Subscription>, bound: Option<Duration>) -> Option<Subscription> {
1549 let mut combined = combined?;
1550 if let Some(bound) = bound {
1551 combined.latency_max = combined.latency_max.min(bound);
1552 }
1553 Some(combined)
1554}
1555
1556fn register_subscription(state: kio::Ref<'_, TrackState>, subscription: &kio::Producer<Subscription>) {
1560 if state.is_closed() {
1561 return;
1562 }
1563 let subs = state.subscriptions.clone();
1564 drop(state);
1565 subs.lock().push(subscription.consume());
1566}
1567
1568#[derive(Clone)]
1570pub(crate) struct TrackWeak {
1571 name: Arc<str>,
1572 state: kio::ProducerWeak<TrackState>,
1573}
1574
1575impl TrackWeak {
1576 pub fn consume(&self) -> Consumer {
1577 Consumer::plain(self.name.clone(), self.state.consume())
1578 }
1579
1580 pub(crate) fn name(&self) -> &Arc<str> {
1583 &self.name
1584 }
1585
1586 pub(crate) fn is_used(&self) -> bool {
1589 !self.state.is_closed() && self.state.is_used()
1590 }
1591
1592 pub(crate) fn poll_used(&self, waiter: &kio::Waiter) {
1595 let _ = self.state.poll_used(waiter);
1596 }
1597
1598 pub(crate) fn poll_unused(&self, waiter: &kio::Waiter) {
1601 let _ = self.state.poll_unused(waiter);
1602 }
1603}
1604
1605impl super::WeakEntry for TrackWeak {
1606 fn is_closed(&self) -> bool {
1607 self.state.is_closed()
1608 }
1609
1610 fn same_channel(&self, other: &Self) -> bool {
1611 self.state.same_channel(&other.state)
1612 }
1613}
1614
1615#[derive(Clone)]
1624pub struct Demand {
1625 name: Arc<str>,
1626 state: kio::ProducerWeak<TrackState>,
1627}
1628
1629impl Demand {
1630 pub fn name(&self) -> &str {
1632 &self.name
1633 }
1634
1635 pub async fn used(&self) -> Result<()> {
1637 self.state.used().await.map_err(|_| self.abort_reason())
1638 }
1639
1640 pub async fn unused(&self) -> Result<()> {
1642 self.state.unused().await.map_err(|_| self.abort_reason())
1643 }
1644
1645 pub async fn closed(&self) -> Error {
1647 self.state.closed().await;
1648 self.abort_reason()
1649 }
1650
1651 fn abort_reason(&self) -> Error {
1653 self.state.read().abort.clone().unwrap_or(Error::Dropped)
1654 }
1655}
1656
1657#[derive(Clone)]
1668pub struct Consumer {
1669 name: Arc<str>,
1670 inner: ConsumerKind,
1671 stats: stats::Scope,
1674}
1675
1676#[derive(Clone)]
1677enum ConsumerKind {
1678 Plain(kio::Consumer<TrackState>),
1679 Spliced(super::resume::Consumer),
1680}
1681
1682impl Consumer {
1683 fn plain(name: Arc<str>, state: kio::Consumer<TrackState>) -> Self {
1684 Self {
1685 name,
1686 inner: ConsumerKind::Plain(state),
1687 stats: stats::Scope::default(),
1688 }
1689 }
1690
1691 pub(crate) fn spliced(name: Arc<str>, resume: super::resume::Consumer) -> Self {
1693 Self {
1694 name,
1695 inner: ConsumerKind::Spliced(resume),
1696 stats: stats::Scope::default(),
1697 }
1698 }
1699
1700 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
1703 self.stats = scope;
1704 self
1705 }
1706
1707 pub fn name(&self) -> &str {
1709 &self.name
1710 }
1711
1712 pub fn subscribe(&self, subscription: impl Into<Option<Subscription>>) -> kio::Pending<Subscribing> {
1718 let subscription = kio::Producer::new(subscription.into().unwrap_or_default());
1719
1720 let inner = match &self.inner {
1721 ConsumerKind::Plain(state) => {
1722 register_subscription(state.read(), &subscription);
1725 SubscribingKind::Plain(state.clone())
1726 }
1727 ConsumerKind::Spliced(resume) => SubscribingKind::Spliced(resume.clone()),
1729 };
1730
1731 kio::Pending::new(Subscribing {
1732 name: self.name.clone(),
1733 inner,
1734 subscription,
1735 stats: self.stats.clone(),
1736 })
1737 }
1738
1739 #[cfg(test)]
1743 pub(crate) fn peek_group(&self, sequence: u64) -> Option<group::Consumer> {
1744 match &self.inner {
1745 ConsumerKind::Plain(state) => state.read().cached_group(sequence),
1746 ConsumerKind::Spliced(_) => None,
1749 }
1750 }
1751
1752 pub fn fetch_group(&self, sequence: u64, options: impl Into<Option<group::Fetch>>) -> kio::Pending<Fetching> {
1764 let options = options.into().unwrap_or_default();
1765
1766 self.stats.fetch();
1770
1771 let state = match &self.inner {
1772 ConsumerKind::Plain(state) => state,
1773 ConsumerKind::Spliced(resume) => {
1776 return kio::Pending::new(Fetching {
1777 inner: FetchingKind::Spliced(resume.fetch_group(sequence, options)),
1778 stats: self.stats.clone(),
1779 });
1780 }
1781 };
1782
1783 let mut result = None;
1784
1785 let (fetch, unresolved) = {
1789 let state = state.read();
1790 (state.fetch.clone(), state.poll_fetch_cached(sequence).is_pending())
1791 };
1792
1793 if unresolved {
1794 let mut fetch = fetch.lock();
1795 if let Some(pending) = fetch.join(&sequence) {
1796 pending.priority = pending.priority.max(options.priority);
1799 result = Some(pending.result.consume());
1800 } else {
1801 let producer = kio::Producer::<FetchOutcome>::default();
1805 let consumer = producer.consume();
1806 let attempt = PendingFetch {
1807 priority: options.priority,
1808 result: producer,
1809 };
1810 if fetch.insert(sequence, attempt).is_ok() {
1811 result = Some(consumer);
1812 }
1813 }
1814 }
1815
1816 kio::Pending::new(Fetching {
1817 inner: FetchingKind::Plain {
1818 state: state.clone(),
1819 fetch,
1820 sequence,
1821 result,
1822 },
1823 stats: self.stats.clone(),
1824 })
1825 }
1826
1827 pub fn info(&self) -> kio::Pending<Querying> {
1834 kio::Pending::new(Querying {
1835 inner: match &self.inner {
1836 ConsumerKind::Plain(state) => QueryingKind::Plain(state.clone()),
1837 ConsumerKind::Spliced(resume) => QueryingKind::Spliced(resume.clone()),
1838 },
1839 })
1840 }
1841
1842 pub fn latest(&self) -> Option<u64> {
1844 match &self.inner {
1845 ConsumerKind::Plain(state) => state.read().max_sequence,
1846 ConsumerKind::Spliced(resume) => resume.latest(),
1847 }
1848 }
1849
1850 pub(crate) fn poll_complete(&self, waiter: &kio::Waiter) -> Poll<Result<()>> {
1855 let ConsumerKind::Plain(state) = &self.inner else {
1856 return Poll::Pending;
1858 };
1859 match ready!(state.poll(waiter, |state| {
1860 if state.is_complete() {
1861 Poll::Ready(())
1862 } else {
1863 Poll::Pending
1864 }
1865 })) {
1866 Ok(_) => Poll::Ready(Ok(())),
1867 Err(closed) => Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped))),
1870 }
1871 }
1872}
1873
1874pub struct Subscribing {
1877 name: Arc<str>,
1878 inner: SubscribingKind,
1879 subscription: kio::Producer<Subscription>,
1880 stats: stats::Scope,
1881}
1882
1883enum SubscribingKind {
1884 Plain(kio::Consumer<TrackState>),
1885 Spliced(super::resume::Consumer),
1886}
1887
1888impl Subscribing {
1889 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<Subscriber>> {
1892 match &self.inner {
1893 SubscribingKind::Plain(state) => {
1894 let info = ready!(state.poll(waiter, |state| state.poll_info()))
1896 .map_err(|e| e.abort.clone().unwrap_or(Error::Dropped))??;
1897
1898 Poll::Ready(Ok(Subscriber {
1899 name: self.name.clone(),
1900 info,
1901 inner: SubscriberKind::Plain(PlainSubscriber {
1902 state: state.clone(),
1903 subscription: self.subscription.clone(),
1904 index: 0,
1905 datagram_index: 0,
1906 min_sequence: 0,
1907 next_sequence: 0,
1908 end_sequence: None,
1909 parked: BTreeMap::new(),
1910 }),
1911 stats: self.stats.clone(),
1912 _stats_sub: self.stats.subscribe(),
1913 }))
1914 }
1915 SubscribingKind::Spliced(resume) => {
1916 let info = ready!(resume.poll_info(waiter))?;
1919
1920 Poll::Ready(Ok(Subscriber {
1921 name: self.name.clone(),
1922 info,
1923 inner: SubscriberKind::Spliced(Box::new(resume.subscribe_shared(self.subscription.clone()))),
1924 stats: self.stats.clone(),
1925 _stats_sub: self.stats.subscribe(),
1926 }))
1927 }
1928 }
1929 }
1930
1931 pub fn update(&mut self, subscription: Subscription) -> Result<()> {
1936 let mut state = self.subscription.write().map_err(|_| Error::Closed)?;
1937 *state = subscription;
1938 Ok(())
1939 }
1940}
1941
1942impl kio::Pollable for Subscribing {
1943 type Output = Result<Subscriber>;
1944
1945 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
1946 self.poll_ok(waiter)
1947 }
1948}
1949
1950pub struct Querying {
1953 inner: QueryingKind,
1954}
1955
1956enum QueryingKind {
1957 Plain(kio::Consumer<TrackState>),
1958 Spliced(super::resume::Consumer),
1959}
1960
1961impl Querying {
1962 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<Info>> {
1964 match &self.inner {
1965 QueryingKind::Plain(state) => {
1966 let info = ready!(state.poll(waiter, |state| state.poll_info()))
1968 .map_err(|e| e.abort.clone().unwrap_or(Error::Dropped))??;
1969 Poll::Ready(Ok(info))
1970 }
1971 QueryingKind::Spliced(resume) => resume.poll_info(waiter),
1972 }
1973 }
1974}
1975
1976impl kio::Pollable for Querying {
1977 type Output = Result<Info>;
1978
1979 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
1980 self.poll_ok(waiter)
1981 }
1982}
1983
1984pub struct GroupRequest {
1993 state: kio::Producer<TrackState>,
1994 fetch: kio::Shared<FetchState>,
1996 sequence: u64,
1997 priority: u8,
1998 result: kio::Producer<FetchOutcome>,
2000 done: bool,
2001}
2002
2003impl GroupRequest {
2004 pub fn sequence(&self) -> u64 {
2006 self.sequence
2007 }
2008
2009 pub fn priority(&self) -> u8 {
2011 self.priority
2012 }
2013
2014 pub fn accept(mut self, info: impl Into<Option<Info>>) -> Result<group::Producer> {
2022 self.done = true;
2023 let res = TrackState::modify(&self.state)
2027 .and_then(|mut state| state.insert_group_request(self.sequence, info.into()));
2028 self.remove();
2029 res
2030 }
2031
2032 pub fn reject(mut self, err: Error) {
2034 self.done = true;
2035 self.remove();
2038 if let Ok(mut outcome) = self.result.write() {
2039 outcome.rejected = Some(err);
2040 }
2041 }
2042
2043 fn remove(&self) {
2046 self.fetch
2047 .lock()
2048 .remove_if(&self.sequence, |pending| pending.result.same_channel(&self.result));
2049 }
2050}
2051
2052impl Drop for GroupRequest {
2053 fn drop(&mut self) {
2054 if self.done {
2055 return;
2056 }
2057 self.remove();
2058 if let Ok(mut outcome) = self.result.write() {
2059 outcome.rejected = Some(Error::Dropped);
2060 }
2061 }
2062}
2063
2064pub struct Fetching {
2070 inner: FetchingKind,
2071 stats: stats::Scope,
2074}
2075
2076enum FetchingKind {
2077 Plain {
2078 state: kio::Consumer<TrackState>,
2079 fetch: kio::Shared<FetchState>,
2080 sequence: u64,
2081 result: Option<kio::Consumer<FetchOutcome>>,
2083 },
2084 Spliced(kio::Pending<super::resume::Fetching>),
2086}
2087
2088impl kio::Pollable for Fetching {
2089 type Output = Result<group::Consumer>;
2090
2091 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
2092 let (state, fetch, sequence, result) = match &self.inner {
2093 FetchingKind::Plain {
2094 state,
2095 fetch,
2096 sequence,
2097 result,
2098 } => (state, fetch, *sequence, result.as_ref()),
2099 FetchingKind::Spliced(spliced) => {
2100 return kio::Pollable::poll(&**spliced, waiter)
2103 .map(|res| res.map(|group| group.with_meter(self.stats.meter())));
2104 }
2105 };
2106
2107 match state.poll(waiter, |state| state.poll_fetch_cached(sequence)) {
2110 Poll::Ready(Ok(res)) => return Poll::Ready(res.map(|group| group.with_meter(self.stats.meter()))),
2111 Poll::Ready(Err(closed)) => {
2112 return Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped)));
2113 }
2114 Poll::Pending => {}
2115 }
2116
2117 let Some(result) = result else {
2119 return match fetch.poll(waiter, |fetch| match fetch.has_handlers() {
2122 false => Poll::Ready(()),
2123 true => Poll::Pending,
2124 }) {
2125 Poll::Ready(_guard) => Poll::Ready(Err(Error::NotFound)),
2126 Poll::Pending => Poll::Pending,
2127 };
2128 };
2129
2130 match result.poll(waiter, |outcome| match &outcome.rejected {
2133 Some(err) => Poll::Ready(err.clone()),
2134 None => Poll::Pending,
2135 }) {
2136 Poll::Ready(Ok(err)) => Poll::Ready(Err(err)),
2137 Poll::Ready(Err(_closed)) => Poll::Ready(Err(Error::NotFound)),
2138 Poll::Pending => Poll::Pending,
2139 }
2140 }
2141}
2142
2143pub struct Subscriber {
2166 name: Arc<str>,
2167 info: Info,
2168 inner: SubscriberKind,
2169 stats: stats::Scope,
2172 _stats_sub: stats::Subscription,
2175}
2176
2177enum SubscriberKind {
2178 Plain(PlainSubscriber),
2179 Spliced(Box<super::resume::Subscriber>),
2181}
2182
2183struct PlainSubscriber {
2185 state: kio::Consumer<TrackState>,
2186
2187 subscription: kio::Producer<Subscription>,
2188 index: usize,
2190 datagram_index: usize,
2192 min_sequence: u64,
2194 next_sequence: u64,
2197 end_sequence: Option<u64>,
2202 parked: BTreeMap<u64, group::Consumer>,
2207}
2208
2209impl PlainSubscriber {
2210 fn poll<F, R>(&self, waiter: &kio::Waiter, f: F) -> Poll<Result<R>>
2212 where
2213 F: Fn(&kio::Ref<'_, TrackState>) -> Poll<Result<R>>,
2214 {
2215 Poll::Ready(match ready!(self.state.poll(waiter, f)) {
2216 Ok(res) => res,
2217 Err(state) => Err(state.abort.clone().unwrap_or(Error::Dropped)),
2219 })
2220 }
2221
2222 fn poll_recv_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2223 let watch = |group: &group::Consumer| match group.poll_closed(waiter) {
2231 Poll::Pending => true,
2232 Poll::Ready(()) => !group.is_aborted(),
2233 };
2234
2235 let min_sequence = self.min_sequence;
2240 self.parked
2241 .retain(|sequence, group| *sequence >= min_sequence && watch(group));
2242
2243 if let Some(&sequence) = self.parked.keys().next()
2245 && self.end_sequence.is_none_or(|end| sequence <= end)
2246 {
2247 return Poll::Ready(Ok(self.parked.remove(&sequence)));
2248 }
2249
2250 loop {
2251 let Some((consumer, found_index)) =
2252 ready!(self.poll(waiter, |state| state.poll_recv_group(self.index, self.min_sequence))?)
2253 else {
2254 if self.parked.is_empty() {
2257 return Poll::Ready(Ok(None));
2258 }
2259 return Poll::Pending;
2260 };
2261 self.index = found_index + 1;
2262
2263 if self.end_sequence.is_some_and(|end| consumer.sequence > end) {
2266 if watch(&consumer) {
2270 self.parked.insert(consumer.sequence, consumer);
2271 }
2272 continue;
2273 }
2274 return Poll::Ready(Ok(Some(consumer)));
2275 }
2276 }
2277
2278 fn poll_recv_datagram(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Datagram>>> {
2279 let Some((datagram, found_index)) =
2280 ready!(self.poll(waiter, |state| state.poll_recv_datagram(self.datagram_index))?)
2281 else {
2282 return Poll::Ready(Ok(None));
2283 };
2284
2285 self.datagram_index = found_index + 1;
2286 self.next_sequence = self.next_sequence.max(datagram.sequence.saturating_add(1));
2287 Poll::Ready(Ok(Some(datagram)))
2288 }
2289
2290 fn poll_next_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2291 let floor = self.next_sequence.max(self.min_sequence);
2292 let Some(group) = ready!(self.poll(waiter, |state| state.poll_next_in_range(floor, self.end_sequence))?) else {
2293 return Poll::Ready(Ok(None));
2294 };
2295 self.next_sequence = group.sequence.saturating_add(1);
2296 Poll::Ready(Ok(Some(group)))
2297 }
2298
2299 fn poll_read_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Frame>>> {
2300 let lower = self.min_sequence.max(self.next_sequence);
2301 let Some((frame, found_index, sequence)) =
2302 ready!(self.poll(waiter, |state| { state.poll_read_frame(self.index, lower, waiter) })?)
2303 else {
2304 return Poll::Ready(Ok(None));
2305 };
2306
2307 self.index = found_index + 1;
2308 self.next_sequence = sequence.saturating_add(1);
2309 Poll::Ready(Ok(Some(frame)))
2310 }
2311}
2312
2313#[derive(Clone)]
2319pub struct SubscriberControl {
2320 subscription: kio::Producer<Subscription>,
2321}
2322
2323impl SubscriberControl {
2324 pub fn subscription(&self) -> Subscription {
2326 self.subscription.read().clone()
2327 }
2328
2329 pub fn update(&self, subscription: Subscription) -> Result<()> {
2334 let mut state = self.subscription.write().map_err(|_| Error::Closed)?;
2335 *state = subscription;
2336 Ok(())
2337 }
2338}
2339
2340impl Subscriber {
2341 pub fn info(&self) -> &Info {
2346 &self.info
2347 }
2348
2349 pub fn name(&self) -> &str {
2351 &self.name
2352 }
2353
2354 pub fn control(&self) -> SubscriberControl {
2356 SubscriberControl {
2357 subscription: match &self.inner {
2358 SubscriberKind::Plain(plain) => plain.subscription.clone(),
2359 SubscriberKind::Spliced(spliced) => spliced.prefs(),
2360 },
2361 }
2362 }
2363
2364 pub fn poll_recv_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2381 let meter = self.stats.meter();
2382 let res = match &mut self.inner {
2383 SubscriberKind::Plain(plain) => plain.poll_recv_group(waiter),
2384 SubscriberKind::Spliced(spliced) => spliced.poll_recv_group(waiter),
2385 };
2386 res.map(|res| res.map(|group| group.map(|group| group.with_meter(meter))))
2387 }
2388
2389 pub async fn recv_group(&mut self) -> Result<Option<group::Consumer>> {
2396 kio::wait(|waiter| self.poll_recv_group(waiter)).await
2397 }
2398
2399 pub fn poll_recv_datagram(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Datagram>>> {
2410 let meter = self.stats.meter();
2411 let res = match &mut self.inner {
2412 SubscriberKind::Plain(plain) => plain.poll_recv_datagram(waiter),
2413 SubscriberKind::Spliced(spliced) => spliced.poll_recv_datagram(waiter),
2414 };
2415 if let Poll::Ready(Ok(Some(datagram))) = &res {
2418 meter.datagram(datagram.payload.len() as u64);
2419 }
2420 res
2421 }
2422
2423 pub async fn recv_datagram(&mut self) -> Result<Option<Datagram>> {
2430 kio::wait(|waiter| self.poll_recv_datagram(waiter)).await
2431 }
2432
2433 pub fn poll_next_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2442 let meter = self.stats.meter();
2443 let res = match &mut self.inner {
2444 SubscriberKind::Plain(plain) => plain.poll_next_group(waiter),
2445 SubscriberKind::Spliced(spliced) => spliced.poll_next_group(waiter),
2446 };
2447 res.map(|res| res.map(|group| group.map(|group| group.with_meter(meter))))
2448 }
2449
2450 pub async fn next_group(&mut self) -> Result<Option<group::Consumer>> {
2456 kio::wait(|waiter| self.poll_next_group(waiter)).await
2457 }
2458
2459 pub fn poll_read_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Frame>>> {
2463 let meter = self.stats.meter();
2464 let res = match &mut self.inner {
2465 SubscriberKind::Plain(plain) => plain.poll_read_frame(waiter),
2466 SubscriberKind::Spliced(spliced) => spliced.poll_read_frame(waiter),
2467 };
2468 if let Poll::Ready(Ok(Some(frame))) = &res {
2471 meter.group();
2472 meter.frames(1);
2473 meter.bytes(frame.payload.len() as u64);
2474 }
2475 res
2476 }
2477
2478 pub async fn read_frame(&mut self) -> Result<Option<frame::Frame>> {
2483 kio::wait(|waiter| self.poll_read_frame(waiter)).await
2484 }
2485
2486 pub fn is_clone(&self, other: &Self) -> bool {
2488 match (&self.inner, &other.inner) {
2489 (SubscriberKind::Plain(a), SubscriberKind::Plain(b)) => a.state.same_channel(&b.state),
2490 (SubscriberKind::Spliced(a), SubscriberKind::Spliced(b)) => a.is_clone(b),
2491 _ => false,
2492 }
2493 }
2494
2495 pub fn poll_finished(&mut self, waiter: &kio::Waiter) -> Poll<Result<u64>> {
2497 match &mut self.inner {
2498 SubscriberKind::Plain(plain) => plain.poll(waiter, |state| state.poll_finished()),
2499 SubscriberKind::Spliced(spliced) => spliced.poll_finished(waiter),
2500 }
2501 }
2502
2503 pub async fn finished(&mut self) -> Result<u64> {
2511 kio::wait(|waiter| self.poll_finished(waiter)).await
2512 }
2513
2514 pub fn start_at(&mut self, sequence: u64) {
2521 match &mut self.inner {
2522 SubscriberKind::Plain(plain) => plain.min_sequence = sequence,
2523 SubscriberKind::Spliced(spliced) => spliced.start_at(sequence),
2524 }
2525 }
2526
2527 pub fn end_at(&mut self, sequence: impl Into<Option<u64>>) {
2540 match &mut self.inner {
2541 SubscriberKind::Plain(plain) => plain.end_sequence = sequence.into(),
2542 SubscriberKind::Spliced(spliced) => spliced.end_at(sequence),
2543 }
2544 }
2545
2546 pub fn subscription(&self) -> Subscription {
2548 self.control().subscription()
2549 }
2550
2551 pub fn update(&mut self, subscription: Subscription) -> Result<()> {
2557 match &mut self.inner {
2558 SubscriberKind::Plain(plain) => {
2559 let mut state = plain.subscription.write().map_err(|_| Error::Closed)?;
2560 *state = subscription;
2561 }
2562 SubscriberKind::Spliced(spliced) => spliced.update(subscription),
2563 }
2564 Ok(())
2565 }
2566
2567 pub fn latest(&self) -> Option<u64> {
2569 match &self.inner {
2570 SubscriberKind::Plain(plain) => plain.state.read().max_sequence,
2571 SubscriberKind::Spliced(spliced) => spliced.latest(),
2572 }
2573 }
2574}
2575
2576pub struct Request {
2588 name: Arc<str>,
2589 broadcast: Arc<broadcast::Info>,
2591 state: kio::Producer<TrackState>,
2592
2593 prev_subscription: Option<Subscription>,
2595
2596 alive: Arc<Alive>,
2599
2600 _dynamic: Dynamic,
2605
2606 stats: stats::Scope,
2609}
2610
2611impl Request {
2612 pub(crate) fn new(broadcast: Arc<broadcast::Info>, name: impl Into<Arc<str>>) -> Self {
2613 let name = name.into();
2614 let state = TrackState::spawn(broadcast.clone());
2615 let alive = Alive::new(name.clone(), state.clone());
2616 let dynamic = Dynamic::new(name.clone(), state.clone(), alive.clone());
2617 Self {
2618 name,
2619 broadcast,
2620 state,
2621 prev_subscription: None,
2622 alive,
2623 _dynamic: dynamic,
2624 stats: stats::Scope::default(),
2625 }
2626 }
2627
2628 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
2631 self.stats = scope;
2632 self
2633 }
2634
2635 pub fn name(&self) -> &str {
2637 &self.name
2638 }
2639
2640 pub fn consume(&self) -> Consumer {
2642 Consumer::plain(self.name.clone(), self.state.consume())
2643 }
2644
2645 pub fn dynamic(&self) -> Dynamic {
2649 Dynamic::new(self.name.clone(), self.state.clone(), self.alive.clone())
2650 }
2651
2652 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
2655 self.state.poll_unused(waiter).map(|_| ())
2656 }
2657
2658 pub fn accept(self, info: impl Into<Option<Info>>) -> Producer {
2665 if let Ok(mut state) = self.state.write() {
2668 state.install(info.into().unwrap_or_default());
2669 }
2670 self.alive.publish(Some(&self.stats));
2673 Producer {
2674 name: self.name,
2675 broadcast: self.broadcast,
2676 state: self.state,
2677 prev_subscription: None,
2678 alive: self.alive,
2679 stats: self.stats,
2680 }
2681 }
2682
2683 pub fn reject(self, err: Error) {
2685 if let Ok(mut state) = self.state.write() {
2686 state.abort = Some(err);
2687 }
2688 }
2689
2690 pub fn subscription(&self) -> Option<Subscription> {
2693 let state = self.state.read();
2694 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
2695 drop(state);
2696 snapshot_subscription(&subs, bound)
2697 }
2698
2699 pub async fn subscription_changed(&mut self) -> Option<Subscription> {
2702 kio::wait(|waiter| self.poll_subscription_changed(waiter)).await
2703 }
2704
2705 pub fn poll_subscription_changed(&mut self, waiter: &kio::Waiter) -> Poll<Option<Subscription>> {
2707 let state = self.state.read();
2708 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
2709 drop(state);
2710
2711 let prev = &self.prev_subscription;
2712 let mut combined = None;
2713 let mut guard = ready!(subs.poll(waiter, |subs| {
2714 let next = combined_subscription(subs, bound, waiter);
2715 if &next == prev {
2716 Poll::Pending
2717 } else {
2718 combined = next;
2719 Poll::Ready(())
2720 }
2721 }));
2722 guard.retain(|sub| !sub.is_closed());
2724 drop(guard);
2725 self.prev_subscription = combined.clone();
2726 Poll::Ready(combined)
2727 }
2728
2729 pub(super) fn weak(&self) -> TrackWeak {
2730 TrackWeak {
2731 name: self.name.clone(),
2732 state: self.state.weak(),
2733 }
2734 }
2735}
2736
2737#[cfg(test)]
2738use futures::FutureExt;
2739
2740#[cfg(test)]
2741#[allow(missing_docs)] impl Subscriber {
2743 pub fn assert_group(&mut self) -> group::Consumer {
2744 self.recv_group()
2745 .now_or_never()
2746 .expect("group would have blocked")
2747 .expect("would have errored")
2748 .expect("track was closed")
2749 }
2750
2751 pub fn assert_no_group(&mut self) {
2752 assert!(
2753 self.recv_group().now_or_never().is_none(),
2754 "recv_group would not have blocked"
2755 );
2756 }
2757
2758 pub fn assert_not_closed(&mut self) {
2759 assert!(self.finished().now_or_never().is_none(), "should not be closed");
2760 }
2761
2762 pub fn assert_closed(&mut self) {
2763 assert!(self.finished().now_or_never().is_some(), "should be closed");
2764 }
2765
2766 pub fn assert_error(&mut self) {
2768 assert!(
2769 self.finished().now_or_never().expect("should not block").is_err(),
2770 "should be error"
2771 );
2772 }
2773
2774 pub fn assert_is_clone(&self, other: &Self) {
2775 assert!(self.is_clone(other), "should be clone");
2776 }
2777
2778 pub fn assert_not_clone(&self, other: &Self) {
2779 assert!(!self.is_clone(other), "should not be clone");
2780 }
2781}
2782
2783#[cfg(test)]
2784mod test {
2785 use super::*;
2786 use crate::model::test_tracing::count_drop_warnings;
2787
2788 fn track_producer(name: impl Into<Arc<str>>, info: impl Into<Option<Info>>) -> Producer {
2791 Producer::new(Arc::new(broadcast::Info::default()), name, info)
2792 }
2793
2794 fn live_groups(state: &TrackState) -> usize {
2796 state.lookup.len()
2797 }
2798
2799 fn first_live_sequence(state: &TrackState) -> u64 {
2801 state
2802 .arrival
2803 .iter()
2804 .find(|(sequence, stamp)| state.lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp))
2805 .map(|(sequence, _)| *sequence)
2806 .unwrap()
2807 }
2808
2809 fn recv_datagram(dg: &mut Subscriber) -> Datagram {
2811 dg.recv_datagram()
2812 .now_or_never()
2813 .expect("datagram would have blocked")
2814 .expect("would have errored")
2815 .expect("track was closed")
2816 }
2817
2818 #[tokio::test]
2819 async fn append_datagram_shares_group_sequence() {
2820 let mut producer = track_producer("test", None);
2821 let ts = Timestamp::from_millis(10).unwrap();
2822
2823 assert_eq!(producer.append_group().unwrap().sequence, 0);
2825 assert_eq!(producer.append_datagram(ts, &b"a"[..]).unwrap(), 1);
2826 assert_eq!(producer.append_group().unwrap().sequence, 2);
2827 assert_eq!(producer.append_datagram(ts, &b"b"[..]).unwrap(), 3);
2828 assert_eq!(producer.latest(), Some(3));
2829 }
2830
2831 #[tokio::test]
2832 async fn append_datagram_roundtrip() {
2833 let mut producer = track_producer("test", None);
2834 let mut dg = producer.subscribe(None);
2835
2836 let ts = Timestamp::from_millis(42).unwrap();
2837 let seq = producer.append_datagram(ts, &b"hello"[..]).unwrap();
2838
2839 let got = recv_datagram(&mut dg);
2840 assert_eq!(got.sequence, seq);
2841 assert_eq!(got.timestamp, ts);
2842 assert_eq!(&got.payload[..], b"hello");
2843 }
2844
2845 #[tokio::test]
2846 async fn write_datagram_preserves_sequence() {
2847 let mut producer = track_producer("test", None);
2848 let mut dg = producer.subscribe(None);
2849
2850 let ts = Timestamp::from_millis(5).unwrap();
2851 producer
2853 .write_datagram(Datagram {
2854 sequence: 100,
2855 timestamp: ts,
2856 payload: bytes::Bytes::from_static(b"x"),
2857 })
2858 .unwrap();
2859
2860 assert_eq!(recv_datagram(&mut dg).sequence, 100);
2861 assert_eq!(producer.append_group().unwrap().sequence, 101);
2863 }
2864
2865 #[tokio::test]
2866 async fn recv_datagram_advances_ordered_group_cursor() {
2867 let mut producer = track_producer("test", None);
2868 let mut subscriber = producer.subscribe(None);
2869 let ts = Timestamp::from_millis(5).unwrap();
2870
2871 producer
2872 .write_datagram(Datagram {
2873 sequence: 5,
2874 timestamp: ts,
2875 payload: bytes::Bytes::from_static(b"x"),
2876 })
2877 .unwrap();
2878 assert_eq!(recv_datagram(&mut subscriber).sequence, 5);
2879
2880 producer.create_group(group::Info { sequence: 3 }).unwrap();
2881 producer.create_group(group::Info { sequence: 6 }).unwrap();
2882
2883 let group = subscriber
2884 .next_group()
2885 .now_or_never()
2886 .expect("group would have blocked")
2887 .expect("would have errored")
2888 .expect("track was closed");
2889 assert_eq!(group.sequence, 6);
2890 }
2891
2892 #[tokio::test]
2893 async fn datagram_normalized_to_track_timescale() {
2894 let info = Info::default().with_timescale(Timescale::MICRO);
2895 let mut producer = track_producer("test", info);
2896 let mut dg = producer.subscribe(None);
2897
2898 producer
2900 .append_datagram(Timestamp::from_millis(2).unwrap(), &b"z"[..])
2901 .unwrap();
2902 let got = recv_datagram(&mut dg);
2903 assert_eq!(got.timestamp.scale(), Timescale::MICRO);
2904 assert_eq!(got.timestamp.value(), 2_000);
2905 }
2906
2907 #[tokio::test]
2908 async fn datagram_rejects_oversized() {
2909 let mut producer = track_producer("test", None);
2910 let big = bytes::Bytes::from(vec![0u8; crate::model::datagram::MAX_DATAGRAM_PAYLOAD + 1]);
2911 let ts = Timestamp::from_millis(0).unwrap();
2912 assert!(matches!(
2913 producer.append_datagram(ts, big.clone()),
2914 Err(Error::FrameTooLarge)
2915 ));
2916 assert!(matches!(
2917 producer.write_datagram(Datagram {
2918 sequence: 0,
2919 timestamp: ts,
2920 payload: big,
2921 }),
2922 Err(Error::FrameTooLarge)
2923 ));
2924 }
2925
2926 #[tokio::test]
2927 async fn datagram_fanout_to_subscribers() {
2928 let mut producer = track_producer("test", None);
2929 let mut a = producer.subscribe(None);
2931 let mut b = producer.subscribe(None);
2932 let ts = Timestamp::from_millis(1).unwrap();
2933
2934 producer.append_datagram(ts, &b"first"[..]).unwrap();
2935 producer.append_datagram(ts, &b"second"[..]).unwrap();
2936
2937 assert_eq!(&recv_datagram(&mut a).payload[..], b"first");
2939 assert_eq!(&recv_datagram(&mut a).payload[..], b"second");
2940 assert_eq!(&recv_datagram(&mut b).payload[..], b"first");
2941 assert_eq!(&recv_datagram(&mut b).payload[..], b"second");
2942 }
2943
2944 #[tokio::test]
2945 async fn datagram_evicts_stale() {
2946 tokio::time::pause();
2947
2948 let mut producer = track_producer("test", None);
2949 let mut dg = producer.subscribe(None);
2950 let ts = Timestamp::from_millis(0).unwrap();
2951
2952 producer.append_datagram(ts, &b"old"[..]).unwrap(); tokio::time::advance(MAX_DATAGRAM_AGE + Duration::from_millis(10)).await;
2956 producer.append_datagram(ts, &b"new"[..]).unwrap(); let got = recv_datagram(&mut dg);
2960 assert_eq!(got.sequence, 1);
2961 assert_eq!(&got.payload[..], b"new");
2962 }
2963
2964 #[tokio::test]
2965 async fn datagram_recv_pends_until_written() {
2966 let mut producer = track_producer("test", None);
2967 let mut dg = producer.subscribe(None);
2968
2969 assert!(
2970 dg.recv_datagram().now_or_never().is_none(),
2971 "should block with no datagrams"
2972 );
2973
2974 producer
2975 .append_datagram(Timestamp::from_millis(0).unwrap(), &b"go"[..])
2976 .unwrap();
2977 assert_eq!(&recv_datagram(&mut dg).payload[..], b"go");
2978 }
2979
2980 #[tokio::test]
2984 async fn datagram_wire_roundtrip_between_tracks() {
2985 use crate::coding::{Decode, Encode};
2986 use crate::lite;
2987
2988 let version = lite::Version::Lite05;
2989
2990 let mut origin = track_producer("test", None);
2992 let mut origin_dg = origin.subscribe(None);
2993 let ts = Timestamp::from_millis(7).unwrap();
2994 let seq = origin.append_datagram(ts, &b"payload"[..]).unwrap();
2995
2996 let d = recv_datagram(&mut origin_dg);
2997 let body = lite::Datagram {
2998 subscribe: 5,
2999 sequence: d.sequence,
3000 timestamp: d.timestamp.value(),
3001 payload: d.payload.clone(),
3002 }
3003 .encode_bytes(version)
3004 .unwrap();
3005
3006 let mut slice = &body[..];
3008 let wire = lite::Datagram::decode(&mut slice, version).unwrap();
3009 let mut downstream = track_producer("test", None);
3010 let mut downstream_dg = downstream.subscribe(None);
3011 downstream
3012 .write_datagram(Datagram {
3013 sequence: wire.sequence,
3014 timestamp: Timestamp::new(wire.timestamp, Timescale::MILLI).unwrap(),
3015 payload: wire.payload,
3016 })
3017 .unwrap();
3018
3019 let got = recv_datagram(&mut downstream_dg);
3020 assert_eq!(got.sequence, seq);
3021 assert_eq!(got.timestamp, ts);
3022 assert_eq!(&got.payload[..], b"payload");
3023 }
3024
3025 #[tokio::test]
3026 async fn evict_expired_groups() {
3027 tokio::time::pause();
3028
3029 let mut producer = track_producer("test", None);
3030
3031 producer.append_group().unwrap(); producer.append_group().unwrap(); producer.append_group().unwrap(); {
3037 let state = producer.state.read();
3038 assert_eq!(live_groups(&state), 3);
3039 assert_eq!(state.offset, 0);
3040 }
3041
3042 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3044
3045 producer.append_group().unwrap(); {
3052 let state = producer.state.read();
3053 assert_eq!(live_groups(&state), 1);
3054 assert_eq!(first_live_sequence(&state), 3);
3055 assert_eq!(state.offset, 3);
3056 assert!(!state.lookup.contains_key(&0));
3057 assert!(!state.lookup.contains_key(&1));
3058 assert!(!state.lookup.contains_key(&2));
3059 assert!(state.lookup.contains_key(&3));
3060 }
3061 }
3062
3063 #[tokio::test]
3067 async fn aging_out_a_finished_group_keeps_the_clean_end() {
3068 tokio::time::pause();
3069
3070 let mut producer = track_producer("test", None);
3071 let mut group = producer.create_group(group::Info { sequence: 0 }).unwrap();
3072 let mut consumer = group.consume();
3073
3074 group
3075 .write_frame(Timestamp::from_millis(0).unwrap(), b"hello".as_slice())
3076 .unwrap();
3077 assert_eq!(consumer.next_frame().await.unwrap().unwrap().size, 5);
3078
3079 tokio::time::advance(DEFAULT_LATENCY_MAX * 12).await;
3081 group.finish().unwrap();
3082 let _next = producer.create_group(group::Info { sequence: 1 }).unwrap();
3083
3084 assert!(consumer.next_frame().await.unwrap().is_none());
3085 }
3086
3087 #[tokio::test]
3088 async fn evict_keeps_max_sequence() {
3089 tokio::time::pause();
3090
3091 let mut producer = track_producer("test", None);
3092 producer.append_group().unwrap(); tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3096
3097 producer.append_group().unwrap(); {
3101 let state = producer.state.read();
3102 assert_eq!(live_groups(&state), 1);
3103 assert_eq!(first_live_sequence(&state), 1);
3104 assert_eq!(state.offset, 1);
3105 }
3106 }
3107
3108 #[tokio::test]
3109 async fn no_eviction_when_fresh() {
3110 tokio::time::pause();
3111
3112 let mut producer = track_producer("test", None);
3113 producer.append_group().unwrap(); producer.append_group().unwrap(); producer.append_group().unwrap(); {
3118 let state = producer.state.read();
3119 assert_eq!(live_groups(&state), 3);
3120 assert_eq!(state.offset, 0);
3121 }
3122 }
3123
3124 #[tokio::test]
3125 async fn consumer_skips_evicted_groups() {
3126 tokio::time::pause();
3127
3128 let mut producer = track_producer("test", None);
3129 producer.append_group().unwrap(); let mut consumer = producer.subscribe(None);
3132
3133 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3134 producer.append_group().unwrap(); let group = consumer.assert_group();
3138 assert_eq!(group.sequence, 1);
3139 }
3140
3141 #[tokio::test]
3142 async fn cache_age_controls_eviction() {
3143 tokio::time::pause();
3144
3145 let mut producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(1)));
3147 producer.append_group().unwrap(); tokio::time::advance(Duration::from_secs(2)).await;
3151 producer.append_group().unwrap(); let state = producer.state.read();
3155 assert_eq!(live_groups(&state), 1);
3156 assert_eq!(first_live_sequence(&state), 1);
3157 }
3158
3159 #[test]
3160 fn latency_max_clamped_to_cache() {
3161 let producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(2)));
3162
3163 let mut subscriber = producer.subscribe(Subscription::default().with_latency_max(Duration::from_secs(10)));
3167 assert_eq!(subscriber.subscription().latency_max, Duration::from_secs(10));
3168 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3169
3170 subscriber
3172 .update(Subscription::default().with_latency_max(Duration::from_millis(500)))
3173 .unwrap();
3174 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_millis(500));
3175
3176 subscriber
3177 .update(Subscription::default().with_latency_max(Duration::ZERO))
3178 .unwrap();
3179 assert_eq!(producer.subscription().unwrap().latency_max, Duration::ZERO);
3180 }
3181
3182 fn track_producer_capped(name: impl Into<Arc<str>>, info: Info, cap: Duration) -> Producer {
3185 let origin = crate::origin::Info::default().with_cache_duration(cap);
3186 Producer::new(Arc::new(broadcast::Info { origin }), name, info)
3187 }
3188
3189 #[test]
3190 fn origin_cache_duration_clamps_latency_max() {
3191 let capped = track_producer_capped(
3194 "test",
3195 Info::default().with_latency_max(Duration::from_secs(60)),
3196 Duration::from_secs(1),
3197 );
3198 assert_eq!(capped.state.read().latency_bound(), Some(Duration::from_secs(1)));
3199
3200 let under = track_producer_capped(
3201 "test",
3202 Info::default().with_latency_max(Duration::from_millis(500)),
3203 Duration::from_secs(1),
3204 );
3205 assert_eq!(under.state.read().latency_bound(), Some(Duration::from_millis(500)));
3206 }
3207
3208 #[tokio::test]
3209 async fn origin_cache_duration_caps_eviction() {
3210 tokio::time::pause();
3211
3212 let mut producer = track_producer_capped(
3214 "test",
3215 Info::default().with_latency_max(Duration::from_secs(60)),
3216 Duration::from_secs(1),
3217 );
3218 producer.append_group().unwrap(); tokio::time::advance(Duration::from_secs(2)).await;
3222 producer.append_group().unwrap(); let state = producer.state.read();
3226 assert_eq!(live_groups(&state), 1);
3227 assert_eq!(first_live_sequence(&state), 1);
3228 }
3229
3230 #[test]
3231 fn latency_max_clamped_via_every_update_path() {
3232 let producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(2)));
3233 let over = Subscription::default().with_latency_max(Duration::from_secs(10));
3234
3235 let mut subscriber = producer.subscribe(over.clone());
3238 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3239
3240 subscriber.control().update(over.clone()).unwrap();
3241 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3242
3243 subscriber.update(over).unwrap();
3244 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3245 }
3246
3247 #[test]
3248 fn latency_max_aggregate_clamps_the_max_across_subscribers() {
3249 let producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(2)));
3250
3251 let _a = producer.subscribe(Subscription::default().with_latency_max(Duration::from_millis(500)));
3254 let _b = producer.subscribe(Subscription::default().with_latency_max(Duration::from_secs(10)));
3255
3256 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3257 }
3258
3259 #[test]
3260 fn subscriber_control_updates_while_read_future_is_pending() {
3261 let producer = track_producer("test", None);
3262 let mut subscriber = producer.subscribe(None);
3263 let control = subscriber.control();
3264
3265 let mut recv = Box::pin(subscriber.recv_group());
3266 assert!(recv.as_mut().now_or_never().is_none());
3267
3268 control
3269 .update(Subscription::default().with_priority(7).with_ordered(false))
3270 .unwrap();
3271
3272 let aggregate = producer.subscription().expect("expected an active subscription");
3273 assert_eq!(aggregate.priority, 7);
3274 assert!(!aggregate.ordered);
3275 }
3276
3277 #[test]
3278 fn dropped_subscriber_leaves_no_ghost_in_aggregate() {
3279 let mut producer = track_producer("test", None);
3284 let a = producer.subscribe(Subscription::default().with_priority(5));
3285
3286 let waiter = kio::Waiter::noop();
3288 assert!(
3289 matches!(producer.poll_subscription_changed(&waiter), Poll::Ready(Ok(Some(_)))),
3290 "one live subscriber should aggregate to Some",
3291 );
3292
3293 drop(a);
3295
3296 assert!(
3298 matches!(producer.poll_subscription_changed(&waiter), Poll::Ready(Ok(None))),
3299 "a dropped subscriber must not linger in the aggregate",
3300 );
3301
3302 assert!(
3304 producer.subscription().is_none(),
3305 "snapshot must exclude a dropped subscriber",
3306 );
3307 }
3308
3309 #[test]
3310 fn dropped_subscriber_wakes_the_aggregate() {
3311 use std::sync::atomic::{AtomicBool, Ordering};
3318
3319 let mut producer = track_producer("test", None);
3320 let a = producer.subscribe(Subscription::default().with_priority(5));
3321
3322 let woken = Arc::new(AtomicBool::new(false));
3323 let waiter = kio::Waiter::new(futures::task::waker(Arc::new(FlagWake(woken.clone()))));
3324
3325 assert!(matches!(
3327 producer.poll_subscription_changed(&waiter),
3328 Poll::Ready(Ok(Some(_)))
3329 ));
3330 assert!(
3331 producer.poll_subscription_changed(&waiter).is_pending(),
3332 "the aggregate is unchanged, so this poll must park",
3333 );
3334 assert!(!woken.load(Ordering::SeqCst), "nothing happened yet");
3335
3336 drop(a);
3337 assert!(
3338 woken.load(Ordering::SeqCst),
3339 "the last subscriber leaving must wake the aggregate watcher",
3340 );
3341 }
3342
3343 struct FlagWake(Arc<std::sync::atomic::AtomicBool>);
3345
3346 impl futures::task::ArcWake for FlagWake {
3347 fn wake_by_ref(arc_self: &Arc<Self>) {
3348 arc_self.0.store(true, std::sync::atomic::Ordering::SeqCst);
3349 }
3350 }
3351
3352 #[tokio::test]
3353 async fn out_of_order_max_sequence_at_front() {
3354 tokio::time::pause();
3355
3356 let mut producer = track_producer("test", None);
3357
3358 producer.create_group(group::Info { sequence: 5 }).unwrap();
3360 producer.create_group(group::Info { sequence: 3 }).unwrap();
3361 producer.create_group(group::Info { sequence: 4 }).unwrap();
3362
3363 {
3365 let state = producer.state.read();
3366 assert_eq!(state.max_sequence, Some(5));
3367 }
3368
3369 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3371
3372 producer.append_group().unwrap(); {
3378 let state = producer.state.read();
3379 assert_eq!(live_groups(&state), 1);
3380 assert_eq!(first_live_sequence(&state), 6);
3381 assert!(!state.lookup.contains_key(&3));
3382 assert!(!state.lookup.contains_key(&4));
3383 assert!(!state.lookup.contains_key(&5));
3384 assert!(state.lookup.contains_key(&6));
3385 }
3386 }
3387
3388 #[tokio::test]
3389 async fn max_sequence_at_front_blocks_trim() {
3390 tokio::time::pause();
3391
3392 let mut producer = track_producer("test", None);
3393
3394 producer.create_group(group::Info { sequence: 5 }).unwrap();
3396
3397 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3398
3399 producer.create_group(group::Info { sequence: 3 }).unwrap();
3401
3402 {
3405 let state = producer.state.read();
3406 assert_eq!(live_groups(&state), 2);
3407 assert_eq!(state.offset, 0);
3408 }
3409
3410 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3412
3413 producer.create_group(group::Info { sequence: 2 }).unwrap();
3415
3416 {
3421 let state = producer.state.read();
3422 assert_eq!(live_groups(&state), 2);
3423 assert_eq!(state.offset, 0);
3424 assert!(state.lookup.contains_key(&5));
3425 assert!(!state.lookup.contains_key(&3));
3426 assert!(state.lookup.contains_key(&2));
3427 }
3428
3429 let mut consumer = producer.subscribe(None);
3431 let group = consumer.assert_group();
3432 assert_eq!(group.sequence, 5);
3434 }
3435
3436 #[tokio::test]
3437 async fn abort_clears_cached_groups() {
3438 let mut producer = track_producer("test", None);
3439 producer.append_group().unwrap();
3440 producer.append_group().unwrap();
3441
3442 let mut consumer = producer.subscribe(None);
3444 assert_eq!(live_groups(&producer.state.read()), 2);
3445
3446 producer.clone().abort(Error::Cancel).unwrap();
3447
3448 {
3449 let state = producer.state.read();
3450 assert!(state.lookup.is_empty(), "cached groups should be dropped on abort");
3451 assert!(state.arrival.is_empty());
3452 assert!(state.evict.is_empty());
3453 }
3454
3455 let result = consumer.recv_group().now_or_never().expect("should not block");
3457 assert!(matches!(result, Err(Error::Cancel)));
3458 }
3459
3460 #[tokio::test]
3461 async fn drop_unfinished_clears_cached_groups() {
3462 let producer = track_producer("test", None);
3463 let mut writer = producer.clone();
3464 writer.append_group().unwrap();
3465
3466 let mut consumer = producer.subscribe(None);
3468 assert_eq!(live_groups(&producer.state.read()), 1);
3469
3470 drop(writer);
3472 drop(producer);
3473
3474 let result = consumer.recv_group().now_or_never().expect("should not block");
3475 assert!(matches!(result, Err(Error::Dropped)));
3476 }
3477
3478 #[tokio::test]
3479 async fn drop_after_abort_does_not_warn() {
3480 let warns = count_drop_warnings("track::Producer dropped without finish", || {
3483 let producer = track_producer("test", None);
3484 let keep = producer.clone();
3485 let mut writer = producer.clone();
3486 let mut group = writer.append_group().unwrap();
3487 group.finish().unwrap();
3488 let _consumer = producer.subscribe(None);
3489 writer.abort(Error::Cancel).unwrap();
3490 drop(keep);
3491 });
3492 assert_eq!(warns, 0, "abort-then-drop must not emit unfinished-producer WARN");
3493 }
3494
3495 #[tokio::test]
3496 async fn drop_unfinished_warns() {
3497 let warns = count_drop_warnings("track::Producer dropped without finish", || {
3498 let producer = track_producer("test", None);
3499 let mut writer = producer.clone();
3500 writer.append_group().unwrap();
3501 let _consumer = producer.subscribe(None);
3502 drop(writer);
3503 drop(producer);
3504 });
3505 assert!(warns >= 1, "unfinished drop must emit unfinished-producer WARN");
3506 }
3507
3508 #[tokio::test]
3509 async fn drop_finished_keeps_cached_groups() {
3510 let mut producer = track_producer("test", None);
3511 producer.append_group().unwrap();
3512 producer.finish().unwrap();
3513
3514 let mut consumer = producer.subscribe(None);
3515 drop(producer);
3516
3517 assert_eq!(consumer.assert_group().sequence, 0);
3519 let done = consumer.recv_group().now_or_never().expect("should not block").unwrap();
3520 assert!(done.is_none(), "consumer should drain then see clean finish");
3521 }
3522
3523 #[test]
3524 fn append_finish_cannot_be_rewritten() {
3525 let mut producer = track_producer("test", None);
3526
3527 assert!(producer.finish().is_ok());
3529 assert!(producer.finish().is_err());
3530 assert!(producer.append_group().is_err());
3531 }
3532
3533 #[test]
3534 fn finish_after_groups() {
3535 let mut producer = track_producer("test", None);
3536
3537 producer.append_group().unwrap();
3538 assert!(producer.finish().is_ok());
3539 assert!(producer.finish().is_err());
3540 assert!(producer.append_group().is_err());
3541 }
3542
3543 #[test]
3544 fn finish_at_rejects_a_boundary_at_or_below_the_live_edge() {
3545 let mut producer = track_producer("test", None);
3546 producer.create_group(group::Info { sequence: 5 }).unwrap();
3547
3548 assert!(producer.finish_at(4).is_err());
3551 assert!(producer.finish_at(5).is_err());
3552 assert!(producer.finish_at(6).is_ok());
3553
3554 {
3555 let state = producer.state.read();
3556 assert_eq!(state.final_sequence, Some(6));
3557 }
3558
3559 assert!(producer.finish_at(6).is_err());
3561 assert!(producer.create_group(group::Info { sequence: 4 }).is_ok());
3562 assert!(producer.create_group(group::Info { sequence: 6 }).is_err());
3563 }
3564
3565 #[test]
3566 fn final_sequence_reports_the_declared_boundary() {
3567 let mut producer = track_producer("test", None);
3568 assert_eq!(producer.final_sequence(), None);
3569
3570 producer.create_group(group::Info { sequence: 5 }).unwrap();
3571 assert_eq!(producer.final_sequence(), None, "a group does not declare a boundary");
3572
3573 producer.finish_at(9).unwrap();
3574 assert_eq!(producer.final_sequence(), Some(9));
3575
3576 assert!(producer.finish().is_err());
3578 }
3579
3580 #[test]
3581 fn final_sequence_reports_the_live_edge_after_finish() {
3582 let mut producer = track_producer("test", None);
3583 producer.create_group(group::Info { sequence: 5 }).unwrap();
3584 producer.finish().unwrap();
3585 assert_eq!(producer.final_sequence(), Some(6));
3586 }
3587
3588 #[tokio::test]
3589 async fn finish_at_declares_a_future_boundary() {
3590 let mut producer = track_producer("test", None);
3591 producer.create_group(group::Info { sequence: 5 }).unwrap();
3592
3593 producer.finish_at(7).unwrap();
3595
3596 let mut consumer = producer.subscribe(None);
3597 assert_eq!(consumer.assert_group().sequence, 5);
3598
3599 let boundary = consumer
3602 .finished()
3603 .now_or_never()
3604 .expect("boundary is known immediately")
3605 .expect("would have errored");
3606 assert_eq!(boundary, 7);
3607 assert!(
3608 consumer.recv_group().now_or_never().is_none(),
3609 "should wait for the outstanding group"
3610 );
3611
3612 producer.create_group(group::Info { sequence: 6 }).unwrap();
3614 assert_eq!(consumer.assert_group().sequence, 6);
3615 let done = consumer
3616 .recv_group()
3617 .now_or_never()
3618 .expect("should not block")
3619 .expect("would have errored");
3620 assert!(done.is_none(), "track completes once the boundary is reached");
3621 }
3622
3623 #[tokio::test]
3624 async fn recv_group_finishes_without_waiting_for_gaps() {
3625 let mut producer = track_producer("test", None);
3626 producer.create_group(group::Info { sequence: 1 }).unwrap();
3627 producer.finish().unwrap();
3628
3629 let mut consumer = producer.subscribe(None);
3630 assert_eq!(consumer.assert_group().sequence, 1);
3631
3632 let done = consumer
3633 .recv_group()
3634 .now_or_never()
3635 .expect("should not block")
3636 .expect("would have errored");
3637 assert!(done.is_none(), "track should finish without waiting for gaps");
3638 }
3639
3640 #[tokio::test]
3641 async fn next_group_skips_late_arrivals() {
3642 let mut producer = track_producer("test", None);
3643 let mut consumer = producer.subscribe(None);
3644
3645 producer.create_group(group::Info { sequence: 5 }).unwrap();
3647 let group = consumer
3648 .next_group()
3649 .now_or_never()
3650 .expect("should not block")
3651 .expect("would have errored")
3652 .expect("track should not be closed");
3653 assert_eq!(group.sequence, 5);
3654
3655 producer.create_group(group::Info { sequence: 3 }).unwrap();
3657 producer.create_group(group::Info { sequence: 4 }).unwrap();
3659 producer.create_group(group::Info { sequence: 7 }).unwrap();
3661
3662 let group = consumer
3663 .next_group()
3664 .now_or_never()
3665 .expect("should not block")
3666 .expect("would have errored")
3667 .expect("track should not be closed");
3668 assert_eq!(group.sequence, 7);
3669
3670 assert!(
3672 consumer.next_group().now_or_never().is_none(),
3673 "should block waiting for a higher sequence"
3674 );
3675 }
3676
3677 #[tokio::test]
3678 async fn next_group_returns_arrivals_in_order() {
3679 let mut producer = track_producer("test", None);
3680 let mut consumer = producer.subscribe(None);
3681
3682 producer.create_group(group::Info { sequence: 3 }).unwrap();
3684 producer.create_group(group::Info { sequence: 5 }).unwrap();
3685
3686 let group = consumer
3687 .next_group()
3688 .now_or_never()
3689 .expect("should not block")
3690 .expect("would have errored")
3691 .expect("track should not be closed");
3692 assert_eq!(group.sequence, 3);
3693
3694 let group = consumer
3695 .next_group()
3696 .now_or_never()
3697 .expect("should not block")
3698 .expect("would have errored")
3699 .expect("track should not be closed");
3700 assert_eq!(group.sequence, 5);
3701 }
3702
3703 #[tokio::test]
3704 async fn next_group_and_recv_group_use_independent_cursors() {
3705 let mut producer = track_producer("test", None);
3706 let mut consumer = producer.subscribe(None);
3707
3708 producer.create_group(group::Info { sequence: 5 }).unwrap();
3710 producer.create_group(group::Info { sequence: 3 }).unwrap();
3711
3712 let group = consumer
3715 .next_group()
3716 .now_or_never()
3717 .expect("should not block")
3718 .expect("would have errored")
3719 .expect("track should not be closed");
3720 assert_eq!(group.sequence, 3);
3721
3722 assert_eq!(consumer.assert_group().sequence, 5);
3725 }
3726
3727 #[tokio::test]
3728 async fn end_at_caps_next_group() {
3729 let mut producer = track_producer("test", None);
3730 let mut consumer = producer.subscribe(None);
3731
3732 for s in 0..6 {
3733 producer.create_group(group::Info { sequence: s }).unwrap();
3734 }
3735
3736 consumer.end_at(2);
3737
3738 assert_eq!(
3740 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3741 0
3742 );
3743 assert_eq!(
3744 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3745 1
3746 );
3747 assert_eq!(
3748 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3749 2
3750 );
3751
3752 assert!(
3754 consumer.next_group().now_or_never().is_none(),
3755 "capped consumer must block instead of returning out-of-range groups"
3756 );
3757 }
3758
3759 #[tokio::test]
3760 async fn end_at_release_drains_cached_groups() {
3761 let mut producer = track_producer("test", None);
3762 let mut consumer = producer.subscribe(None);
3763
3764 for s in 0..6 {
3765 producer.create_group(group::Info { sequence: s }).unwrap();
3766 }
3767
3768 consumer.end_at(1);
3769 assert_eq!(
3770 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3771 0
3772 );
3773 assert_eq!(
3774 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3775 1
3776 );
3777 assert!(consumer.next_group().now_or_never().is_none(), "capped at 1");
3778
3779 consumer.end_at(4);
3781 assert_eq!(
3782 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3783 2
3784 );
3785 assert_eq!(
3786 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3787 3
3788 );
3789 assert_eq!(
3790 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3791 4
3792 );
3793 assert!(consumer.next_group().now_or_never().is_none(), "capped at 4");
3794
3795 consumer.end_at(None);
3797 assert_eq!(
3798 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3799 5
3800 );
3801 assert!(consumer.next_group().now_or_never().is_none(), "no more groups");
3802 }
3803
3804 #[tokio::test]
3805 async fn end_at_lower_than_cursor_parks_consumer() {
3806 let mut producer = track_producer("test", None);
3807 let mut consumer = producer.subscribe(None);
3808
3809 for s in 0..3 {
3810 producer.create_group(group::Info { sequence: s }).unwrap();
3811 }
3812
3813 assert_eq!(
3815 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3816 0
3817 );
3818 assert_eq!(
3819 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3820 1
3821 );
3822 assert_eq!(
3823 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3824 2
3825 );
3826
3827 consumer.end_at(1);
3829 producer.create_group(group::Info { sequence: 3 }).unwrap();
3830 producer.create_group(group::Info { sequence: 4 }).unwrap();
3831 assert!(
3832 consumer.next_group().now_or_never().is_none(),
3833 "cap is below cursor; nothing returnable until cap rises"
3834 );
3835
3836 consumer.end_at(None);
3838 assert_eq!(
3839 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3840 3
3841 );
3842 assert_eq!(
3843 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3844 4
3845 );
3846 }
3847
3848 #[tokio::test]
3849 async fn end_at_toggling_around_late_arrivals() {
3850 let mut producer = track_producer("test", None);
3851 let mut consumer = producer.subscribe(None);
3852
3853 consumer.end_at(5);
3854
3855 producer.create_group(group::Info { sequence: 2 }).unwrap();
3857 producer.create_group(group::Info { sequence: 5 }).unwrap();
3858 producer.create_group(group::Info { sequence: 3 }).unwrap();
3859 producer.create_group(group::Info { sequence: 8 }).unwrap();
3861 producer.create_group(group::Info { sequence: 4 }).unwrap();
3862
3863 assert_eq!(
3865 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3866 2
3867 );
3868 assert_eq!(
3869 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3870 3
3871 );
3872 assert_eq!(
3873 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3874 4
3875 );
3876 assert_eq!(
3877 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3878 5
3879 );
3880 assert!(consumer.next_group().now_or_never().is_none());
3882
3883 consumer.end_at(10);
3885 assert_eq!(
3886 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3887 8
3888 );
3889 }
3890
3891 #[tokio::test]
3895 async fn end_at_parks_recv_group() {
3896 let mut producer = track_producer("test", None);
3897 let mut consumer = producer.subscribe(None);
3898
3899 for s in 0..3 {
3900 producer.create_group(group::Info { sequence: s }).unwrap();
3901 }
3902
3903 consumer.end_at(1);
3904 assert_eq!(
3905 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3906 0
3907 );
3908 assert_eq!(
3909 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3910 1
3911 );
3912 assert!(consumer.recv_group().now_or_never().is_none(), "capped at 1");
3913
3914 producer.finish().unwrap();
3916 assert!(
3917 consumer.recv_group().now_or_never().is_none(),
3918 "still parked after finish"
3919 );
3920
3921 consumer.end_at(None);
3922 assert_eq!(
3923 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3924 2
3925 );
3926 assert!(
3927 matches!(consumer.recv_group().now_or_never(), Some(Ok(None))),
3928 "finished once the parked group drains"
3929 );
3930 }
3931
3932 #[tokio::test]
3935 async fn recv_group_serves_arrivals_behind_the_cap() {
3936 let mut producer = track_producer("test", None);
3937 let mut consumer = producer.subscribe(None);
3938
3939 consumer.end_at(1);
3940
3941 producer.create_group(group::Info { sequence: 2 }).unwrap();
3943 producer.create_group(group::Info { sequence: 0 }).unwrap();
3944 producer.create_group(group::Info { sequence: 1 }).unwrap();
3945
3946 assert_eq!(
3947 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3948 0
3949 );
3950 assert_eq!(
3951 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3952 1
3953 );
3954 assert!(consumer.recv_group().now_or_never().is_none(), "capped at 1");
3955
3956 consumer.end_at(2);
3957 assert_eq!(
3958 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3959 2
3960 );
3961 }
3962
3963 #[tokio::test]
3966 async fn start_at_drops_parked_recv_groups() {
3967 let mut producer = track_producer("test", None);
3968 let mut consumer = producer.subscribe(None);
3969
3970 consumer.end_at(0);
3971 producer.create_group(group::Info { sequence: 1 }).unwrap();
3972 assert!(
3973 consumer.recv_group().now_or_never().is_none(),
3974 "group 1 parked at the cap"
3975 );
3976
3977 consumer.start_at(2);
3978 consumer.end_at(None);
3979 producer.create_group(group::Info { sequence: 2 }).unwrap();
3980 assert_eq!(
3981 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3982 2,
3983 "the overtaken parked group is dropped, not re-offered"
3984 );
3985 }
3986
3987 #[tokio::test]
3991 async fn evicted_parked_recv_groups_are_dropped() {
3992 let mut producer = track_producer("test", None);
3993 let mut consumer = producer.subscribe(None);
3994
3995 producer.create_group(group::Info { sequence: 0 }).unwrap();
3996 assert_eq!(
3997 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3998 0
3999 );
4000
4001 consumer.end_at(0);
4002 let straggler = producer.create_group(group::Info { sequence: 1 }).unwrap();
4003 assert!(
4004 consumer.recv_group().now_or_never().is_none(),
4005 "group 1 parked at the cap"
4006 );
4007
4008 straggler.abort(Error::Old).unwrap();
4010 producer.finish().unwrap();
4011
4012 consumer.end_at(None);
4013 assert!(
4014 matches!(consumer.recv_group().now_or_never(), Some(Ok(None))),
4015 "a dead parked group must not be delivered or hold the stream open"
4016 );
4017 }
4018
4019 #[tokio::test]
4023 async fn evicted_parked_group_wakes_the_clean_end() {
4024 use std::sync::atomic::{AtomicUsize, Ordering};
4025 use std::task::{Context, Wake};
4026
4027 struct CountWaker(AtomicUsize);
4030 impl Wake for CountWaker {
4031 fn wake(self: std::sync::Arc<Self>) {
4032 self.0.fetch_add(1, Ordering::SeqCst);
4033 }
4034 }
4035
4036 let mut producer = track_producer("test", None);
4037 let mut consumer = producer.subscribe(None);
4038
4039 producer.create_group(group::Info { sequence: 0 }).unwrap();
4040 assert_eq!(
4041 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4042 0
4043 );
4044
4045 consumer.end_at(0);
4046 let straggler = producer.create_group(group::Info { sequence: 1 }).unwrap();
4047 assert!(consumer.recv_group().now_or_never().is_none(), "parked at the cap");
4048 producer.finish().unwrap();
4049
4050 let counter = std::sync::Arc::new(CountWaker(AtomicUsize::new(0)));
4051 let waker = std::task::Waker::from(counter.clone());
4052 let mut cx = Context::from_waker(&waker);
4053 let mut fut = std::pin::pin!(consumer.recv_group());
4054 assert!(
4055 fut.as_mut().poll(&mut cx).is_pending(),
4056 "the parked group holds it open"
4057 );
4058
4059 straggler.abort(Error::Old).unwrap();
4060 assert!(counter.0.load(Ordering::SeqCst) > 0, "the eviction wakeup was lost");
4061 assert!(matches!(fut.as_mut().poll(&mut cx), Poll::Ready(Ok(None))));
4062 }
4063
4064 #[tokio::test]
4065 async fn read_frame_returns_single_frame_per_group() {
4066 let mut producer = track_producer("test", None);
4067 let mut consumer = producer.subscribe(None);
4068
4069 producer.write_frame(Timestamp::ZERO, b"hello".as_slice()).unwrap();
4070 producer.write_frame(Timestamp::ZERO, b"world".as_slice()).unwrap();
4071
4072 let frame = consumer
4073 .read_frame()
4074 .now_or_never()
4075 .expect("should not block")
4076 .expect("would have errored")
4077 .expect("track should not be closed");
4078 assert_eq!(&frame.payload[..], b"hello");
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"world");
4087 }
4088
4089 #[test]
4090 fn write_frame_rejects_an_oversized_frame_before_appending_its_group() {
4091 let mut producer = track_producer("test", None);
4092 let frame = bytes::Bytes::from(vec![0; group::MAX_GROUP_CACHE as usize + 1]);
4093
4094 assert!(matches!(
4095 producer.write_frame(Timestamp::ZERO, frame),
4096 Err(Error::FrameTooLarge)
4097 ));
4098 assert_eq!(producer.latest(), None, "the rejected frame did not publish a group");
4099 }
4100
4101 #[tokio::test]
4102 async fn read_frame_preserves_timestamp() {
4103 let mut producer = track_producer("test", None);
4104 let mut consumer = producer.subscribe(None);
4105
4106 producer
4107 .write_frame(Timestamp::from_micros(20_000).unwrap(), b"hello".as_slice())
4108 .unwrap();
4109
4110 let frame = consumer
4111 .read_frame()
4112 .now_or_never()
4113 .expect("should not block")
4114 .expect("would have errored")
4115 .expect("track should not be closed");
4116 assert_eq!(frame.timestamp.as_micros(), 20_000);
4117 assert_eq!(&frame.payload[..], b"hello");
4118 }
4119
4120 #[tokio::test]
4121 async fn read_frame_skips_stalled_group_for_newer_ready_frame() {
4122 let mut producer = track_producer("test", None);
4123 let mut consumer = producer.subscribe(None);
4124
4125 let _stalled = producer.create_group(group::Info { sequence: 3 }).unwrap();
4127 let mut g5 = producer.create_group(group::Info { sequence: 5 }).unwrap();
4129 g5.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"later"))
4130 .unwrap();
4131 g5.finish().unwrap();
4132
4133 let frame = consumer
4135 .read_frame()
4136 .now_or_never()
4137 .expect("should not block on stalled earlier group")
4138 .expect("would have errored")
4139 .expect("track should not be closed");
4140 assert_eq!(&frame.payload[..], b"later");
4141 }
4142
4143 #[tokio::test]
4144 async fn read_frame_discards_rest_of_multi_frame_group() {
4145 let mut producer = track_producer("test", None);
4146 let mut consumer = producer.subscribe(None);
4147
4148 let mut g0 = producer.create_group(group::Info { sequence: 0 }).unwrap();
4150 g0.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"one"))
4151 .unwrap();
4152 g0.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"two"))
4153 .unwrap();
4154 g0.finish().unwrap();
4155
4156 producer.write_frame(Timestamp::ZERO, b"next".as_slice()).unwrap();
4158
4159 let frame = consumer
4160 .read_frame()
4161 .now_or_never()
4162 .expect("should not block")
4163 .expect("would have errored")
4164 .expect("track should not be closed");
4165 assert_eq!(&frame.payload[..], b"one");
4166
4167 let frame = consumer
4169 .read_frame()
4170 .now_or_never()
4171 .expect("should not block")
4172 .expect("would have errored")
4173 .expect("track should not be closed");
4174 assert_eq!(&frame.payload[..], b"next");
4175 }
4176
4177 #[tokio::test]
4178 async fn read_frame_waits_for_pending_group_after_finish() {
4179 let mut producer = track_producer("test", None);
4182 let mut consumer = producer.subscribe(None);
4183
4184 let mut g0 = producer.create_group(group::Info { sequence: 0 }).unwrap();
4185 producer.finish().unwrap();
4186
4187 assert!(
4189 consumer.read_frame().now_or_never().is_none(),
4190 "read_frame must block on a pending group even after finish()"
4191 );
4192
4193 g0.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"late"))
4195 .unwrap();
4196 let frame = consumer
4197 .read_frame()
4198 .now_or_never()
4199 .expect("should not block once a frame is written")
4200 .expect("would have errored")
4201 .expect("track should not be closed");
4202 assert_eq!(&frame.payload[..], b"late");
4203 }
4204
4205 #[tokio::test]
4206 async fn read_frame_respects_start_at() {
4207 let mut producer = track_producer("test", None);
4210 let mut consumer = producer.subscribe(None);
4211 consumer.start_at(5);
4212
4213 let mut g3 = producer.create_group(group::Info { sequence: 3 }).unwrap();
4215 g3.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"skip-me"))
4216 .unwrap();
4217 g3.finish().unwrap();
4218
4219 let mut g5 = producer.create_group(group::Info { sequence: 5 }).unwrap();
4220 g5.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"keep"))
4221 .unwrap();
4222 g5.finish().unwrap();
4223
4224 let frame = consumer
4225 .read_frame()
4226 .now_or_never()
4227 .expect("should not block")
4228 .expect("would have errored")
4229 .expect("track should not be closed");
4230 assert_eq!(&frame.payload[..], b"keep");
4231 }
4232
4233 #[tokio::test]
4234 async fn read_frame_returns_none_when_finished() {
4235 let mut producer = track_producer("test", None);
4236 let mut consumer = producer.subscribe(None);
4237
4238 producer.write_frame(Timestamp::ZERO, b"only".as_slice()).unwrap();
4239 producer.finish().unwrap();
4240
4241 let frame = consumer
4242 .read_frame()
4243 .now_or_never()
4244 .expect("should not block")
4245 .expect("would have errored")
4246 .expect("track should not be closed");
4247 assert_eq!(&frame.payload[..], b"only");
4248
4249 let done = consumer
4250 .read_frame()
4251 .now_or_never()
4252 .expect("should not block")
4253 .expect("would have errored");
4254 assert!(done.is_none());
4255 }
4256
4257 #[test]
4258 fn append_group_returns_bounds_exceeded_on_sequence_overflow() {
4259 let mut producer = track_producer("test", None);
4260 {
4261 let mut state = producer.state.write().ok().unwrap();
4262 state.max_sequence = Some(u64::MAX);
4263 }
4264
4265 assert!(matches!(producer.append_group(), Err(Error::BoundsExceeded(_))));
4266 }
4267
4268 #[tokio::test]
4269 async fn fetch_cache_hit() {
4270 let mut producer = track_producer("test", None);
4271
4272 let mut group = producer.append_group().unwrap(); group
4275 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hello"))
4276 .unwrap();
4277 group.finish().unwrap();
4278
4279 let dynamic = producer.dynamic();
4282 let consumer = producer.consume();
4283 assert!(consumer.peek_group(0).is_some());
4284 let mut g = consumer.fetch_group(0, None).await.unwrap();
4285 assert_eq!(g.sequence, 0);
4286 assert_eq!(&g.read_frame().await.unwrap().unwrap().payload[..], b"hello");
4287
4288 assert!(dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending());
4290 }
4291
4292 #[tokio::test]
4293 async fn fetch_miss_signals_dynamic() {
4294 let producer = track_producer("test", None);
4295 let dynamic = producer.dynamic();
4296 let consumer = producer.consume();
4297
4298 assert!(consumer.peek_group(5).is_none());
4302 let pending = consumer.fetch_group(5, group::Fetch::default().with_priority(7));
4303 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
4304
4305 let req = dynamic
4306 .requested_group()
4307 .now_or_never()
4308 .expect("should not block")
4309 .unwrap();
4310 assert_eq!(req.sequence(), 5);
4311 assert_eq!(req.priority(), 7);
4312
4313 let mut group = req.accept(None).unwrap();
4315 group
4316 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hi"))
4317 .unwrap();
4318 group.finish().unwrap();
4319
4320 let mut g = pending.await.unwrap();
4321 assert_eq!(g.sequence, 5);
4322 assert_eq!(&g.read_frame().await.unwrap().unwrap().payload[..], b"hi");
4323 }
4324
4325 #[tokio::test]
4326 async fn fetch_miss_rejects() {
4327 let producer = track_producer("test", None);
4328 let dynamic = producer.dynamic();
4329 let consumer = producer.consume();
4330
4331 let pending = consumer.fetch_group(5, None);
4332 let req = dynamic
4333 .requested_group()
4334 .now_or_never()
4335 .expect("should not block")
4336 .unwrap();
4337
4338 req.reject(Error::Cancel);
4339 assert!(matches!(pending.await, Err(Error::Cancel)));
4340 let fetch = producer.state.read().fetch.clone();
4341 assert!(fetch.read().is_empty());
4342 }
4343
4344 #[tokio::test]
4345 async fn fetch_miss_drop_rejects() {
4346 let producer = track_producer("test", None);
4347 let dynamic = producer.dynamic();
4348 let consumer = producer.consume();
4349
4350 let pending = consumer.fetch_group(5, None);
4351 let req = dynamic
4352 .requested_group()
4353 .now_or_never()
4354 .expect("should not block")
4355 .unwrap();
4356
4357 drop(req);
4358 assert!(matches!(pending.await, Err(Error::Dropped)));
4359 }
4360
4361 #[tokio::test]
4362 async fn fetch_reject_does_not_poison_retry() {
4363 let producer = track_producer("test", None);
4364 let dynamic = producer.dynamic();
4365 let consumer = producer.consume();
4366
4367 let pending = consumer.fetch_group(5, None);
4368 let req = dynamic
4369 .requested_group()
4370 .now_or_never()
4371 .expect("should not block")
4372 .unwrap();
4373 req.reject(Error::Cancel);
4374 assert!(matches!(pending.await, Err(Error::Cancel)));
4375
4376 let retry = consumer.fetch_group(5, None);
4377 let req = dynamic
4378 .requested_group()
4379 .now_or_never()
4380 .expect("should not block")
4381 .unwrap();
4382 let mut group = req.accept(None).unwrap();
4383 group
4384 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"retry"))
4385 .unwrap();
4386 group.finish().unwrap();
4387
4388 let mut group = retry.await.unwrap();
4389 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"retry");
4390 }
4391
4392 #[tokio::test]
4393 async fn fetch_coalesces_concurrent() {
4394 let producer = track_producer("test", None);
4395 let dynamic = producer.dynamic();
4396 let consumer = producer.consume();
4397
4398 let first = consumer.fetch_group(5, group::Fetch::default().with_priority(1));
4401 let second = consumer.fetch_group(5, group::Fetch::default().with_priority(7));
4402 assert!(kio::Pollable::poll(&*first, &kio::Waiter::noop()).is_pending());
4403
4404 let req = dynamic
4405 .requested_group()
4406 .now_or_never()
4407 .expect("should not block")
4408 .unwrap();
4409 assert_eq!(req.sequence(), 5);
4410 assert_eq!(req.priority(), 7);
4411 assert!(
4412 dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending(),
4413 "the second fetch queued a duplicate request"
4414 );
4415
4416 let third = consumer.fetch_group(5, None);
4418
4419 let mut group = req.accept(None).unwrap();
4421 group
4422 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hi"))
4423 .unwrap();
4424 group.finish().unwrap();
4425
4426 assert_eq!(first.await.unwrap().sequence, 5);
4427 assert_eq!(second.await.unwrap().sequence, 5);
4428 assert_eq!(third.await.unwrap().sequence, 5);
4429 }
4430
4431 #[tokio::test]
4432 async fn fetch_coalesced_reject_fails_all() {
4433 let producer = track_producer("test", None);
4434 let dynamic = producer.dynamic();
4435 let consumer = producer.consume();
4436
4437 let first = consumer.fetch_group(5, None);
4438 let second = consumer.fetch_group(5, None);
4439 let req = dynamic
4440 .requested_group()
4441 .now_or_never()
4442 .expect("should not block")
4443 .unwrap();
4444 req.reject(Error::Cancel);
4445
4446 assert!(matches!(first.await, Err(Error::Cancel)));
4447 assert!(matches!(second.await, Err(Error::Cancel)));
4448
4449 let retry = consumer.fetch_group(5, None);
4451 assert!(kio::Pollable::poll(&*retry, &kio::Waiter::noop()).is_pending());
4452 let req = dynamic
4453 .requested_group()
4454 .now_or_never()
4455 .expect("should not block")
4456 .unwrap();
4457 assert_eq!(req.sequence(), 5);
4458 }
4459
4460 #[tokio::test]
4461 async fn fetch_queued_fails_when_handlers_leave() {
4462 let producer = track_producer("test", None);
4463 let dynamic = producer.dynamic();
4464 let consumer = producer.consume();
4465
4466 let pending = consumer.fetch_group(5, None);
4468 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
4469 drop(dynamic);
4470 assert!(matches!(pending.await, Err(Error::NotFound)));
4471
4472 let fetch = producer.state.read().fetch.clone();
4474 assert!(fetch.read().is_empty());
4475 }
4476
4477 #[tokio::test]
4478 async fn fetch_miss_no_dynamic_not_found() {
4479 let mut producer = track_producer("test", None);
4482 producer.append_group().unwrap(); let consumer = producer.consume();
4484 assert!(matches!(consumer.fetch_group(5, None).await, Err(Error::NotFound)));
4485 }
4486
4487 #[tokio::test]
4488 async fn fetch_past_final_not_found() {
4489 let mut producer = track_producer("test", None);
4490 producer.append_group().unwrap(); producer.finish().unwrap(); let dynamic = producer.dynamic();
4496 let consumer = producer.consume();
4497 assert!(matches!(consumer.fetch_group(5, None).await, Err(Error::NotFound)));
4498
4499 assert!(dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending());
4501 }
4502
4503 fn pooled_producer(capacity: u64) -> (Producer, cache::Pool) {
4505 let pool = cache::Pool::new(capacity);
4506 let broadcast = broadcast::Info {
4507 origin: crate::origin::Info::default().with_pool(pool.clone()),
4508 ..Default::default()
4509 };
4510 let producer = Producer::new(Arc::new(broadcast), "test", None);
4511 (producer, pool)
4512 }
4513
4514 fn finished_group(producer: &mut Producer, size: usize) -> u64 {
4515 let mut group = producer.append_group().unwrap();
4516 group
4517 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; size]))
4518 .unwrap();
4519 group.finish().unwrap();
4520 group.sequence
4521 }
4522
4523 #[tokio::test]
4526 async fn debt_evicts_oldest_group() {
4527 tokio::time::pause();
4528
4529 let (mut producer, pool) = pooled_producer(10_000);
4531
4532 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let consumer = producer.consume();
4537 assert!(consumer.peek_group(0).is_none(), "oldest group is evicted");
4538 assert!(consumer.peek_group(2).is_some(), "latest group survives");
4539 assert!(pool.used() <= 21_000, "usage hovers near capacity: {}", pool.used());
4542
4543 let mut subscriber = producer.subscribe(None);
4545 assert!(subscriber.assert_group().sequence > 0, "evicted group is not delivered");
4546 }
4547
4548 #[tokio::test]
4550 async fn latest_group_never_evicted() {
4551 tokio::time::pause();
4552
4553 let (mut producer, pool) = pooled_producer(100);
4555 finished_group(&mut producer, 1000); assert!(pool.used() > 100, "the latest may exceed the budget");
4557
4558 finished_group(&mut producer, 1000); finished_group(&mut producer, 1000); let consumer = producer.consume();
4563 assert!(consumer.peek_group(0).is_none());
4564 let mut group = consumer.peek_group(2).expect("latest survives");
4565 assert_eq!(group.read_frame().await.unwrap().unwrap().payload.len(), 1000);
4566 }
4567
4568 #[tokio::test]
4572 async fn fetch_refresh_survives_eviction() {
4573 tokio::time::pause();
4574
4575 let (mut producer, _pool) = pooled_producer(10_000);
4576 let consumer = producer.consume();
4577
4578 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_secs(1)).await;
4580 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_secs(1)).await;
4582 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_millis(500)).await;
4584
4585 let mut fetched = consumer.fetch_group(0, None).await.unwrap();
4587 assert_eq!(fetched.read_frame().await.unwrap().unwrap().payload.len(), 3_000);
4588 tokio::time::advance(Duration::from_millis(500)).await;
4589
4590 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_secs(1)).await;
4594 finished_group(&mut producer, 3_000); assert!(consumer.peek_group(0).is_some(), "refreshed group survives");
4597 assert!(consumer.peek_group(1).is_none(), "unread group is evicted instead");
4598 }
4599
4600 #[tokio::test]
4603 async fn eviction_aborts_readers() {
4604 tokio::time::pause();
4605
4606 let (mut producer, _pool) = pooled_producer(10_000);
4607 let mut subscriber = producer.subscribe(None);
4608
4609 finished_group(&mut producer, 10_000); let mut group0 = subscriber.assert_group();
4611
4612 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let read = group0.read_frame().await;
4616 assert!(matches!(read, Err(Error::Evicted)), "expected Evicted, got {read:?}");
4617 }
4618
4619 #[tokio::test]
4623 async fn small_writes_carry_debt() {
4624 tokio::time::pause();
4625
4626 let (mut producer, pool) = pooled_producer(22_000);
4627 let consumer = producer.consume();
4628
4629 finished_group(&mut producer, 20_000); for _ in 0..3 {
4634 finished_group(&mut producer, 1_000);
4635 }
4636 assert!(consumer.peek_group(0).is_some(), "debt smaller than the victim carries");
4637
4638 for _ in 0..20 {
4640 finished_group(&mut producer, 1_000);
4641 }
4642 assert!(
4643 consumer.peek_group(0).is_none(),
4644 "accumulated debt evicts the large group"
4645 );
4646 assert!(pool.used() <= 24_000, "usage hovers near capacity: {}", pool.used());
4649 }
4650
4651 #[tokio::test]
4655 async fn payment_capped_per_write() {
4656 tokio::time::pause();
4657
4658 let (mut producer, pool) = pooled_producer(1 << 40);
4659 for _ in 0..10 {
4660 finished_group(&mut producer, 1_000);
4661 }
4662
4663 pool.resize(100);
4665 let before = pool.used();
4666
4667 finished_group(&mut producer, 1_000);
4669
4670 let consumer = producer.consume();
4671 assert!(consumer.peek_group(0).is_none(), "the oldest groups are evicted");
4672 assert!(consumer.peek_group(1).is_none());
4673 assert!(consumer.peek_group(2).is_some(), "the backlog drains gradually");
4674 assert!(pool.used() > before - 4_000, "one write must not dump the backlog");
4675 }
4676
4677 #[tokio::test]
4681 async fn accept_preserves_write_accounting() {
4682 tokio::time::pause();
4683
4684 let pool = cache::Pool::new(12_000);
4685 let broadcast = broadcast::Info {
4686 origin: crate::origin::Info::default().with_pool(pool.clone()),
4687 ..Default::default()
4688 };
4689 let request = Request::new(Arc::new(broadcast), "test");
4690 let dynamic = request.dynamic();
4691 let consumer = request.consume();
4692
4693 let pending = consumer.fetch_group(0, None);
4695 let req = dynamic
4696 .requested_group()
4697 .now_or_never()
4698 .expect("should not block")
4699 .unwrap();
4700 let mut backfill = req.accept(None).unwrap();
4701 pending.await.unwrap();
4702 backfill
4703 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 30_000]))
4704 .unwrap();
4705
4706 let mut producer = request.accept(None);
4709 producer.append_group().unwrap().finish().unwrap();
4710 producer.append_group().unwrap().finish().unwrap();
4711
4712 assert!(
4713 producer.consume().peek_group(0).is_none(),
4714 "pre-accept backfill growth is reclaimed after accept"
4715 );
4716 assert!(pool.used() <= 13_000, "usage converges: {}", pool.used());
4717 }
4718
4719 #[tokio::test]
4722 async fn recreated_sequence_bounds_eviction_hints() {
4723 let (mut producer, _pool) = pooled_producer(1 << 40);
4724 producer.create_group(5u64.into()).unwrap().finish().unwrap();
4725
4726 for _ in 0..200 {
4727 let group = producer.create_group(1u64.into()).unwrap();
4728 group.abort(Error::Cancel).unwrap();
4729 }
4730
4731 let state = producer.state.read();
4732 assert!(
4733 state.evict.len() <= 2 * state.lookup.len() + EVICT_SLACK,
4734 "stale hints are compacted: {} entries for {} slots",
4735 state.evict.len(),
4736 state.lookup.len()
4737 );
4738 }
4739
4740 #[tokio::test]
4743 async fn same_tick_write_outranks_inserted() {
4744 tokio::time::pause();
4745
4746 let (mut producer, _pool) = pooled_producer(10_000);
4748
4749 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();
4756 assert!(consumer.peek_group(0).is_none(), "insert-only content pays first");
4757 assert!(consumer.peek_group(1).is_some(), "same-tick written content survives");
4758 }
4759
4760 #[tokio::test]
4763 async fn frame_only_writer_pays() {
4764 tokio::time::pause();
4765
4766 let (mut producer, pool) = pooled_producer(2_000);
4767 let mut demoted = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); demoted
4773 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 300_000]))
4774 .unwrap();
4775
4776 assert!(
4777 pool.used() <= 5_000,
4778 "the frame write settled the debt: {}",
4779 pool.used()
4780 );
4781 assert!(matches!(demoted.finish(), Err(Error::Evicted)));
4782 }
4783
4784 #[tokio::test]
4787 async fn each_track_owns_its_account() {
4788 let broadcast = Arc::new(broadcast::Info::default());
4789 let info = Info::default();
4790 let a = Producer::new(broadcast.clone(), "a", info.clone());
4791 let b = Producer::new(broadcast, "b", info);
4792
4793 let a = a.state.read().cache.clone();
4794 let b = b.state.read().cache.clone();
4795 assert!(!Arc::ptr_eq(&a, &b), "each track owns its account");
4796 }
4797
4798 #[tokio::test]
4801 async fn a_dynamic_defers_teardown() {
4802 let (mut producer, pool) = pooled_producer(1 << 40);
4803 let dynamic = producer.dynamic();
4804 finished_group(&mut producer, 100);
4805
4806 drop(producer);
4807 assert!(pool.used() > 0, "the handler still serves the cache");
4808
4809 drop(dynamic);
4810 assert_eq!(pool.used(), 0, "the last handle tears it down");
4811 }
4812
4813 #[tokio::test]
4819 async fn finished_track_frees_its_cache() {
4820 let (mut producer, pool) = pooled_producer(1 << 40);
4821 finished_group(&mut producer, 100);
4822 producer.finish().unwrap();
4823
4824 let state = producer.state.downgrade();
4825 drop(producer);
4826
4827 assert!(state.upgrade().is_none(), "the track state is freed");
4828 assert_eq!(pool.used(), 0, "so are its cached bytes");
4829 }
4830
4831 #[tokio::test]
4835 async fn teardown_ignores_a_settling_group() {
4836 let (mut producer, pool) = pooled_producer(1 << 40);
4837 finished_group(&mut producer, 100);
4838
4839 let settling = producer.state.downgrade().upgrade().expect("open");
4841 drop(producer);
4842
4843 assert_eq!(pool.used(), 0, "the abrupt teardown still released the cache");
4844 drop(settling);
4845 }
4846
4847 #[tokio::test]
4850 async fn cached_group_outlives_its_track() {
4851 let (mut producer, pool) = pooled_producer(1 << 40);
4852 let sequence = finished_group(&mut producer, 100);
4853 let group = producer.consume().peek_group(sequence).expect("cached");
4854 producer.finish().unwrap();
4855
4856 let state = producer.state.downgrade();
4857 drop(producer);
4858 assert!(state.upgrade().is_none(), "the track state is freed");
4859 assert!(pool.used() > 0, "the retained group keeps its own bytes");
4860
4861 drop(group);
4862 assert_eq!(pool.used(), 0, "which it releases when dropped");
4863 }
4864
4865 #[tokio::test]
4869 async fn pre_accept_backfill_settles_late_writes() {
4870 tokio::time::pause();
4871
4872 let pool = cache::Pool::new(2_000);
4873 let broadcast = broadcast::Info {
4874 origin: crate::origin::Info::default().with_pool(pool.clone()),
4875 ..Default::default()
4876 };
4877 let request = Request::new(Arc::new(broadcast), "test");
4878 let dynamic = request.dynamic();
4879 let consumer = request.consume();
4880
4881 let pending = consumer.fetch_group(0, None);
4883 let req = dynamic
4884 .requested_group()
4885 .now_or_never()
4886 .expect("should not block")
4887 .unwrap();
4888 let mut backfill = req.accept(None).unwrap();
4889 pending.await.unwrap();
4890
4891 let mut producer = request.accept(None);
4893 producer.append_group().unwrap().finish().unwrap();
4894
4895 backfill
4898 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 300_000]))
4899 .unwrap();
4900
4901 assert!(
4902 pool.used() <= 5_000,
4903 "the frame write settled the debt: {}",
4904 pool.used()
4905 );
4906 }
4907
4908 #[tokio::test]
4912 async fn write_restarts_retention_clock() {
4913 tokio::time::pause();
4914
4915 let (mut producer, _pool) = pooled_producer(1 << 40);
4916 let mut straggler = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
4921 straggler
4922 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 100]))
4923 .unwrap();
4924 producer.append_group().unwrap().finish().unwrap(); let consumer = producer.consume();
4927 assert!(consumer.peek_group(0).is_some(), "the write restarted the clock");
4928
4929 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
4931 producer.append_group().unwrap().finish().unwrap(); assert!(consumer.peek_group(0).is_none(), "idle content still expires");
4933 }
4934
4935 #[tokio::test]
4938 async fn refreshed_front_does_not_starve_expiry() {
4939 tokio::time::pause();
4940
4941 let (mut producer, _pool) = pooled_producer(1 << 40);
4942 let dynamic = producer.dynamic();
4943 let consumer = producer.consume();
4944
4945 producer.create_group(10u64.into()).unwrap().finish().unwrap();
4946 for sequence in 1..=5u64 {
4947 let pending = consumer.fetch_group(sequence, None);
4948 let req = dynamic
4949 .requested_group()
4950 .now_or_never()
4951 .expect("should not block")
4952 .unwrap();
4953 let mut group = req.accept(None).unwrap();
4954 group
4955 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 100]))
4956 .unwrap();
4957 group.finish().unwrap();
4958 pending.await.unwrap();
4959 }
4960
4961 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
4964 for sequence in 1..=4u64 {
4965 consumer.fetch_group(sequence, None).await.unwrap();
4966 }
4967
4968 for _ in 0..3 {
4970 producer.append_group().unwrap().finish().unwrap();
4971 }
4972 assert!(consumer.peek_group(5).is_none(), "expired backfill is reclaimed");
4973 assert!(consumer.peek_group(1).is_some(), "refreshed backfill survives");
4974 }
4975
4976 #[tokio::test]
4979 async fn recreated_sequence_delivered_once() {
4980 let (mut producer, _pool) = pooled_producer(1 << 40);
4981
4982 producer.create_group(0u64.into()).unwrap().finish().unwrap();
4983 let aborted = producer.create_group(1u64.into()).unwrap();
4984 aborted.abort(Error::Cancel).unwrap();
4985 producer.create_group(2u64.into()).unwrap().finish().unwrap();
4986 producer.create_group(1u64.into()).unwrap().finish().unwrap();
4987
4988 let mut subscriber = producer.subscribe(None);
4989 assert_eq!(subscriber.assert_group().sequence, 0);
4990 assert_eq!(subscriber.assert_group().sequence, 2);
4991 assert_eq!(
4992 subscriber.assert_group().sequence,
4993 1,
4994 "replacement arrives at its own position"
4995 );
4996 subscriber.assert_no_group();
4997 }
4998
4999 #[tokio::test]
5003 async fn datagrams_do_not_block_eviction() {
5004 tokio::time::pause();
5005
5006 let (mut producer, pool) = pooled_producer(1_000);
5007 for _ in 0..10 {
5008 finished_group(&mut producer, 1_000);
5009 producer.append_datagram(Timestamp::ZERO, &b"beat"[..]).unwrap();
5010 }
5011
5012 let consumer = producer.consume();
5013 assert!(consumer.peek_group(0).is_none(), "old groups still evict");
5014 assert!(
5015 pool.used() < 4 * 1_256,
5016 "interleaved datagrams must not bypass the budget: {}",
5017 pool.used()
5018 );
5019 }
5020
5021 #[tokio::test]
5025 async fn aborted_group_leaves_no_ghost_sample() {
5026 tokio::time::pause();
5027
5028 let (mut producer, pool) = pooled_producer(1 << 40);
5029 let group0 = producer.append_group().unwrap();
5030 producer.append_group().unwrap(); assert!(pool.average().is_some(), "demoted group is sampled");
5033 group0.abort(Error::Cancel).unwrap();
5034 assert_eq!(pool.average(), None, "the abort must remove the sample");
5035 }
5036
5037 #[tokio::test]
5040 async fn empty_groups_repay_overhead() {
5041 tokio::time::pause();
5042
5043 let (mut producer, pool) = pooled_producer(1_000);
5044 for _ in 0..100 {
5045 let mut group = producer.append_group().unwrap();
5046 group.finish().unwrap();
5047 }
5048
5049 assert!(
5050 pool.used() <= 3_000,
5051 "empty-group overhead must stay near the budget: {}",
5052 pool.used()
5053 );
5054 }
5055
5056 #[tokio::test]
5059 async fn growth_on_demoted_group_is_billed() {
5060 tokio::time::pause();
5061
5062 let (mut producer, pool) = pooled_producer(2_000);
5063 let mut straggler = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); straggler
5068 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 10_000]))
5069 .unwrap();
5070
5071 producer.append_group().unwrap().finish().unwrap(); let consumer = producer.consume();
5075 assert!(consumer.peek_group(0).is_none(), "the ballooned group is evicted");
5076 assert!(pool.used() <= 3_000, "growth is reclaimed: {}", pool.used());
5077 }
5078
5079 #[tokio::test]
5082 async fn refilled_sequence_stays_out_of_subscriptions() {
5083 let (mut producer, _pool) = pooled_producer(1 << 40);
5084 let dynamic = producer.dynamic();
5085 let consumer = producer.consume();
5086
5087 producer.create_group(0u64.into()).unwrap().finish().unwrap();
5088 let aborted = producer.create_group(1u64.into()).unwrap();
5089 aborted.abort(Error::Cancel).unwrap();
5090 producer.create_group(2u64.into()).unwrap().finish().unwrap();
5091
5092 let pending = consumer.fetch_group(1, None);
5095 let req = dynamic
5096 .requested_group()
5097 .now_or_never()
5098 .expect("should not block")
5099 .unwrap();
5100 let mut group = req.accept(None).unwrap();
5101 group
5102 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"backfill"))
5103 .unwrap();
5104 group.finish().unwrap();
5105 pending.await.unwrap();
5106
5107 assert!(consumer.peek_group(1).is_some());
5109 let mut subscriber = producer.subscribe(None);
5110 assert_eq!(subscriber.assert_group().sequence, 0);
5111 assert_eq!(subscriber.assert_group().sequence, 2);
5112 subscriber.assert_no_group();
5113 }
5114
5115 #[tokio::test]
5118 async fn expired_backfill_behind_refreshed_reclaimed() {
5119 tokio::time::pause();
5120
5121 let (mut producer, _pool) = pooled_producer(1 << 40);
5122 let dynamic = producer.dynamic();
5123 let consumer = producer.consume();
5124
5125 producer.create_group(5u64.into()).unwrap().finish().unwrap();
5126 for sequence in [2u64, 3u64] {
5127 let pending = consumer.fetch_group(sequence, None);
5128 let req = dynamic
5129 .requested_group()
5130 .now_or_never()
5131 .expect("should not block")
5132 .unwrap();
5133 let mut group = req.accept(None).unwrap();
5134 group
5135 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 100]))
5136 .unwrap();
5137 group.finish().unwrap();
5138 pending.await.unwrap();
5139 }
5140
5141 tokio::time::advance(Duration::from_secs(4)).await;
5143 consumer.fetch_group(2, None).await.unwrap();
5144 tokio::time::advance(DEFAULT_LATENCY_MAX - Duration::from_secs(2)).await;
5145 producer.create_group(6u64.into()).unwrap().finish().unwrap();
5146
5147 let consumer = producer.consume();
5148 assert!(consumer.peek_group(2).is_some(), "refreshed backfill survives");
5149 assert!(consumer.peek_group(3).is_none(), "expired backfill is reclaimed");
5150 }
5151
5152 #[tokio::test]
5155 async fn same_tick_fetch_protects() {
5156 tokio::time::pause();
5157
5158 let (mut producer, _pool) = pooled_producer(10_000);
5160 let consumer = producer.consume();
5161
5162 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();
5167
5168 finished_group(&mut producer, 3_000); finished_group(&mut producer, 3_000); assert!(consumer.peek_group(0).is_some(), "same-tick refresh protects");
5172 assert!(consumer.peek_group(1).is_none(), "the unread group dies instead");
5173 }
5174
5175 #[tokio::test]
5179 async fn refetched_latest_stays_protected() {
5180 tokio::time::pause();
5181
5182 let (mut producer, _pool) = pooled_producer(10_000);
5183 let dynamic = producer.dynamic();
5184 let consumer = producer.consume();
5185
5186 let straggler = producer.append_group().unwrap(); let latest = producer.append_group().unwrap(); latest.abort(Error::Cancel).unwrap();
5191
5192 let pending = consumer.fetch_group(1, None);
5194 let req = dynamic
5195 .requested_group()
5196 .now_or_never()
5197 .expect("should not block")
5198 .unwrap();
5199 let mut group = req.accept(None).unwrap();
5200 group
5201 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 1000]))
5202 .unwrap();
5203 group.finish().unwrap();
5204 pending.await.unwrap();
5205
5206 {
5209 let state = producer.state.read();
5210 assert!(state.lookup.contains_key(&1), "refetched group is cached");
5211 assert!(
5212 state.evict.iter().all(|(sequence, _)| *sequence != 1),
5213 "the live edge must not be an eviction candidate"
5214 );
5215 }
5216 drop(straggler);
5217 }
5218
5219 #[tokio::test]
5222 async fn eviction_allows_refetch() {
5223 tokio::time::pause();
5224
5225 let (mut producer, _pool) = pooled_producer(10_000);
5226 let dynamic = producer.dynamic();
5227
5228 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let consumer = producer.consume();
5233 assert!(consumer.peek_group(0).is_none());
5234 let pending = consumer.fetch_group(0, None);
5235
5236 let req = dynamic
5237 .requested_group()
5238 .now_or_never()
5239 .expect("should not block")
5240 .unwrap();
5241 assert_eq!(req.sequence(), 0);
5242
5243 let mut group = req.accept(None).unwrap();
5244 group
5245 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"refetched"))
5246 .unwrap();
5247 group.finish().unwrap();
5248
5249 let mut group = pending.await.unwrap();
5250 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"refetched");
5251 }
5252
5253 #[tokio::test]
5256 async fn fetched_backfill_not_subscribed() {
5257 let (mut producer, _pool) = pooled_producer(1 << 40);
5258 let dynamic = producer.dynamic();
5259 let consumer = producer.consume();
5260
5261 producer.create_group(5u64.into()).unwrap().finish().unwrap();
5263 producer.create_group(6u64.into()).unwrap().finish().unwrap();
5264
5265 let pending = consumer.fetch_group(2, None);
5267 let req = dynamic
5268 .requested_group()
5269 .now_or_never()
5270 .expect("should not block")
5271 .unwrap();
5272 let mut group = req.accept(None).unwrap();
5273 group
5274 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"backfill"))
5275 .unwrap();
5276 group.finish().unwrap();
5277 let mut fetched = pending.await.unwrap();
5278 assert_eq!(&fetched.read_frame().await.unwrap().unwrap().payload[..], b"backfill");
5279 assert!(consumer.peek_group(2).is_some(), "backfill is cached for later fetches");
5280
5281 let mut subscriber = producer.subscribe(None);
5283 assert_eq!(subscriber.assert_group().sequence, 5);
5284 assert_eq!(subscriber.assert_group().sequence, 6);
5285 subscriber.assert_no_group();
5286 }
5287
5288 #[tokio::test]
5291 async fn expired_backfill_reclaimed() {
5292 tokio::time::pause();
5293
5294 let (mut producer, pool) = pooled_producer(1 << 40);
5295 let dynamic = producer.dynamic();
5296 let consumer = producer.consume();
5297
5298 producer.create_group(5u64.into()).unwrap().finish().unwrap();
5299
5300 let pending = consumer.fetch_group(2, None);
5302 let req = dynamic
5303 .requested_group()
5304 .now_or_never()
5305 .expect("should not block")
5306 .unwrap();
5307 let mut group = req.accept(None).unwrap();
5308 group
5309 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 1000]))
5310 .unwrap();
5311 group.finish().unwrap();
5312 pending.await.unwrap();
5313 let used = pool.used();
5314
5315 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
5317 producer.create_group(6u64.into()).unwrap().finish().unwrap();
5318
5319 assert!(consumer.peek_group(2).is_none(), "expired backfill is reclaimed");
5320 assert!(pool.used() < used, "its bytes are released");
5321 }
5322
5323 #[tokio::test]
5324 async fn fetch_aborts_with_track() {
5325 let producer = track_producer("test", None);
5326 let dynamic = producer.dynamic();
5327 let consumer = producer.consume();
5328
5329 let pending = consumer.fetch_group(3, None);
5330 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
5331
5332 producer.abort(Error::Cancel).unwrap();
5333 assert!(pending.await.is_err());
5334 drop(dynamic);
5335 }
5336}