1use crate::{Error, Result, Timescale, Timestamp, coding};
17use crate::{broadcast, cache, frame, group, stats};
18
19use super::{Datagram, Requests};
20
21pub use super::subscription::Subscription;
22
23use std::{
24 collections::{BTreeMap, HashMap, VecDeque},
25 sync::Arc,
26 sync::OnceLock,
27 sync::atomic::{AtomicBool, Ordering},
28 task::{Poll, ready},
29 time::Duration,
30};
31
32pub const DEFAULT_LATENCY_MAX: Duration = Duration::from_secs(5);
34
35const MAX_DATAGRAM_AGE: Duration = Duration::from_millis(50);
41
42const EVICT_SLACK: usize = 64;
45
46const EVICT_SCAN: usize = 4;
50
51#[derive(Clone, Debug)]
62#[non_exhaustive]
63pub struct Info {
64 pub timescale: Timescale,
71 pub latency_max: Duration,
80 pub priority: u8,
83 pub ordered: bool,
87}
88
89impl Default for Info {
90 fn default() -> Self {
91 Self {
92 timescale: Timescale::default(),
93 latency_max: DEFAULT_LATENCY_MAX,
94 priority: 0,
95 ordered: false,
96 }
97 }
98}
99
100impl Info {
101 pub fn with_timescale(mut self, timescale: Timescale) -> Self {
106 self.timescale = timescale;
107 self
108 }
109
110 pub fn with_latency_max(mut self, latency_max: Duration) -> Self {
112 self.latency_max = latency_max;
113 self
114 }
115
116 pub fn with_priority(mut self, priority: u8) -> Self {
118 self.priority = priority;
119 self
120 }
121
122 pub fn with_ordered(mut self, ordered: bool) -> Self {
126 self.ordered = ordered;
127 self
128 }
129}
130
131#[derive(Default)]
132pub(crate) struct TrackState {
133 info: Option<Info>,
136
137 broadcast: Arc<broadcast::Info>,
140
141 cache: Arc<cache::Track>,
145
146 lookup: HashMap<u64, Slot>,
150
151 arrival: VecDeque<(u64, u32)>,
156
157 evict: VecDeque<(u64, u32)>,
165
166 debt: u64,
171
172 datagrams: VecDeque<(Datagram, web_async::time::Instant)>,
176
177 datagram_offset: usize,
180
181 offset: usize,
184
185 max_sequence: Option<u64>,
188
189 latest_group: Option<u64>,
194
195 next_stamp: u32,
197
198 expire_cursor: usize,
201
202 final_sequence: Option<u64>,
204
205 abort: Option<Error>,
207
208 subscriptions: kio::Shared<Subscriptions>,
212
213 fetch: kio::Shared<FetchState>,
216}
217
218struct Slot {
224 group: group::Producer,
225
226 stamp: u32,
231}
232
233type Subscriptions = Vec<kio::Consumer<Subscription>>;
235
236type FetchState = Requests<u64, PendingFetch>;
241
242struct PendingFetch {
244 priority: u8,
246
247 result: kio::Producer<FetchOutcome>,
252}
253
254#[derive(Default)]
257struct FetchOutcome {
258 rejected: Option<Error>,
259}
260
261impl TrackState {
262 fn poll_info(&self) -> Poll<Result<Info>> {
263 if let Some(info) = &self.info {
264 Poll::Ready(Ok(info.clone()))
265 } else {
266 Poll::Pending
267 }
268 }
269
270 fn poll_recv_group(&self, index: usize, min_sequence: u64) -> Poll<Result<Option<(group::Consumer, usize)>>> {
274 let start = index.saturating_sub(self.offset);
275 for (i, (sequence, stamp)) in self.arrival.iter().enumerate().skip(start) {
276 if *sequence >= min_sequence
277 && let Some(slot) = self.lookup.get(sequence)
278 && slot.stamp == *stamp
279 && !slot.group.is_aborted()
280 {
281 return Poll::Ready(Ok(Some((slot.group.consume(), self.offset + i))));
282 }
283 }
284
285 if self.is_complete() {
287 Poll::Ready(Ok(None))
288 } else if let Some(err) = &self.abort {
289 Poll::Ready(Err(err.clone()))
290 } else {
291 Poll::Pending
292 }
293 }
294
295 fn poll_recv_datagram(&self, index: usize) -> Poll<Result<Option<(Datagram, usize)>>> {
301 let start = index.saturating_sub(self.datagram_offset);
302 if let Some((datagram, _)) = self.datagrams.get(start) {
303 return Poll::Ready(Ok(Some((datagram.clone(), self.datagram_offset + start))));
304 }
305
306 if self.is_complete() {
308 Poll::Ready(Ok(None))
309 } else if let Some(err) = &self.abort {
310 Poll::Ready(Err(err.clone()))
311 } else {
312 Poll::Pending
313 }
314 }
315
316 fn push_datagram(&mut self, datagram: Datagram) {
318 let now = web_async::time::Instant::now();
319 self.datagrams.push_back((datagram, now));
320 while let Some((_, at)) = self.datagrams.front() {
321 if now.duration_since(*at) <= MAX_DATAGRAM_AGE {
322 break;
323 }
324 self.datagrams.pop_front();
325 self.datagram_offset += 1;
326 }
327 }
328
329 fn poll_read_frame(
333 &self,
334 index: usize,
335 next_sequence: u64,
336 waiter: &kio::Waiter,
337 ) -> Poll<Result<Option<(frame::Frame, usize, u64)>>> {
338 let start = index.saturating_sub(self.offset);
339 let mut pending_seen = false;
340 for (i, (sequence, stamp)) in self.arrival.iter().enumerate().skip(start) {
341 if *sequence < next_sequence {
342 continue;
343 }
344 let Some(slot) = self.lookup.get(sequence) else {
345 continue;
346 };
347 if slot.stamp != *stamp {
348 continue;
351 }
352
353 let mut consumer = slot.group.consume();
354 match consumer.poll_read_frame(waiter) {
355 Poll::Ready(Ok(Some(frame))) => {
356 return Poll::Ready(Ok(Some((frame, self.offset + i, *sequence))));
357 }
358 Poll::Ready(Ok(None)) => continue,
359 Poll::Ready(Err(_)) => continue,
362 Poll::Pending => {
363 pending_seen = true;
364 continue;
365 }
366 }
367 }
368
369 if pending_seen {
372 Poll::Pending
373 } else if self.is_complete() {
374 Poll::Ready(Ok(None))
375 } else if let Some(err) = &self.abort {
376 Poll::Ready(Err(err.clone()))
377 } else {
378 Poll::Pending
379 }
380 }
381
382 fn poll_next_in_range(
392 &self,
393 next_sequence: u64,
394 end_sequence: Option<u64>,
395 ) -> Poll<Result<Option<group::Consumer>>> {
396 if let Some(end) = end_sequence
400 && end < next_sequence
401 {
402 if let Some(err) = &self.abort {
403 return Poll::Ready(Err(err.clone()));
404 }
405 return Poll::Pending;
406 }
407
408 let mut best: Option<&group::Producer> = None;
409 for slot in self.lookup.values() {
410 let group = &slot.group;
411 if group.sequence < next_sequence {
412 continue;
413 }
414 if let Some(end) = end_sequence
415 && group.sequence > end
416 {
417 continue;
418 }
419 if group.is_aborted() {
420 continue;
421 }
422 if best.is_none_or(|b| group.sequence < b.sequence) {
423 best = Some(group);
424 }
425 }
426
427 if let Some(group) = best {
428 return Poll::Ready(Ok(Some(group.consume())));
429 }
430
431 if let Some(err) = &self.abort {
433 return Poll::Ready(Err(err.clone()));
434 }
435 if let Some(fin) = self.final_sequence
438 && next_sequence >= fin
439 {
440 return Poll::Ready(Ok(None));
441 }
442 Poll::Pending
443 }
444
445 #[cfg(test)]
450 fn cached_group(&self, sequence: u64) -> Option<group::Consumer> {
451 let slot = self.lookup.get(&sequence)?;
452 if slot.group.is_aborted() {
453 return None;
454 }
455 Some(slot.group.consume())
456 }
457
458 fn latency_bound(&self) -> Option<Duration> {
461 self.info.as_ref().map(|info| info.latency_max)
462 }
463
464 fn poll_fetch_cached(&self, sequence: u64) -> Poll<Result<group::Consumer>> {
469 if let Some(slot) = self.lookup.get(&sequence)
470 && !slot.group.is_aborted()
471 {
472 slot.group.cache_refresh();
476 return Poll::Ready(Ok(slot.group.consume()));
477 }
478
479 if let Some(err) = &self.abort {
480 return Poll::Ready(Err(err.clone()));
481 }
482
483 if self.final_sequence.is_some_and(|fin| sequence >= fin) {
485 return Poll::Ready(Err(Error::NotFound));
486 }
487
488 Poll::Pending
489 }
490
491 fn evict_expired(&mut self, max_age: Duration) {
500 let now = self.cache.pool().now();
501 let max_ticks = cache::Pool::ticks(max_age);
502
503 let len = self.evict.len();
504 if len > 0 {
505 let start = self.expire_cursor % len;
506 for step in 0..len.min(EVICT_SCAN) {
507 let (sequence, stamp) = self.evict[(start + step) % len];
508 let Some(slot) = self.lookup.get(&sequence) else {
509 continue;
510 };
511 if slot.stamp != stamp {
512 continue;
514 }
515 if slot.group.is_aborted() {
518 self.lookup.remove(&sequence);
519 continue;
520 }
521 if Some(sequence) == self.latest_group || now.saturating_sub(slot.group.cache_accessed()) <= max_ticks {
522 continue;
523 }
524 let slot = self.lookup.remove(&sequence).unwrap();
528 let _ = slot.group.abort(Error::Old);
529 }
530 self.expire_cursor = (start + EVICT_SCAN) % len;
531 }
532
533 while let Some((sequence, stamp)) = self.arrival.front() {
536 if self.lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp) {
537 break;
538 }
539 self.arrival.pop_front();
540 self.offset += 1;
541 }
542
543 while let Some((sequence, stamp)) = self.evict.front() {
545 if self.lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp) {
546 break;
547 }
548 self.evict.pop_front();
549 }
550
551 if self.evict.len() > 2 * self.lookup.len() + EVICT_SLACK {
554 let lookup = &self.lookup;
555 self.evict
556 .retain(|(sequence, stamp)| lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp));
557 }
558 }
559
560 fn clear_cache(&mut self) {
563 self.lookup.clear();
564 self.arrival.clear();
565 self.evict.clear();
566 self.latest_group = None;
567 self.debt = 0;
568 }
569
570 fn install(&mut self, mut info: Info) {
576 info.latency_max = info.latency_max.min(self.broadcast.origin.cache_duration);
577 self.info = Some(info);
578 }
579
580 fn spawn(broadcast: Arc<broadcast::Info>) -> kio::Producer<Self> {
587 let state = kio::Producer::new(Self {
588 broadcast: broadcast.clone(),
589 ..Default::default()
590 });
591 let cache = cache::Track::new(broadcast.origin.pool.clone(), state.downgrade());
592 state.write().ok().expect("a new track is open").cache = cache;
593 state
594 }
595
596 fn claim_sequence(&mut self, sequence: u64) -> Result<()> {
602 if let Some(slot) = self.lookup.get(&sequence) {
603 if !slot.group.is_aborted() {
604 return Err(Error::Duplicate);
605 }
606 self.lookup.remove(&sequence);
607 }
608 Ok(())
609 }
610
611 fn insert_group(&mut self, group: &group::Producer, visible: bool) {
618 let sequence = group.sequence;
619 self.next_stamp = self.next_stamp.wrapping_add(1);
620 let stamp = self.next_stamp;
621
622 if self.latest_group.is_none_or(|latest| sequence >= latest) {
626 if let Some(latest) = self.latest_group
629 && sequence > latest
630 && let Some(prev) = self.lookup.get(&latest)
631 {
632 prev.group.cache_demote();
633 self.evict.push_back((latest, prev.stamp));
634 }
635 self.latest_group = Some(sequence);
636 } else {
637 group.cache_demote();
638 self.evict.push_back((sequence, stamp));
639 }
640
641 self.max_sequence = Some(self.max_sequence.map_or(sequence, |max| max.max(sequence)));
642 self.lookup.insert(
643 sequence,
644 Slot {
645 group: group.clone(),
646 stamp,
647 },
648 );
649 if visible {
650 self.arrival.push_back((sequence, stamp));
651 }
652 }
653
654 fn commit_group(&mut self, group: &group::Producer, visible: bool, latency_max: Duration) {
658 self.charge_debt();
659 self.insert_group(group, visible);
660 self.evict_expired(latency_max);
661 }
662
663 pub(super) fn charge_debt(&mut self) {
676 let written = self.cache.take_written();
677 let pool = self.cache.pool().clone();
678 match pool.accrue(written) {
679 Some(mut accrued) => {
680 if self.oldest_is_stale(&pool) {
681 accrued = accrued.saturating_mul(2);
682 }
683 self.debt = self.debt.saturating_add(accrued).min(pool.used());
686 self.pay_debt(&pool, written.saturating_mul(2));
689 }
690 None => self.debt = 0,
693 }
694 }
695
696 fn oldest_is_stale(&self, pool: &cache::Pool) -> bool {
700 let Some(average) = pool.average() else {
701 return false;
702 };
703 let Some((sequence, stamp)) = self.evict.front() else {
704 return false;
705 };
706 let Some(slot) = self.lookup.get(sequence) else {
707 return false;
708 };
709 slot.stamp == *stamp && !slot.group.is_aborted() && slot.group.cache_accessed() <= average
710 }
711
712 fn pay_debt(&mut self, pool: &cache::Pool, cap: u64) {
725 let average = pool.average().unwrap_or(0);
726 let mut paid = 0u64;
727 let mut scanned = 0usize;
728 for _ in 0..self.evict.len() {
729 if self.debt == 0 || paid >= cap || scanned >= EVICT_SCAN {
730 return;
731 }
732 let Some((sequence, stamp)) = self.evict.pop_front() else {
733 return;
734 };
735 let Some(slot) = self.lookup.get(&sequence) else {
736 continue;
738 };
739 if slot.stamp != stamp {
740 continue;
742 }
743 if slot.group.is_aborted() {
744 self.lookup.remove(&sequence);
746 continue;
747 }
748 if Some(sequence) == self.latest_group {
749 self.evict.push_back((sequence, stamp));
751 continue;
752 }
753
754 scanned += 1;
755 if slot.group.cache_accessed() > average {
759 self.evict.push_back((sequence, stamp));
760 continue;
761 }
762 let size = slot.group.cache_size();
765 if size > self.debt {
766 self.evict.push_front((sequence, stamp));
767 return;
768 }
769
770 self.debt -= size;
771 paid = paid.saturating_add(size);
772 let slot = self.lookup.remove(&sequence).unwrap();
773 let _ = slot.group.abort(Error::Evicted);
774 }
775 }
776
777 fn set_final(&mut self, final_sequence: u64) -> Result<()> {
780 if self.final_sequence.is_some() {
781 return Err(Error::Closed);
782 }
783 if let Some(max) = self.max_sequence
784 && final_sequence <= max
785 {
786 return Err(Error::ProtocolViolation);
787 }
788 self.final_sequence = Some(final_sequence);
789 Ok(())
790 }
791
792 fn is_complete(&self) -> bool {
798 self.final_sequence
799 .is_some_and(|fin| self.max_sequence.map_or(0, |max| max.saturating_add(1)) >= fin)
800 }
801
802 fn poll_finished(&self) -> Poll<Result<u64>> {
803 if let Some(fin) = self.final_sequence {
804 Poll::Ready(Ok(fin))
805 } else if let Some(err) = &self.abort {
806 Poll::Ready(Err(err.clone()))
807 } else {
808 Poll::Pending
809 }
810 }
811
812 fn modify(producer: &kio::Producer<Self>) -> Result<kio::Mut<'_, Self>> {
813 producer.write().map_err(|r| r.abort.clone().unwrap_or(Error::Dropped))
814 }
815
816 fn insert_group_request(&mut self, sequence: u64, info: Option<Info>) -> Result<group::Producer> {
822 if let Some(err) = &self.abort {
823 return Err(err.clone());
824 }
825 if let Some(fin) = self.final_sequence
826 && sequence >= fin
827 {
828 return Err(Error::Closed);
829 }
830
831 if self.info.is_none() {
835 self.install(info.unwrap_or_default());
836 }
837 let info = self.info.clone().unwrap();
838
839 self.claim_sequence(sequence)?;
841
842 let latency_max = info.latency_max;
843 let group = group::Producer::new(group::Info { sequence }, info, self.cache.clone());
844 group.cache_refresh();
849 self.commit_group(&group, false, latency_max);
850 Ok(group)
851 }
852}
853
854#[derive(Clone)]
856pub struct Producer {
857 name: Arc<str>,
858 broadcast: Arc<broadcast::Info>,
861 state: kio::Producer<TrackState>,
862 prev_subscription: Option<Subscription>,
863 alive: Arc<Alive>,
865 stats: stats::Scope,
869}
870
871impl Producer {
872 pub(crate) fn new(
880 broadcast: Arc<broadcast::Info>,
881 name: impl Into<Arc<str>>,
882 info: impl Into<Option<Info>>,
883 ) -> Self {
884 let name = name.into();
885 let state = TrackState::spawn(broadcast.clone());
886 state
887 .write()
888 .ok()
889 .expect("a new track is open")
890 .install(info.into().unwrap_or_default());
891 let alive = Alive::new(name.clone(), state.clone());
892 alive.publish(None);
893 Self {
894 name,
895 state,
896 broadcast,
897 prev_subscription: None,
898 alive,
899 stats: stats::Scope::default(),
900 }
901 }
902
903 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
907 self.alive.publish(Some(&scope));
908 self.stats = scope;
909 self
910 }
911
912 pub fn name(&self) -> &str {
914 &self.name
915 }
916
917 pub fn broadcast(&self) -> &broadcast::Info {
919 &self.broadcast
920 }
921
922 pub fn create_group(&mut self, group: group::Info) -> Result<group::Producer> {
924 let mut state = self.modify()?;
925 if let Some(fin) = state.final_sequence
926 && group.sequence >= fin
927 {
928 return Err(Error::Closed);
929 }
930 let track = state.info.clone().unwrap();
931 let latency_max = track.latency_max;
932
933 state.claim_sequence(group.sequence)?;
935
936 let group = group::Producer::new(group, track, state.cache.clone()).with_meter(self.stats.meter());
937 state.commit_group(&group, true, latency_max);
938
939 Ok(group)
940 }
941
942 pub fn append_group(&mut self) -> Result<group::Producer> {
944 let mut state = self.modify()?;
945 let sequence = match state.max_sequence {
946 Some(s) => s.checked_add(1).ok_or(coding::BoundsExceeded)?,
947 None => 0,
948 };
949 if let Some(fin) = state.final_sequence
950 && sequence >= fin
951 {
952 return Err(Error::Closed);
953 }
954
955 let track = state.info.clone().unwrap();
956 let latency_max = track.latency_max;
957
958 let group =
959 group::Producer::new(group::Info { sequence }, track, state.cache.clone()).with_meter(self.stats.meter());
960 state.commit_group(&group, true, latency_max);
961
962 Ok(group)
963 }
964
965 pub fn append_datagram<B: crate::IntoBytes>(&mut self, timestamp: Timestamp, payload: B) -> Result<u64> {
977 let payload = payload.into_bytes();
978 if payload.len() > super::datagram::MAX_DATAGRAM_PAYLOAD {
979 return Err(Error::FrameTooLarge);
980 }
981 let meter = self.stats.meter();
983 let mut state = self.modify()?;
984 let timescale = state.info.as_ref().unwrap().timescale;
986 let timestamp = timestamp.convert(timescale).map_err(|_| Error::TimestampMismatch)?;
987 let sequence = match state.max_sequence {
988 Some(s) => s.checked_add(1).ok_or(coding::BoundsExceeded)?,
989 None => 0,
990 };
991 if let Some(fin) = state.final_sequence
992 && sequence >= fin
993 {
994 return Err(Error::Closed);
995 }
996 state.max_sequence = Some(sequence);
997 meter.datagram(payload.len() as u64);
998 state.push_datagram(Datagram {
999 sequence,
1000 timestamp,
1001 payload,
1002 });
1003 Ok(sequence)
1004 }
1005
1006 pub fn write_datagram(&mut self, mut datagram: Datagram) -> Result<()> {
1012 if datagram.payload.len() > super::datagram::MAX_DATAGRAM_PAYLOAD {
1013 return Err(Error::FrameTooLarge);
1014 }
1015 let meter = self.stats.meter();
1017 let mut state = self.modify()?;
1018 let timescale = state.info.as_ref().unwrap().timescale;
1020 datagram.timestamp = datagram
1021 .timestamp
1022 .convert(timescale)
1023 .map_err(|_| Error::TimestampMismatch)?;
1024 if let Some(fin) = state.final_sequence
1025 && datagram.sequence >= fin
1026 {
1027 return Err(Error::Closed);
1028 }
1029 state.max_sequence = Some(state.max_sequence.unwrap_or(0).max(datagram.sequence));
1030 meter.datagram(datagram.payload.len() as u64);
1031 state.push_datagram(datagram);
1032 Ok(())
1033 }
1034
1035 pub fn write_frame<B: crate::IntoBytes>(&mut self, timestamp: Timestamp, frame: B) -> Result<()> {
1040 let mut group = self.append_group()?;
1041 group.write_frame(timestamp, frame)?;
1042 group.finish()?;
1043 Ok(())
1044 }
1045
1046 pub fn finish(&mut self) -> Result<()> {
1052 let mut state = self.modify()?;
1053 let final_sequence = match state.max_sequence {
1054 Some(max) => max.checked_add(1).ok_or(coding::BoundsExceeded)?,
1055 None => 0,
1056 };
1057 state.set_final(final_sequence)
1058 }
1059
1060 pub fn finish_at(&mut self, final_sequence: u64) -> Result<()> {
1073 self.modify()?.set_final(final_sequence)
1074 }
1075
1076 pub fn final_sequence(&self) -> Option<u64> {
1081 self.state.read().final_sequence
1082 }
1083
1084 pub fn abort(self, err: Error) -> Result<()> {
1095 let mut guard = self.modify()?;
1096 guard.abort = Some(err);
1097 guard.clear_cache();
1098 guard.datagrams.clear();
1099 guard.close();
1100 Ok(())
1101 }
1102
1103 pub async fn unused(&self) -> Result<()> {
1105 self.state.unused().await.map_err(|_| self.abort_reason())
1106 }
1107
1108 pub async fn used(&self) -> Result<()> {
1110 self.state.used().await.map_err(|_| self.abort_reason())
1111 }
1112
1113 pub async fn closed(&self) -> Error {
1115 kio::wait(|waiter| self.poll_closed(waiter)).await
1116 }
1117
1118 pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<Error> {
1120 self.state.poll_closed(waiter).map(|()| self.abort_reason())
1121 }
1122
1123 fn abort_reason(&self) -> Error {
1125 self.state.read().abort.clone().unwrap_or(Error::Dropped)
1126 }
1127
1128 pub fn is_closed(&self) -> bool {
1130 self.state.read().is_closed()
1131 }
1132
1133 pub fn latest(&self) -> Option<u64> {
1135 self.state.read().max_sequence
1136 }
1137
1138 pub fn is_clone(&self, other: &Self) -> bool {
1140 self.state.same_channel(&other.state)
1141 }
1142
1143 pub(crate) fn weak(&self) -> TrackWeak {
1145 TrackWeak {
1146 name: self.name.clone(),
1147 state: self.state.weak(),
1148 }
1149 }
1150
1151 pub fn demand(&self) -> Demand {
1159 Demand {
1160 name: self.name.clone(),
1161 state: self.state.weak(),
1162 }
1163 }
1164
1165 pub fn consume(&self) -> Consumer {
1170 Consumer::plain(self.name.clone(), self.state.consume())
1171 }
1172
1173 pub fn subscribe(&self, subscription: impl Into<Option<Subscription>>) -> Subscriber {
1178 let preferences = subscription.into().unwrap_or_default();
1179
1180 let info = self.state.read().info.clone().expect("producer always has info");
1185 let subscription = kio::Producer::new(preferences);
1186 register_subscription(self.state.read(), &subscription);
1187
1188 Subscriber {
1189 name: self.name.clone(),
1190 info,
1191 inner: SubscriberKind::Plain(PlainSubscriber {
1192 state: self.state.consume(),
1193 subscription,
1194 index: 0,
1195 datagram_index: 0,
1196 min_sequence: 0,
1197 next_sequence: 0,
1198 end_sequence: None,
1199 parked: BTreeMap::new(),
1200 }),
1201 stats: stats::Scope::default(),
1203 _stats_sub: stats::Subscription::default(),
1204 }
1205 }
1206
1207 pub async fn subscription_changed(&mut self) -> Result<Option<Subscription>> {
1213 kio::wait(|waiter| self.poll_subscription_changed(waiter)).await
1214 }
1215
1216 pub fn subscription(&self) -> Option<Subscription> {
1224 let state = self.state.read();
1225 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
1226 drop(state);
1227 snapshot_subscription(&subs, bound)
1228 }
1229
1230 pub fn poll_subscription_changed(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Subscription>>> {
1234 if self.state.poll_closed(waiter).is_ready() {
1237 let abort = self.state.read().abort.clone();
1238 return Poll::Ready(Err(abort.unwrap_or(Error::Dropped)));
1239 }
1240
1241 let state = self.state.read();
1243 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
1244 drop(state);
1245
1246 let prev = &self.prev_subscription;
1247 let mut combined = None;
1248 let mut guard = match subs.poll(waiter, |subs| {
1249 let next = combined_subscription(subs, bound, waiter);
1250 if &next == prev {
1251 Poll::Pending
1252 } else {
1253 combined = next;
1254 Poll::Ready(())
1255 }
1256 }) {
1257 Poll::Ready(guard) => guard,
1258 Poll::Pending => return Poll::Pending,
1259 };
1260 guard.retain(|sub| !sub.is_closed());
1262 drop(guard);
1263 self.prev_subscription = combined.clone();
1264 Poll::Ready(Ok(combined))
1265 }
1266
1267 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
1269 self.state.poll_unused(waiter).map(|_| ())
1270 }
1271
1272 pub fn dynamic(&self) -> Dynamic {
1276 Dynamic::new(self.name.clone(), self.state.clone(), self.alive.clone())
1277 }
1278
1279 fn modify(&self) -> Result<kio::Mut<'_, TrackState>> {
1280 TrackState::modify(&self.state)
1281 }
1282}
1283
1284fn poll_requested_group(
1288 state: &kio::Producer<TrackState>,
1289 fetch: &kio::Shared<FetchState>,
1290 waiter: &kio::Waiter,
1291) -> Poll<Result<GroupRequest>> {
1292 if let Poll::Ready(mut guard) = fetch.poll(waiter, |fetch| {
1294 if fetch.has_queued() {
1295 Poll::Ready(())
1296 } else {
1297 Poll::Pending
1298 }
1299 }) {
1300 let sequence = guard.pop().expect("predicate guaranteed a request");
1301 let pending = guard.get(&sequence).expect("popped key must be pending");
1305 let priority = pending.priority;
1306 let result = pending.result.clone();
1307 drop(guard);
1308 return Poll::Ready(Ok(GroupRequest {
1309 state: state.clone(),
1310 fetch: fetch.clone(),
1311 sequence,
1312 priority,
1313 result,
1314 done: false,
1315 }));
1316 }
1317
1318 match state.poll_ref(waiter, |state| match &state.abort {
1320 Some(err) => Poll::Ready(err.clone()),
1321 None => Poll::Pending,
1322 }) {
1323 Poll::Ready(Ok(err)) => Poll::Ready(Err(err)),
1324 Poll::Ready(Err(closed)) => Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped))),
1325 Poll::Pending => Poll::Pending,
1326 }
1327}
1328
1329pub struct Dynamic {
1339 name: Arc<str>,
1340 state: kio::Producer<TrackState>,
1342 fetch: kio::Shared<FetchState>,
1344 alive: Arc<Alive>,
1347}
1348
1349impl Dynamic {
1350 fn new(name: Arc<str>, state: kio::Producer<TrackState>, alive: Arc<Alive>) -> Self {
1351 let fetch = state.read().fetch.clone();
1352 fetch.lock().add_handler();
1353 Self {
1354 name,
1355 state,
1356 fetch,
1357 alive,
1358 }
1359 }
1360
1361 pub fn name(&self) -> &str {
1363 &self.name
1364 }
1365
1366 pub async fn requested_group(&self) -> Result<GroupRequest> {
1372 kio::wait(|waiter| self.poll_requested_group(waiter)).await
1373 }
1374
1375 pub fn poll_requested_group(&self, waiter: &kio::Waiter) -> Poll<Result<GroupRequest>> {
1377 poll_requested_group(&self.state, &self.fetch, waiter)
1378 }
1379
1380 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
1382 self.state.poll_unused(waiter).map(|_| ())
1383 }
1384}
1385
1386impl Clone for Dynamic {
1387 fn clone(&self) -> Self {
1388 self.fetch.lock().add_handler();
1390 Self {
1391 name: self.name.clone(),
1392 state: self.state.clone(),
1393 fetch: self.fetch.clone(),
1394 alive: self.alive.clone(),
1395 }
1396 }
1397}
1398
1399impl Drop for Dynamic {
1400 fn drop(&mut self) {
1401 let mut fetch = self.fetch.lock();
1407 if fetch.remove_handler() {
1408 fetch.drain_queued();
1409 }
1410 }
1411}
1412
1413struct Alive {
1422 name: Arc<str>,
1423 state: kio::Producer<TrackState>,
1424
1425 published: AtomicBool,
1428
1429 stats: OnceLock<stats::Subscription>,
1432}
1433
1434impl Alive {
1435 fn new(name: Arc<str>, state: kio::Producer<TrackState>) -> Arc<Self> {
1436 Arc::new(Self {
1437 name,
1438 state,
1439 published: Default::default(),
1440 stats: Default::default(),
1441 })
1442 }
1443
1444 fn publish(&self, stats: Option<&stats::Scope>) {
1448 self.published.store(true, Ordering::Relaxed);
1449 if let Some(scope) = stats {
1450 let _ = self.stats.set(scope.subscribe());
1453 }
1454 }
1455}
1456
1457impl Drop for Alive {
1458 fn drop(&mut self) {
1459 if !self.published.load(Ordering::Relaxed) {
1461 return;
1462 }
1463 match self.state.write() {
1471 Ok(mut state) => {
1472 if state.final_sequence.is_some() || state.abort.is_some() {
1473 return;
1474 }
1475 tracing::warn!(
1476 track = %self.name,
1477 "track::Producer dropped without finish() or abort()"
1478 );
1479 state.clear_cache();
1480 state.datagrams.clear();
1481 }
1482 Err(state) => {
1483 if state.final_sequence.is_some() || state.abort.is_some() {
1484 return;
1485 }
1486 tracing::warn!(
1487 track = %self.name,
1488 "track::Producer dropped without finish() or abort()"
1489 );
1490 }
1491 }
1492 }
1493}
1494
1495fn combined_subscription(subs: &Subscriptions, bound: Option<Duration>, waiter: &kio::Waiter) -> Option<Subscription> {
1501 let mut combined = None;
1502 for sub in subs.iter() {
1503 if sub.is_closed() {
1508 continue;
1509 }
1510 let _ = sub.poll_closed(waiter);
1515 if let Poll::Ready(Ok(sub)) = sub.poll(waiter, |sub| sub.poll_combined(&combined)) {
1516 combined = Some(sub);
1517 }
1518 }
1519 clamp_combined(combined, bound)
1520}
1521
1522fn snapshot_subscription(subs: &kio::Shared<Subscriptions>, bound: Option<Duration>) -> Option<Subscription> {
1524 let mut combined: Option<Subscription> = None;
1525 for sub in subs.read().iter() {
1526 if sub.is_closed() {
1528 continue;
1529 }
1530 if let Poll::Ready(merged) = sub.read().poll_combined(&combined) {
1531 combined = Some(merged);
1532 }
1533 }
1534 clamp_combined(combined, bound)
1535}
1536
1537fn clamp_combined(combined: Option<Subscription>, bound: Option<Duration>) -> Option<Subscription> {
1545 let mut combined = combined?;
1546 if let Some(bound) = bound {
1547 combined.latency_max = combined.latency_max.min(bound);
1548 }
1549 Some(combined)
1550}
1551
1552fn register_subscription(state: kio::Ref<'_, TrackState>, subscription: &kio::Producer<Subscription>) {
1556 if state.is_closed() {
1557 return;
1558 }
1559 let subs = state.subscriptions.clone();
1560 drop(state);
1561 subs.lock().push(subscription.consume());
1562}
1563
1564#[derive(Clone)]
1566pub(crate) struct TrackWeak {
1567 name: Arc<str>,
1568 state: kio::ProducerWeak<TrackState>,
1569}
1570
1571impl TrackWeak {
1572 pub fn consume(&self) -> Consumer {
1573 Consumer::plain(self.name.clone(), self.state.consume())
1574 }
1575
1576 pub(crate) fn name(&self) -> &Arc<str> {
1579 &self.name
1580 }
1581
1582 pub(crate) fn is_used(&self) -> bool {
1585 !self.state.is_closed() && self.state.is_used()
1586 }
1587
1588 pub(crate) fn poll_used(&self, waiter: &kio::Waiter) {
1591 let _ = self.state.poll_used(waiter);
1592 }
1593
1594 pub(crate) fn poll_unused(&self, waiter: &kio::Waiter) {
1597 let _ = self.state.poll_unused(waiter);
1598 }
1599}
1600
1601impl super::WeakEntry for TrackWeak {
1602 fn is_closed(&self) -> bool {
1603 self.state.is_closed()
1604 }
1605
1606 fn same_channel(&self, other: &Self) -> bool {
1607 self.state.same_channel(&other.state)
1608 }
1609}
1610
1611#[derive(Clone)]
1620pub struct Demand {
1621 name: Arc<str>,
1622 state: kio::ProducerWeak<TrackState>,
1623}
1624
1625impl Demand {
1626 pub fn name(&self) -> &str {
1628 &self.name
1629 }
1630
1631 pub async fn used(&self) -> Result<()> {
1633 self.state.used().await.map_err(|_| self.abort_reason())
1634 }
1635
1636 pub async fn unused(&self) -> Result<()> {
1638 self.state.unused().await.map_err(|_| self.abort_reason())
1639 }
1640
1641 pub async fn closed(&self) -> Error {
1643 self.state.closed().await;
1644 self.abort_reason()
1645 }
1646
1647 fn abort_reason(&self) -> Error {
1649 self.state.read().abort.clone().unwrap_or(Error::Dropped)
1650 }
1651}
1652
1653#[derive(Clone)]
1664pub struct Consumer {
1665 name: Arc<str>,
1666 inner: ConsumerKind,
1667 stats: stats::Scope,
1670}
1671
1672#[derive(Clone)]
1673enum ConsumerKind {
1674 Plain(kio::Consumer<TrackState>),
1675 Spliced(super::resume::Consumer),
1676}
1677
1678impl Consumer {
1679 fn plain(name: Arc<str>, state: kio::Consumer<TrackState>) -> Self {
1680 Self {
1681 name,
1682 inner: ConsumerKind::Plain(state),
1683 stats: stats::Scope::default(),
1684 }
1685 }
1686
1687 pub(crate) fn spliced(name: Arc<str>, resume: super::resume::Consumer) -> Self {
1689 Self {
1690 name,
1691 inner: ConsumerKind::Spliced(resume),
1692 stats: stats::Scope::default(),
1693 }
1694 }
1695
1696 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
1699 self.stats = scope;
1700 self
1701 }
1702
1703 pub fn name(&self) -> &str {
1705 &self.name
1706 }
1707
1708 pub fn subscribe(&self, subscription: impl Into<Option<Subscription>>) -> kio::Pending<Subscribing> {
1714 let subscription = kio::Producer::new(subscription.into().unwrap_or_default());
1715
1716 let inner = match &self.inner {
1717 ConsumerKind::Plain(state) => {
1718 register_subscription(state.read(), &subscription);
1721 SubscribingKind::Plain(state.clone())
1722 }
1723 ConsumerKind::Spliced(resume) => SubscribingKind::Spliced(resume.clone()),
1725 };
1726
1727 kio::Pending::new(Subscribing {
1728 name: self.name.clone(),
1729 inner,
1730 subscription,
1731 stats: self.stats.clone(),
1732 })
1733 }
1734
1735 #[cfg(test)]
1739 pub(crate) fn peek_group(&self, sequence: u64) -> Option<group::Consumer> {
1740 match &self.inner {
1741 ConsumerKind::Plain(state) => state.read().cached_group(sequence),
1742 ConsumerKind::Spliced(_) => None,
1745 }
1746 }
1747
1748 pub fn fetch_group(&self, sequence: u64, options: impl Into<Option<group::Fetch>>) -> kio::Pending<Fetching> {
1760 let options = options.into().unwrap_or_default();
1761
1762 self.stats.fetch();
1766
1767 let state = match &self.inner {
1768 ConsumerKind::Plain(state) => state,
1769 ConsumerKind::Spliced(resume) => {
1772 return kio::Pending::new(Fetching {
1773 inner: FetchingKind::Spliced(resume.fetch_group(sequence, options)),
1774 stats: self.stats.clone(),
1775 });
1776 }
1777 };
1778
1779 let mut result = None;
1780
1781 let (fetch, unresolved) = {
1785 let state = state.read();
1786 (state.fetch.clone(), state.poll_fetch_cached(sequence).is_pending())
1787 };
1788
1789 if unresolved {
1790 let mut fetch = fetch.lock();
1791 if let Some(pending) = fetch.join(&sequence) {
1792 pending.priority = pending.priority.max(options.priority);
1795 result = Some(pending.result.consume());
1796 } else {
1797 let producer = kio::Producer::<FetchOutcome>::default();
1801 let consumer = producer.consume();
1802 let attempt = PendingFetch {
1803 priority: options.priority,
1804 result: producer,
1805 };
1806 if fetch.insert(sequence, attempt).is_ok() {
1807 result = Some(consumer);
1808 }
1809 }
1810 }
1811
1812 kio::Pending::new(Fetching {
1813 inner: FetchingKind::Plain {
1814 state: state.clone(),
1815 fetch,
1816 sequence,
1817 result,
1818 },
1819 stats: self.stats.clone(),
1820 })
1821 }
1822
1823 pub fn info(&self) -> kio::Pending<Querying> {
1830 kio::Pending::new(Querying {
1831 inner: match &self.inner {
1832 ConsumerKind::Plain(state) => QueryingKind::Plain(state.clone()),
1833 ConsumerKind::Spliced(resume) => QueryingKind::Spliced(resume.clone()),
1834 },
1835 })
1836 }
1837
1838 pub fn latest(&self) -> Option<u64> {
1840 match &self.inner {
1841 ConsumerKind::Plain(state) => state.read().max_sequence,
1842 ConsumerKind::Spliced(resume) => resume.latest(),
1843 }
1844 }
1845
1846 pub(crate) fn poll_complete(&self, waiter: &kio::Waiter) -> Poll<Result<()>> {
1851 let ConsumerKind::Plain(state) = &self.inner else {
1852 return Poll::Pending;
1854 };
1855 match ready!(state.poll(waiter, |state| {
1856 if state.is_complete() {
1857 Poll::Ready(())
1858 } else {
1859 Poll::Pending
1860 }
1861 })) {
1862 Ok(_) => Poll::Ready(Ok(())),
1863 Err(closed) => Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped))),
1866 }
1867 }
1868}
1869
1870pub struct Subscribing {
1873 name: Arc<str>,
1874 inner: SubscribingKind,
1875 subscription: kio::Producer<Subscription>,
1876 stats: stats::Scope,
1877}
1878
1879enum SubscribingKind {
1880 Plain(kio::Consumer<TrackState>),
1881 Spliced(super::resume::Consumer),
1882}
1883
1884impl Subscribing {
1885 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<Subscriber>> {
1888 match &self.inner {
1889 SubscribingKind::Plain(state) => {
1890 let info = ready!(state.poll(waiter, |state| state.poll_info()))
1892 .map_err(|e| e.abort.clone().unwrap_or(Error::Dropped))??;
1893
1894 Poll::Ready(Ok(Subscriber {
1895 name: self.name.clone(),
1896 info,
1897 inner: SubscriberKind::Plain(PlainSubscriber {
1898 state: state.clone(),
1899 subscription: self.subscription.clone(),
1900 index: 0,
1901 datagram_index: 0,
1902 min_sequence: 0,
1903 next_sequence: 0,
1904 end_sequence: None,
1905 parked: BTreeMap::new(),
1906 }),
1907 stats: self.stats.clone(),
1908 _stats_sub: self.stats.subscribe(),
1909 }))
1910 }
1911 SubscribingKind::Spliced(resume) => {
1912 let info = ready!(resume.poll_info(waiter))?;
1915
1916 Poll::Ready(Ok(Subscriber {
1917 name: self.name.clone(),
1918 info,
1919 inner: SubscriberKind::Spliced(Box::new(resume.subscribe_shared(self.subscription.clone()))),
1920 stats: self.stats.clone(),
1921 _stats_sub: self.stats.subscribe(),
1922 }))
1923 }
1924 }
1925 }
1926
1927 pub fn update(&mut self, subscription: Subscription) -> Result<()> {
1932 let mut state = self.subscription.write().map_err(|_| Error::Closed)?;
1933 *state = subscription;
1934 Ok(())
1935 }
1936}
1937
1938impl kio::Pollable for Subscribing {
1939 type Output = Result<Subscriber>;
1940
1941 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
1942 self.poll_ok(waiter)
1943 }
1944}
1945
1946pub struct Querying {
1949 inner: QueryingKind,
1950}
1951
1952enum QueryingKind {
1953 Plain(kio::Consumer<TrackState>),
1954 Spliced(super::resume::Consumer),
1955}
1956
1957impl Querying {
1958 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<Info>> {
1960 match &self.inner {
1961 QueryingKind::Plain(state) => {
1962 let info = ready!(state.poll(waiter, |state| state.poll_info()))
1964 .map_err(|e| e.abort.clone().unwrap_or(Error::Dropped))??;
1965 Poll::Ready(Ok(info))
1966 }
1967 QueryingKind::Spliced(resume) => resume.poll_info(waiter),
1968 }
1969 }
1970}
1971
1972impl kio::Pollable for Querying {
1973 type Output = Result<Info>;
1974
1975 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
1976 self.poll_ok(waiter)
1977 }
1978}
1979
1980pub struct GroupRequest {
1989 state: kio::Producer<TrackState>,
1990 fetch: kio::Shared<FetchState>,
1992 sequence: u64,
1993 priority: u8,
1994 result: kio::Producer<FetchOutcome>,
1996 done: bool,
1997}
1998
1999impl GroupRequest {
2000 pub fn sequence(&self) -> u64 {
2002 self.sequence
2003 }
2004
2005 pub fn priority(&self) -> u8 {
2007 self.priority
2008 }
2009
2010 pub fn accept(mut self, info: impl Into<Option<Info>>) -> Result<group::Producer> {
2018 self.done = true;
2019 let res = TrackState::modify(&self.state)
2023 .and_then(|mut state| state.insert_group_request(self.sequence, info.into()));
2024 self.remove();
2025 res
2026 }
2027
2028 pub fn reject(mut self, err: Error) {
2030 self.done = true;
2031 self.remove();
2034 if let Ok(mut outcome) = self.result.write() {
2035 outcome.rejected = Some(err);
2036 }
2037 }
2038
2039 fn remove(&self) {
2042 self.fetch
2043 .lock()
2044 .remove_if(&self.sequence, |pending| pending.result.same_channel(&self.result));
2045 }
2046}
2047
2048impl Drop for GroupRequest {
2049 fn drop(&mut self) {
2050 if self.done {
2051 return;
2052 }
2053 self.remove();
2054 if let Ok(mut outcome) = self.result.write() {
2055 outcome.rejected = Some(Error::Dropped);
2056 }
2057 }
2058}
2059
2060pub struct Fetching {
2066 inner: FetchingKind,
2067 stats: stats::Scope,
2070}
2071
2072enum FetchingKind {
2073 Plain {
2074 state: kio::Consumer<TrackState>,
2075 fetch: kio::Shared<FetchState>,
2076 sequence: u64,
2077 result: Option<kio::Consumer<FetchOutcome>>,
2079 },
2080 Spliced(kio::Pending<super::resume::Fetching>),
2082}
2083
2084impl kio::Pollable for Fetching {
2085 type Output = Result<group::Consumer>;
2086
2087 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
2088 let (state, fetch, sequence, result) = match &self.inner {
2089 FetchingKind::Plain {
2090 state,
2091 fetch,
2092 sequence,
2093 result,
2094 } => (state, fetch, *sequence, result.as_ref()),
2095 FetchingKind::Spliced(spliced) => {
2096 return kio::Pollable::poll(&**spliced, waiter)
2099 .map(|res| res.map(|group| group.with_meter(self.stats.meter())));
2100 }
2101 };
2102
2103 match state.poll(waiter, |state| state.poll_fetch_cached(sequence)) {
2106 Poll::Ready(Ok(res)) => return Poll::Ready(res.map(|group| group.with_meter(self.stats.meter()))),
2107 Poll::Ready(Err(closed)) => {
2108 return Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped)));
2109 }
2110 Poll::Pending => {}
2111 }
2112
2113 let Some(result) = result else {
2115 return match fetch.poll(waiter, |fetch| match fetch.has_handlers() {
2118 false => Poll::Ready(()),
2119 true => Poll::Pending,
2120 }) {
2121 Poll::Ready(_guard) => Poll::Ready(Err(Error::NotFound)),
2122 Poll::Pending => Poll::Pending,
2123 };
2124 };
2125
2126 match result.poll(waiter, |outcome| match &outcome.rejected {
2129 Some(err) => Poll::Ready(err.clone()),
2130 None => Poll::Pending,
2131 }) {
2132 Poll::Ready(Ok(err)) => Poll::Ready(Err(err)),
2133 Poll::Ready(Err(_closed)) => Poll::Ready(Err(Error::NotFound)),
2134 Poll::Pending => Poll::Pending,
2135 }
2136 }
2137}
2138
2139pub struct Subscriber {
2162 name: Arc<str>,
2163 info: Info,
2164 inner: SubscriberKind,
2165 stats: stats::Scope,
2168 _stats_sub: stats::Subscription,
2171}
2172
2173enum SubscriberKind {
2174 Plain(PlainSubscriber),
2175 Spliced(Box<super::resume::Subscriber>),
2177}
2178
2179struct PlainSubscriber {
2181 state: kio::Consumer<TrackState>,
2182
2183 subscription: kio::Producer<Subscription>,
2184 index: usize,
2186 datagram_index: usize,
2188 min_sequence: u64,
2190 next_sequence: u64,
2193 end_sequence: Option<u64>,
2198 parked: BTreeMap<u64, group::Consumer>,
2203}
2204
2205impl PlainSubscriber {
2206 fn poll<F, R>(&self, waiter: &kio::Waiter, f: F) -> Poll<Result<R>>
2208 where
2209 F: Fn(&kio::Ref<'_, TrackState>) -> Poll<Result<R>>,
2210 {
2211 Poll::Ready(match ready!(self.state.poll(waiter, f)) {
2212 Ok(res) => res,
2213 Err(state) => Err(state.abort.clone().unwrap_or(Error::Dropped)),
2215 })
2216 }
2217
2218 fn poll_recv_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2219 let watch = |group: &group::Consumer| match group.poll_closed(waiter) {
2227 Poll::Pending => true,
2228 Poll::Ready(()) => !group.is_aborted(),
2229 };
2230
2231 let min_sequence = self.min_sequence;
2236 self.parked
2237 .retain(|sequence, group| *sequence >= min_sequence && watch(group));
2238
2239 if let Some(&sequence) = self.parked.keys().next()
2241 && self.end_sequence.is_none_or(|end| sequence <= end)
2242 {
2243 return Poll::Ready(Ok(self.parked.remove(&sequence)));
2244 }
2245
2246 loop {
2247 let Some((consumer, found_index)) =
2248 ready!(self.poll(waiter, |state| state.poll_recv_group(self.index, self.min_sequence))?)
2249 else {
2250 if self.parked.is_empty() {
2253 return Poll::Ready(Ok(None));
2254 }
2255 return Poll::Pending;
2256 };
2257 self.index = found_index + 1;
2258
2259 if self.end_sequence.is_some_and(|end| consumer.sequence > end) {
2262 if watch(&consumer) {
2266 self.parked.insert(consumer.sequence, consumer);
2267 }
2268 continue;
2269 }
2270 return Poll::Ready(Ok(Some(consumer)));
2271 }
2272 }
2273
2274 fn poll_recv_datagram(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Datagram>>> {
2275 let Some((datagram, found_index)) =
2276 ready!(self.poll(waiter, |state| state.poll_recv_datagram(self.datagram_index))?)
2277 else {
2278 return Poll::Ready(Ok(None));
2279 };
2280
2281 self.datagram_index = found_index + 1;
2282 self.next_sequence = self.next_sequence.max(datagram.sequence.saturating_add(1));
2283 Poll::Ready(Ok(Some(datagram)))
2284 }
2285
2286 fn poll_next_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2287 let floor = self.next_sequence.max(self.min_sequence);
2288 let Some(group) = ready!(self.poll(waiter, |state| state.poll_next_in_range(floor, self.end_sequence))?) else {
2289 return Poll::Ready(Ok(None));
2290 };
2291 self.next_sequence = group.sequence.saturating_add(1);
2292 Poll::Ready(Ok(Some(group)))
2293 }
2294
2295 fn poll_read_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Frame>>> {
2296 let lower = self.min_sequence.max(self.next_sequence);
2297 let Some((frame, found_index, sequence)) =
2298 ready!(self.poll(waiter, |state| { state.poll_read_frame(self.index, lower, waiter) })?)
2299 else {
2300 return Poll::Ready(Ok(None));
2301 };
2302
2303 self.index = found_index + 1;
2304 self.next_sequence = sequence.saturating_add(1);
2305 Poll::Ready(Ok(Some(frame)))
2306 }
2307}
2308
2309#[derive(Clone)]
2315pub struct SubscriberControl {
2316 subscription: kio::Producer<Subscription>,
2317}
2318
2319impl SubscriberControl {
2320 pub fn subscription(&self) -> Subscription {
2322 self.subscription.read().clone()
2323 }
2324
2325 pub fn update(&self, subscription: Subscription) -> Result<()> {
2330 let mut state = self.subscription.write().map_err(|_| Error::Closed)?;
2331 *state = subscription;
2332 Ok(())
2333 }
2334}
2335
2336impl Subscriber {
2337 pub fn info(&self) -> &Info {
2342 &self.info
2343 }
2344
2345 pub fn name(&self) -> &str {
2347 &self.name
2348 }
2349
2350 pub fn control(&self) -> SubscriberControl {
2352 SubscriberControl {
2353 subscription: match &self.inner {
2354 SubscriberKind::Plain(plain) => plain.subscription.clone(),
2355 SubscriberKind::Spliced(spliced) => spliced.prefs(),
2356 },
2357 }
2358 }
2359
2360 pub fn poll_recv_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2377 let meter = self.stats.meter();
2378 let res = match &mut self.inner {
2379 SubscriberKind::Plain(plain) => plain.poll_recv_group(waiter),
2380 SubscriberKind::Spliced(spliced) => spliced.poll_recv_group(waiter),
2381 };
2382 res.map(|res| res.map(|group| group.map(|group| group.with_meter(meter))))
2383 }
2384
2385 pub async fn recv_group(&mut self) -> Result<Option<group::Consumer>> {
2392 kio::wait(|waiter| self.poll_recv_group(waiter)).await
2393 }
2394
2395 pub fn poll_recv_datagram(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Datagram>>> {
2406 let meter = self.stats.meter();
2407 let res = match &mut self.inner {
2408 SubscriberKind::Plain(plain) => plain.poll_recv_datagram(waiter),
2409 SubscriberKind::Spliced(spliced) => spliced.poll_recv_datagram(waiter),
2410 };
2411 if let Poll::Ready(Ok(Some(datagram))) = &res {
2414 meter.datagram(datagram.payload.len() as u64);
2415 }
2416 res
2417 }
2418
2419 pub async fn recv_datagram(&mut self) -> Result<Option<Datagram>> {
2426 kio::wait(|waiter| self.poll_recv_datagram(waiter)).await
2427 }
2428
2429 pub fn poll_next_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2438 let meter = self.stats.meter();
2439 let res = match &mut self.inner {
2440 SubscriberKind::Plain(plain) => plain.poll_next_group(waiter),
2441 SubscriberKind::Spliced(spliced) => spliced.poll_next_group(waiter),
2442 };
2443 res.map(|res| res.map(|group| group.map(|group| group.with_meter(meter))))
2444 }
2445
2446 pub async fn next_group(&mut self) -> Result<Option<group::Consumer>> {
2452 kio::wait(|waiter| self.poll_next_group(waiter)).await
2453 }
2454
2455 pub fn poll_read_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Frame>>> {
2459 let meter = self.stats.meter();
2460 let res = match &mut self.inner {
2461 SubscriberKind::Plain(plain) => plain.poll_read_frame(waiter),
2462 SubscriberKind::Spliced(spliced) => spliced.poll_read_frame(waiter),
2463 };
2464 if let Poll::Ready(Ok(Some(frame))) = &res {
2467 meter.group();
2468 meter.frames(1);
2469 meter.bytes(frame.payload.len() as u64);
2470 }
2471 res
2472 }
2473
2474 pub async fn read_frame(&mut self) -> Result<Option<frame::Frame>> {
2479 kio::wait(|waiter| self.poll_read_frame(waiter)).await
2480 }
2481
2482 pub fn is_clone(&self, other: &Self) -> bool {
2484 match (&self.inner, &other.inner) {
2485 (SubscriberKind::Plain(a), SubscriberKind::Plain(b)) => a.state.same_channel(&b.state),
2486 (SubscriberKind::Spliced(a), SubscriberKind::Spliced(b)) => a.is_clone(b),
2487 _ => false,
2488 }
2489 }
2490
2491 pub fn poll_finished(&mut self, waiter: &kio::Waiter) -> Poll<Result<u64>> {
2493 match &mut self.inner {
2494 SubscriberKind::Plain(plain) => plain.poll(waiter, |state| state.poll_finished()),
2495 SubscriberKind::Spliced(spliced) => spliced.poll_finished(waiter),
2496 }
2497 }
2498
2499 pub async fn finished(&mut self) -> Result<u64> {
2507 kio::wait(|waiter| self.poll_finished(waiter)).await
2508 }
2509
2510 pub fn start_at(&mut self, sequence: u64) {
2517 match &mut self.inner {
2518 SubscriberKind::Plain(plain) => plain.min_sequence = sequence,
2519 SubscriberKind::Spliced(spliced) => spliced.start_at(sequence),
2520 }
2521 }
2522
2523 pub fn end_at(&mut self, sequence: impl Into<Option<u64>>) {
2536 match &mut self.inner {
2537 SubscriberKind::Plain(plain) => plain.end_sequence = sequence.into(),
2538 SubscriberKind::Spliced(spliced) => spliced.end_at(sequence),
2539 }
2540 }
2541
2542 pub fn subscription(&self) -> Subscription {
2544 self.control().subscription()
2545 }
2546
2547 pub fn update(&mut self, subscription: Subscription) -> Result<()> {
2553 match &mut self.inner {
2554 SubscriberKind::Plain(plain) => {
2555 let mut state = plain.subscription.write().map_err(|_| Error::Closed)?;
2556 *state = subscription;
2557 }
2558 SubscriberKind::Spliced(spliced) => spliced.update(subscription),
2559 }
2560 Ok(())
2561 }
2562
2563 pub fn latest(&self) -> Option<u64> {
2565 match &self.inner {
2566 SubscriberKind::Plain(plain) => plain.state.read().max_sequence,
2567 SubscriberKind::Spliced(spliced) => spliced.latest(),
2568 }
2569 }
2570}
2571
2572pub struct Request {
2584 name: Arc<str>,
2585 broadcast: Arc<broadcast::Info>,
2587 state: kio::Producer<TrackState>,
2588
2589 prev_subscription: Option<Subscription>,
2591
2592 alive: Arc<Alive>,
2595
2596 _dynamic: Dynamic,
2601
2602 stats: stats::Scope,
2605}
2606
2607impl Request {
2608 pub(crate) fn new(broadcast: Arc<broadcast::Info>, name: impl Into<Arc<str>>) -> Self {
2609 let name = name.into();
2610 let state = TrackState::spawn(broadcast.clone());
2611 let alive = Alive::new(name.clone(), state.clone());
2612 let dynamic = Dynamic::new(name.clone(), state.clone(), alive.clone());
2613 Self {
2614 name,
2615 broadcast,
2616 state,
2617 prev_subscription: None,
2618 alive,
2619 _dynamic: dynamic,
2620 stats: stats::Scope::default(),
2621 }
2622 }
2623
2624 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
2627 self.stats = scope;
2628 self
2629 }
2630
2631 pub fn name(&self) -> &str {
2633 &self.name
2634 }
2635
2636 pub fn consume(&self) -> Consumer {
2638 Consumer::plain(self.name.clone(), self.state.consume())
2639 }
2640
2641 pub fn dynamic(&self) -> Dynamic {
2645 Dynamic::new(self.name.clone(), self.state.clone(), self.alive.clone())
2646 }
2647
2648 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
2651 self.state.poll_unused(waiter).map(|_| ())
2652 }
2653
2654 pub fn accept(self, info: impl Into<Option<Info>>) -> Producer {
2661 if let Ok(mut state) = self.state.write() {
2664 state.install(info.into().unwrap_or_default());
2665 }
2666 self.alive.publish(Some(&self.stats));
2669 Producer {
2670 name: self.name,
2671 broadcast: self.broadcast,
2672 state: self.state,
2673 prev_subscription: None,
2674 alive: self.alive,
2675 stats: self.stats,
2676 }
2677 }
2678
2679 pub fn reject(self, err: Error) {
2681 if let Ok(mut state) = self.state.write() {
2682 state.abort = Some(err);
2683 }
2684 }
2685
2686 pub fn subscription(&self) -> Option<Subscription> {
2689 let state = self.state.read();
2690 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
2691 drop(state);
2692 snapshot_subscription(&subs, bound)
2693 }
2694
2695 pub async fn subscription_changed(&mut self) -> Option<Subscription> {
2698 kio::wait(|waiter| self.poll_subscription_changed(waiter)).await
2699 }
2700
2701 pub fn poll_subscription_changed(&mut self, waiter: &kio::Waiter) -> Poll<Option<Subscription>> {
2703 let state = self.state.read();
2704 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
2705 drop(state);
2706
2707 let prev = &self.prev_subscription;
2708 let mut combined = None;
2709 let mut guard = ready!(subs.poll(waiter, |subs| {
2710 let next = combined_subscription(subs, bound, waiter);
2711 if &next == prev {
2712 Poll::Pending
2713 } else {
2714 combined = next;
2715 Poll::Ready(())
2716 }
2717 }));
2718 guard.retain(|sub| !sub.is_closed());
2720 drop(guard);
2721 self.prev_subscription = combined.clone();
2722 Poll::Ready(combined)
2723 }
2724
2725 pub(super) fn weak(&self) -> TrackWeak {
2726 TrackWeak {
2727 name: self.name.clone(),
2728 state: self.state.weak(),
2729 }
2730 }
2731}
2732
2733#[cfg(test)]
2734use futures::FutureExt;
2735
2736#[cfg(test)]
2737#[allow(missing_docs)] impl Subscriber {
2739 pub fn assert_group(&mut self) -> group::Consumer {
2740 self.recv_group()
2741 .now_or_never()
2742 .expect("group would have blocked")
2743 .expect("would have errored")
2744 .expect("track was closed")
2745 }
2746
2747 pub fn assert_no_group(&mut self) {
2748 assert!(
2749 self.recv_group().now_or_never().is_none(),
2750 "recv_group would not have blocked"
2751 );
2752 }
2753
2754 pub fn assert_not_closed(&mut self) {
2755 assert!(self.finished().now_or_never().is_none(), "should not be closed");
2756 }
2757
2758 pub fn assert_closed(&mut self) {
2759 assert!(self.finished().now_or_never().is_some(), "should be closed");
2760 }
2761
2762 pub fn assert_error(&mut self) {
2764 assert!(
2765 self.finished().now_or_never().expect("should not block").is_err(),
2766 "should be error"
2767 );
2768 }
2769
2770 pub fn assert_is_clone(&self, other: &Self) {
2771 assert!(self.is_clone(other), "should be clone");
2772 }
2773
2774 pub fn assert_not_clone(&self, other: &Self) {
2775 assert!(!self.is_clone(other), "should not be clone");
2776 }
2777}
2778
2779#[cfg(test)]
2780mod test {
2781 use super::*;
2782 use crate::model::test_tracing::count_drop_warnings;
2783
2784 fn track_producer(name: impl Into<Arc<str>>, info: impl Into<Option<Info>>) -> Producer {
2787 Producer::new(Arc::new(broadcast::Info::default()), name, info)
2788 }
2789
2790 fn live_groups(state: &TrackState) -> usize {
2792 state.lookup.len()
2793 }
2794
2795 fn first_live_sequence(state: &TrackState) -> u64 {
2797 state
2798 .arrival
2799 .iter()
2800 .find(|(sequence, stamp)| state.lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp))
2801 .map(|(sequence, _)| *sequence)
2802 .unwrap()
2803 }
2804
2805 fn recv_datagram(dg: &mut Subscriber) -> Datagram {
2807 dg.recv_datagram()
2808 .now_or_never()
2809 .expect("datagram would have blocked")
2810 .expect("would have errored")
2811 .expect("track was closed")
2812 }
2813
2814 #[tokio::test]
2815 async fn append_datagram_shares_group_sequence() {
2816 let mut producer = track_producer("test", None);
2817 let ts = Timestamp::from_millis(10).unwrap();
2818
2819 assert_eq!(producer.append_group().unwrap().sequence, 0);
2821 assert_eq!(producer.append_datagram(ts, &b"a"[..]).unwrap(), 1);
2822 assert_eq!(producer.append_group().unwrap().sequence, 2);
2823 assert_eq!(producer.append_datagram(ts, &b"b"[..]).unwrap(), 3);
2824 assert_eq!(producer.latest(), Some(3));
2825 }
2826
2827 #[tokio::test]
2828 async fn append_datagram_roundtrip() {
2829 let mut producer = track_producer("test", None);
2830 let mut dg = producer.subscribe(None);
2831
2832 let ts = Timestamp::from_millis(42).unwrap();
2833 let seq = producer.append_datagram(ts, &b"hello"[..]).unwrap();
2834
2835 let got = recv_datagram(&mut dg);
2836 assert_eq!(got.sequence, seq);
2837 assert_eq!(got.timestamp, ts);
2838 assert_eq!(&got.payload[..], b"hello");
2839 }
2840
2841 #[tokio::test]
2842 async fn write_datagram_preserves_sequence() {
2843 let mut producer = track_producer("test", None);
2844 let mut dg = producer.subscribe(None);
2845
2846 let ts = Timestamp::from_millis(5).unwrap();
2847 producer
2849 .write_datagram(Datagram {
2850 sequence: 100,
2851 timestamp: ts,
2852 payload: bytes::Bytes::from_static(b"x"),
2853 })
2854 .unwrap();
2855
2856 assert_eq!(recv_datagram(&mut dg).sequence, 100);
2857 assert_eq!(producer.append_group().unwrap().sequence, 101);
2859 }
2860
2861 #[tokio::test]
2862 async fn recv_datagram_advances_ordered_group_cursor() {
2863 let mut producer = track_producer("test", None);
2864 let mut subscriber = producer.subscribe(None);
2865 let ts = Timestamp::from_millis(5).unwrap();
2866
2867 producer
2868 .write_datagram(Datagram {
2869 sequence: 5,
2870 timestamp: ts,
2871 payload: bytes::Bytes::from_static(b"x"),
2872 })
2873 .unwrap();
2874 assert_eq!(recv_datagram(&mut subscriber).sequence, 5);
2875
2876 producer.create_group(group::Info { sequence: 3 }).unwrap();
2877 producer.create_group(group::Info { sequence: 6 }).unwrap();
2878
2879 let group = subscriber
2880 .next_group()
2881 .now_or_never()
2882 .expect("group would have blocked")
2883 .expect("would have errored")
2884 .expect("track was closed");
2885 assert_eq!(group.sequence, 6);
2886 }
2887
2888 #[tokio::test]
2889 async fn datagram_normalized_to_track_timescale() {
2890 let info = Info::default().with_timescale(Timescale::MICRO);
2891 let mut producer = track_producer("test", info);
2892 let mut dg = producer.subscribe(None);
2893
2894 producer
2896 .append_datagram(Timestamp::from_millis(2).unwrap(), &b"z"[..])
2897 .unwrap();
2898 let got = recv_datagram(&mut dg);
2899 assert_eq!(got.timestamp.scale(), Timescale::MICRO);
2900 assert_eq!(got.timestamp.value(), 2_000);
2901 }
2902
2903 #[tokio::test]
2904 async fn datagram_rejects_oversized() {
2905 let mut producer = track_producer("test", None);
2906 let big = bytes::Bytes::from(vec![0u8; crate::model::datagram::MAX_DATAGRAM_PAYLOAD + 1]);
2907 let ts = Timestamp::from_millis(0).unwrap();
2908 assert!(matches!(
2909 producer.append_datagram(ts, big.clone()),
2910 Err(Error::FrameTooLarge)
2911 ));
2912 assert!(matches!(
2913 producer.write_datagram(Datagram {
2914 sequence: 0,
2915 timestamp: ts,
2916 payload: big,
2917 }),
2918 Err(Error::FrameTooLarge)
2919 ));
2920 }
2921
2922 #[tokio::test]
2923 async fn datagram_fanout_to_subscribers() {
2924 let mut producer = track_producer("test", None);
2925 let mut a = producer.subscribe(None);
2927 let mut b = producer.subscribe(None);
2928 let ts = Timestamp::from_millis(1).unwrap();
2929
2930 producer.append_datagram(ts, &b"first"[..]).unwrap();
2931 producer.append_datagram(ts, &b"second"[..]).unwrap();
2932
2933 assert_eq!(&recv_datagram(&mut a).payload[..], b"first");
2935 assert_eq!(&recv_datagram(&mut a).payload[..], b"second");
2936 assert_eq!(&recv_datagram(&mut b).payload[..], b"first");
2937 assert_eq!(&recv_datagram(&mut b).payload[..], b"second");
2938 }
2939
2940 #[tokio::test]
2941 async fn datagram_evicts_stale() {
2942 tokio::time::pause();
2943
2944 let mut producer = track_producer("test", None);
2945 let mut dg = producer.subscribe(None);
2946 let ts = Timestamp::from_millis(0).unwrap();
2947
2948 producer.append_datagram(ts, &b"old"[..]).unwrap(); tokio::time::advance(MAX_DATAGRAM_AGE + Duration::from_millis(10)).await;
2952 producer.append_datagram(ts, &b"new"[..]).unwrap(); let got = recv_datagram(&mut dg);
2956 assert_eq!(got.sequence, 1);
2957 assert_eq!(&got.payload[..], b"new");
2958 }
2959
2960 #[tokio::test]
2961 async fn datagram_recv_pends_until_written() {
2962 let mut producer = track_producer("test", None);
2963 let mut dg = producer.subscribe(None);
2964
2965 assert!(
2966 dg.recv_datagram().now_or_never().is_none(),
2967 "should block with no datagrams"
2968 );
2969
2970 producer
2971 .append_datagram(Timestamp::from_millis(0).unwrap(), &b"go"[..])
2972 .unwrap();
2973 assert_eq!(&recv_datagram(&mut dg).payload[..], b"go");
2974 }
2975
2976 #[tokio::test]
2980 async fn datagram_wire_roundtrip_between_tracks() {
2981 use crate::coding::{Decode, Encode};
2982 use crate::lite;
2983
2984 let version = lite::Version::Lite05;
2985
2986 let mut origin = track_producer("test", None);
2988 let mut origin_dg = origin.subscribe(None);
2989 let ts = Timestamp::from_millis(7).unwrap();
2990 let seq = origin.append_datagram(ts, &b"payload"[..]).unwrap();
2991
2992 let d = recv_datagram(&mut origin_dg);
2993 let body = lite::Datagram {
2994 subscribe: 5,
2995 sequence: d.sequence,
2996 timestamp: d.timestamp.value(),
2997 payload: d.payload.clone(),
2998 }
2999 .encode_bytes(version)
3000 .unwrap();
3001
3002 let mut slice = &body[..];
3004 let wire = lite::Datagram::decode(&mut slice, version).unwrap();
3005 let mut downstream = track_producer("test", None);
3006 let mut downstream_dg = downstream.subscribe(None);
3007 downstream
3008 .write_datagram(Datagram {
3009 sequence: wire.sequence,
3010 timestamp: Timestamp::new(wire.timestamp, Timescale::MILLI).unwrap(),
3011 payload: wire.payload,
3012 })
3013 .unwrap();
3014
3015 let got = recv_datagram(&mut downstream_dg);
3016 assert_eq!(got.sequence, seq);
3017 assert_eq!(got.timestamp, ts);
3018 assert_eq!(&got.payload[..], b"payload");
3019 }
3020
3021 #[tokio::test]
3022 async fn evict_expired_groups() {
3023 tokio::time::pause();
3024
3025 let mut producer = track_producer("test", None);
3026
3027 producer.append_group().unwrap(); producer.append_group().unwrap(); producer.append_group().unwrap(); {
3033 let state = producer.state.read();
3034 assert_eq!(live_groups(&state), 3);
3035 assert_eq!(state.offset, 0);
3036 }
3037
3038 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3040
3041 producer.append_group().unwrap(); {
3048 let state = producer.state.read();
3049 assert_eq!(live_groups(&state), 1);
3050 assert_eq!(first_live_sequence(&state), 3);
3051 assert_eq!(state.offset, 3);
3052 assert!(!state.lookup.contains_key(&0));
3053 assert!(!state.lookup.contains_key(&1));
3054 assert!(!state.lookup.contains_key(&2));
3055 assert!(state.lookup.contains_key(&3));
3056 }
3057 }
3058
3059 #[tokio::test]
3063 async fn aging_out_a_finished_group_keeps_the_clean_end() {
3064 tokio::time::pause();
3065
3066 let mut producer = track_producer("test", None);
3067 let mut group = producer.create_group(group::Info { sequence: 0 }).unwrap();
3068 let mut consumer = group.consume();
3069
3070 group
3071 .write_frame(Timestamp::from_millis(0).unwrap(), b"hello".as_slice())
3072 .unwrap();
3073 assert_eq!(consumer.next_frame().await.unwrap().unwrap().size, 5);
3074
3075 tokio::time::advance(DEFAULT_LATENCY_MAX * 12).await;
3077 group.finish().unwrap();
3078 let _next = producer.create_group(group::Info { sequence: 1 }).unwrap();
3079
3080 assert!(consumer.next_frame().await.unwrap().is_none());
3081 }
3082
3083 #[tokio::test]
3084 async fn evict_keeps_max_sequence() {
3085 tokio::time::pause();
3086
3087 let mut producer = track_producer("test", None);
3088 producer.append_group().unwrap(); tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3092
3093 producer.append_group().unwrap(); {
3097 let state = producer.state.read();
3098 assert_eq!(live_groups(&state), 1);
3099 assert_eq!(first_live_sequence(&state), 1);
3100 assert_eq!(state.offset, 1);
3101 }
3102 }
3103
3104 #[tokio::test]
3105 async fn no_eviction_when_fresh() {
3106 tokio::time::pause();
3107
3108 let mut producer = track_producer("test", None);
3109 producer.append_group().unwrap(); producer.append_group().unwrap(); producer.append_group().unwrap(); {
3114 let state = producer.state.read();
3115 assert_eq!(live_groups(&state), 3);
3116 assert_eq!(state.offset, 0);
3117 }
3118 }
3119
3120 #[tokio::test]
3121 async fn consumer_skips_evicted_groups() {
3122 tokio::time::pause();
3123
3124 let mut producer = track_producer("test", None);
3125 producer.append_group().unwrap(); let mut consumer = producer.subscribe(None);
3128
3129 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3130 producer.append_group().unwrap(); let group = consumer.assert_group();
3134 assert_eq!(group.sequence, 1);
3135 }
3136
3137 #[tokio::test]
3138 async fn cache_age_controls_eviction() {
3139 tokio::time::pause();
3140
3141 let mut producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(1)));
3143 producer.append_group().unwrap(); tokio::time::advance(Duration::from_secs(2)).await;
3147 producer.append_group().unwrap(); let state = producer.state.read();
3151 assert_eq!(live_groups(&state), 1);
3152 assert_eq!(first_live_sequence(&state), 1);
3153 }
3154
3155 #[test]
3156 fn latency_max_clamped_to_cache() {
3157 let producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(2)));
3158
3159 let mut subscriber = producer.subscribe(Subscription::default().with_latency_max(Duration::from_secs(10)));
3163 assert_eq!(subscriber.subscription().latency_max, Duration::from_secs(10));
3164 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3165
3166 subscriber
3168 .update(Subscription::default().with_latency_max(Duration::from_millis(500)))
3169 .unwrap();
3170 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_millis(500));
3171
3172 subscriber
3173 .update(Subscription::default().with_latency_max(Duration::ZERO))
3174 .unwrap();
3175 assert_eq!(producer.subscription().unwrap().latency_max, Duration::ZERO);
3176 }
3177
3178 fn track_producer_capped(name: impl Into<Arc<str>>, info: Info, cap: Duration) -> Producer {
3181 let origin = crate::origin::Info::default().with_cache_duration(cap);
3182 Producer::new(Arc::new(broadcast::Info { origin }), name, info)
3183 }
3184
3185 #[test]
3186 fn origin_cache_duration_clamps_latency_max() {
3187 let capped = track_producer_capped(
3190 "test",
3191 Info::default().with_latency_max(Duration::from_secs(60)),
3192 Duration::from_secs(1),
3193 );
3194 assert_eq!(capped.state.read().latency_bound(), Some(Duration::from_secs(1)));
3195
3196 let under = track_producer_capped(
3197 "test",
3198 Info::default().with_latency_max(Duration::from_millis(500)),
3199 Duration::from_secs(1),
3200 );
3201 assert_eq!(under.state.read().latency_bound(), Some(Duration::from_millis(500)));
3202 }
3203
3204 #[tokio::test]
3205 async fn origin_cache_duration_caps_eviction() {
3206 tokio::time::pause();
3207
3208 let mut producer = track_producer_capped(
3210 "test",
3211 Info::default().with_latency_max(Duration::from_secs(60)),
3212 Duration::from_secs(1),
3213 );
3214 producer.append_group().unwrap(); tokio::time::advance(Duration::from_secs(2)).await;
3218 producer.append_group().unwrap(); let state = producer.state.read();
3222 assert_eq!(live_groups(&state), 1);
3223 assert_eq!(first_live_sequence(&state), 1);
3224 }
3225
3226 #[test]
3227 fn latency_max_clamped_via_every_update_path() {
3228 let producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(2)));
3229 let over = Subscription::default().with_latency_max(Duration::from_secs(10));
3230
3231 let mut subscriber = producer.subscribe(over.clone());
3234 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3235
3236 subscriber.control().update(over.clone()).unwrap();
3237 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3238
3239 subscriber.update(over).unwrap();
3240 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3241 }
3242
3243 #[test]
3244 fn latency_max_aggregate_clamps_the_max_across_subscribers() {
3245 let producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(2)));
3246
3247 let _a = producer.subscribe(Subscription::default().with_latency_max(Duration::from_millis(500)));
3250 let _b = producer.subscribe(Subscription::default().with_latency_max(Duration::from_secs(10)));
3251
3252 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3253 }
3254
3255 #[test]
3256 fn subscriber_control_updates_while_read_future_is_pending() {
3257 let producer = track_producer("test", None);
3258 let mut subscriber = producer.subscribe(None);
3259 let control = subscriber.control();
3260
3261 let mut recv = Box::pin(subscriber.recv_group());
3262 assert!(recv.as_mut().now_or_never().is_none());
3263
3264 control
3265 .update(Subscription::default().with_priority(7).with_ordered(false))
3266 .unwrap();
3267
3268 let aggregate = producer.subscription().expect("expected an active subscription");
3269 assert_eq!(aggregate.priority, 7);
3270 assert!(!aggregate.ordered);
3271 }
3272
3273 #[test]
3274 fn dropped_subscriber_leaves_no_ghost_in_aggregate() {
3275 let mut producer = track_producer("test", None);
3280 let a = producer.subscribe(Subscription::default().with_priority(5));
3281
3282 let waiter = kio::Waiter::noop();
3284 assert!(
3285 matches!(producer.poll_subscription_changed(&waiter), Poll::Ready(Ok(Some(_)))),
3286 "one live subscriber should aggregate to Some",
3287 );
3288
3289 drop(a);
3291
3292 assert!(
3294 matches!(producer.poll_subscription_changed(&waiter), Poll::Ready(Ok(None))),
3295 "a dropped subscriber must not linger in the aggregate",
3296 );
3297
3298 assert!(
3300 producer.subscription().is_none(),
3301 "snapshot must exclude a dropped subscriber",
3302 );
3303 }
3304
3305 #[test]
3306 fn dropped_subscriber_wakes_the_aggregate() {
3307 use std::sync::atomic::{AtomicBool, Ordering};
3314
3315 let mut producer = track_producer("test", None);
3316 let a = producer.subscribe(Subscription::default().with_priority(5));
3317
3318 let woken = Arc::new(AtomicBool::new(false));
3319 let waiter = kio::Waiter::new(futures::task::waker(Arc::new(FlagWake(woken.clone()))));
3320
3321 assert!(matches!(
3323 producer.poll_subscription_changed(&waiter),
3324 Poll::Ready(Ok(Some(_)))
3325 ));
3326 assert!(
3327 producer.poll_subscription_changed(&waiter).is_pending(),
3328 "the aggregate is unchanged, so this poll must park",
3329 );
3330 assert!(!woken.load(Ordering::SeqCst), "nothing happened yet");
3331
3332 drop(a);
3333 assert!(
3334 woken.load(Ordering::SeqCst),
3335 "the last subscriber leaving must wake the aggregate watcher",
3336 );
3337 }
3338
3339 struct FlagWake(Arc<std::sync::atomic::AtomicBool>);
3341
3342 impl futures::task::ArcWake for FlagWake {
3343 fn wake_by_ref(arc_self: &Arc<Self>) {
3344 arc_self.0.store(true, std::sync::atomic::Ordering::SeqCst);
3345 }
3346 }
3347
3348 #[tokio::test]
3349 async fn out_of_order_max_sequence_at_front() {
3350 tokio::time::pause();
3351
3352 let mut producer = track_producer("test", None);
3353
3354 producer.create_group(group::Info { sequence: 5 }).unwrap();
3356 producer.create_group(group::Info { sequence: 3 }).unwrap();
3357 producer.create_group(group::Info { sequence: 4 }).unwrap();
3358
3359 {
3361 let state = producer.state.read();
3362 assert_eq!(state.max_sequence, Some(5));
3363 }
3364
3365 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3367
3368 producer.append_group().unwrap(); {
3374 let state = producer.state.read();
3375 assert_eq!(live_groups(&state), 1);
3376 assert_eq!(first_live_sequence(&state), 6);
3377 assert!(!state.lookup.contains_key(&3));
3378 assert!(!state.lookup.contains_key(&4));
3379 assert!(!state.lookup.contains_key(&5));
3380 assert!(state.lookup.contains_key(&6));
3381 }
3382 }
3383
3384 #[tokio::test]
3385 async fn max_sequence_at_front_blocks_trim() {
3386 tokio::time::pause();
3387
3388 let mut producer = track_producer("test", None);
3389
3390 producer.create_group(group::Info { sequence: 5 }).unwrap();
3392
3393 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3394
3395 producer.create_group(group::Info { sequence: 3 }).unwrap();
3397
3398 {
3401 let state = producer.state.read();
3402 assert_eq!(live_groups(&state), 2);
3403 assert_eq!(state.offset, 0);
3404 }
3405
3406 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3408
3409 producer.create_group(group::Info { sequence: 2 }).unwrap();
3411
3412 {
3417 let state = producer.state.read();
3418 assert_eq!(live_groups(&state), 2);
3419 assert_eq!(state.offset, 0);
3420 assert!(state.lookup.contains_key(&5));
3421 assert!(!state.lookup.contains_key(&3));
3422 assert!(state.lookup.contains_key(&2));
3423 }
3424
3425 let mut consumer = producer.subscribe(None);
3427 let group = consumer.assert_group();
3428 assert_eq!(group.sequence, 5);
3430 }
3431
3432 #[tokio::test]
3433 async fn abort_clears_cached_groups() {
3434 let mut producer = track_producer("test", None);
3435 producer.append_group().unwrap();
3436 producer.append_group().unwrap();
3437
3438 let mut consumer = producer.subscribe(None);
3440 assert_eq!(live_groups(&producer.state.read()), 2);
3441
3442 producer.clone().abort(Error::Cancel).unwrap();
3443
3444 {
3445 let state = producer.state.read();
3446 assert!(state.lookup.is_empty(), "cached groups should be dropped on abort");
3447 assert!(state.arrival.is_empty());
3448 assert!(state.evict.is_empty());
3449 }
3450
3451 let result = consumer.recv_group().now_or_never().expect("should not block");
3453 assert!(matches!(result, Err(Error::Cancel)));
3454 }
3455
3456 #[tokio::test]
3457 async fn drop_unfinished_clears_cached_groups() {
3458 let producer = track_producer("test", None);
3459 let mut writer = producer.clone();
3460 writer.append_group().unwrap();
3461
3462 let mut consumer = producer.subscribe(None);
3464 assert_eq!(live_groups(&producer.state.read()), 1);
3465
3466 drop(writer);
3468 drop(producer);
3469
3470 let result = consumer.recv_group().now_or_never().expect("should not block");
3471 assert!(matches!(result, Err(Error::Dropped)));
3472 }
3473
3474 #[tokio::test]
3475 async fn drop_after_abort_does_not_warn() {
3476 let warns = count_drop_warnings("track::Producer dropped without finish", || {
3479 let producer = track_producer("test", None);
3480 let keep = producer.clone();
3481 let mut writer = producer.clone();
3482 let mut group = writer.append_group().unwrap();
3483 group.finish().unwrap();
3484 let _consumer = producer.subscribe(None);
3485 writer.abort(Error::Cancel).unwrap();
3486 drop(keep);
3487 });
3488 assert_eq!(warns, 0, "abort-then-drop must not emit unfinished-producer WARN");
3489 }
3490
3491 #[tokio::test]
3492 async fn drop_unfinished_warns() {
3493 let warns = count_drop_warnings("track::Producer dropped without finish", || {
3494 let producer = track_producer("test", None);
3495 let mut writer = producer.clone();
3496 writer.append_group().unwrap();
3497 let _consumer = producer.subscribe(None);
3498 drop(writer);
3499 drop(producer);
3500 });
3501 assert!(warns >= 1, "unfinished drop must emit unfinished-producer WARN");
3502 }
3503
3504 #[tokio::test]
3505 async fn drop_finished_keeps_cached_groups() {
3506 let mut producer = track_producer("test", None);
3507 producer.append_group().unwrap();
3508 producer.finish().unwrap();
3509
3510 let mut consumer = producer.subscribe(None);
3511 drop(producer);
3512
3513 assert_eq!(consumer.assert_group().sequence, 0);
3515 let done = consumer.recv_group().now_or_never().expect("should not block").unwrap();
3516 assert!(done.is_none(), "consumer should drain then see clean finish");
3517 }
3518
3519 #[test]
3520 fn append_finish_cannot_be_rewritten() {
3521 let mut producer = track_producer("test", None);
3522
3523 assert!(producer.finish().is_ok());
3525 assert!(producer.finish().is_err());
3526 assert!(producer.append_group().is_err());
3527 }
3528
3529 #[test]
3530 fn finish_after_groups() {
3531 let mut producer = track_producer("test", None);
3532
3533 producer.append_group().unwrap();
3534 assert!(producer.finish().is_ok());
3535 assert!(producer.finish().is_err());
3536 assert!(producer.append_group().is_err());
3537 }
3538
3539 #[test]
3540 fn finish_at_rejects_a_boundary_at_or_below_the_live_edge() {
3541 let mut producer = track_producer("test", None);
3542 producer.create_group(group::Info { sequence: 5 }).unwrap();
3543
3544 assert!(producer.finish_at(4).is_err());
3547 assert!(producer.finish_at(5).is_err());
3548 assert!(producer.finish_at(6).is_ok());
3549
3550 {
3551 let state = producer.state.read();
3552 assert_eq!(state.final_sequence, Some(6));
3553 }
3554
3555 assert!(producer.finish_at(6).is_err());
3557 assert!(producer.create_group(group::Info { sequence: 4 }).is_ok());
3558 assert!(producer.create_group(group::Info { sequence: 6 }).is_err());
3559 }
3560
3561 #[test]
3562 fn final_sequence_reports_the_declared_boundary() {
3563 let mut producer = track_producer("test", None);
3564 assert_eq!(producer.final_sequence(), None);
3565
3566 producer.create_group(group::Info { sequence: 5 }).unwrap();
3567 assert_eq!(producer.final_sequence(), None, "a group does not declare a boundary");
3568
3569 producer.finish_at(9).unwrap();
3570 assert_eq!(producer.final_sequence(), Some(9));
3571
3572 assert!(producer.finish().is_err());
3574 }
3575
3576 #[test]
3577 fn final_sequence_reports_the_live_edge_after_finish() {
3578 let mut producer = track_producer("test", None);
3579 producer.create_group(group::Info { sequence: 5 }).unwrap();
3580 producer.finish().unwrap();
3581 assert_eq!(producer.final_sequence(), Some(6));
3582 }
3583
3584 #[tokio::test]
3585 async fn finish_at_declares_a_future_boundary() {
3586 let mut producer = track_producer("test", None);
3587 producer.create_group(group::Info { sequence: 5 }).unwrap();
3588
3589 producer.finish_at(7).unwrap();
3591
3592 let mut consumer = producer.subscribe(None);
3593 assert_eq!(consumer.assert_group().sequence, 5);
3594
3595 let boundary = consumer
3598 .finished()
3599 .now_or_never()
3600 .expect("boundary is known immediately")
3601 .expect("would have errored");
3602 assert_eq!(boundary, 7);
3603 assert!(
3604 consumer.recv_group().now_or_never().is_none(),
3605 "should wait for the outstanding group"
3606 );
3607
3608 producer.create_group(group::Info { sequence: 6 }).unwrap();
3610 assert_eq!(consumer.assert_group().sequence, 6);
3611 let done = consumer
3612 .recv_group()
3613 .now_or_never()
3614 .expect("should not block")
3615 .expect("would have errored");
3616 assert!(done.is_none(), "track completes once the boundary is reached");
3617 }
3618
3619 #[tokio::test]
3620 async fn recv_group_finishes_without_waiting_for_gaps() {
3621 let mut producer = track_producer("test", None);
3622 producer.create_group(group::Info { sequence: 1 }).unwrap();
3623 producer.finish().unwrap();
3624
3625 let mut consumer = producer.subscribe(None);
3626 assert_eq!(consumer.assert_group().sequence, 1);
3627
3628 let done = consumer
3629 .recv_group()
3630 .now_or_never()
3631 .expect("should not block")
3632 .expect("would have errored");
3633 assert!(done.is_none(), "track should finish without waiting for gaps");
3634 }
3635
3636 #[tokio::test]
3637 async fn next_group_skips_late_arrivals() {
3638 let mut producer = track_producer("test", None);
3639 let mut consumer = producer.subscribe(None);
3640
3641 producer.create_group(group::Info { sequence: 5 }).unwrap();
3643 let group = consumer
3644 .next_group()
3645 .now_or_never()
3646 .expect("should not block")
3647 .expect("would have errored")
3648 .expect("track should not be closed");
3649 assert_eq!(group.sequence, 5);
3650
3651 producer.create_group(group::Info { sequence: 3 }).unwrap();
3653 producer.create_group(group::Info { sequence: 4 }).unwrap();
3655 producer.create_group(group::Info { sequence: 7 }).unwrap();
3657
3658 let group = consumer
3659 .next_group()
3660 .now_or_never()
3661 .expect("should not block")
3662 .expect("would have errored")
3663 .expect("track should not be closed");
3664 assert_eq!(group.sequence, 7);
3665
3666 assert!(
3668 consumer.next_group().now_or_never().is_none(),
3669 "should block waiting for a higher sequence"
3670 );
3671 }
3672
3673 #[tokio::test]
3674 async fn next_group_returns_arrivals_in_order() {
3675 let mut producer = track_producer("test", None);
3676 let mut consumer = producer.subscribe(None);
3677
3678 producer.create_group(group::Info { sequence: 3 }).unwrap();
3680 producer.create_group(group::Info { sequence: 5 }).unwrap();
3681
3682 let group = consumer
3683 .next_group()
3684 .now_or_never()
3685 .expect("should not block")
3686 .expect("would have errored")
3687 .expect("track should not be closed");
3688 assert_eq!(group.sequence, 3);
3689
3690 let group = consumer
3691 .next_group()
3692 .now_or_never()
3693 .expect("should not block")
3694 .expect("would have errored")
3695 .expect("track should not be closed");
3696 assert_eq!(group.sequence, 5);
3697 }
3698
3699 #[tokio::test]
3700 async fn next_group_and_recv_group_use_independent_cursors() {
3701 let mut producer = track_producer("test", None);
3702 let mut consumer = producer.subscribe(None);
3703
3704 producer.create_group(group::Info { sequence: 5 }).unwrap();
3706 producer.create_group(group::Info { sequence: 3 }).unwrap();
3707
3708 let group = consumer
3711 .next_group()
3712 .now_or_never()
3713 .expect("should not block")
3714 .expect("would have errored")
3715 .expect("track should not be closed");
3716 assert_eq!(group.sequence, 3);
3717
3718 assert_eq!(consumer.assert_group().sequence, 5);
3721 }
3722
3723 #[tokio::test]
3724 async fn end_at_caps_next_group() {
3725 let mut producer = track_producer("test", None);
3726 let mut consumer = producer.subscribe(None);
3727
3728 for s in 0..6 {
3729 producer.create_group(group::Info { sequence: s }).unwrap();
3730 }
3731
3732 consumer.end_at(2);
3733
3734 assert_eq!(
3736 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3737 0
3738 );
3739 assert_eq!(
3740 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3741 1
3742 );
3743 assert_eq!(
3744 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3745 2
3746 );
3747
3748 assert!(
3750 consumer.next_group().now_or_never().is_none(),
3751 "capped consumer must block instead of returning out-of-range groups"
3752 );
3753 }
3754
3755 #[tokio::test]
3756 async fn end_at_release_drains_cached_groups() {
3757 let mut producer = track_producer("test", None);
3758 let mut consumer = producer.subscribe(None);
3759
3760 for s in 0..6 {
3761 producer.create_group(group::Info { sequence: s }).unwrap();
3762 }
3763
3764 consumer.end_at(1);
3765 assert_eq!(
3766 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3767 0
3768 );
3769 assert_eq!(
3770 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3771 1
3772 );
3773 assert!(consumer.next_group().now_or_never().is_none(), "capped at 1");
3774
3775 consumer.end_at(4);
3777 assert_eq!(
3778 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3779 2
3780 );
3781 assert_eq!(
3782 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3783 3
3784 );
3785 assert_eq!(
3786 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3787 4
3788 );
3789 assert!(consumer.next_group().now_or_never().is_none(), "capped at 4");
3790
3791 consumer.end_at(None);
3793 assert_eq!(
3794 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3795 5
3796 );
3797 assert!(consumer.next_group().now_or_never().is_none(), "no more groups");
3798 }
3799
3800 #[tokio::test]
3801 async fn end_at_lower_than_cursor_parks_consumer() {
3802 let mut producer = track_producer("test", None);
3803 let mut consumer = producer.subscribe(None);
3804
3805 for s in 0..3 {
3806 producer.create_group(group::Info { sequence: s }).unwrap();
3807 }
3808
3809 assert_eq!(
3811 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3812 0
3813 );
3814 assert_eq!(
3815 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3816 1
3817 );
3818 assert_eq!(
3819 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3820 2
3821 );
3822
3823 consumer.end_at(1);
3825 producer.create_group(group::Info { sequence: 3 }).unwrap();
3826 producer.create_group(group::Info { sequence: 4 }).unwrap();
3827 assert!(
3828 consumer.next_group().now_or_never().is_none(),
3829 "cap is below cursor; nothing returnable until cap rises"
3830 );
3831
3832 consumer.end_at(None);
3834 assert_eq!(
3835 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3836 3
3837 );
3838 assert_eq!(
3839 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3840 4
3841 );
3842 }
3843
3844 #[tokio::test]
3845 async fn end_at_toggling_around_late_arrivals() {
3846 let mut producer = track_producer("test", None);
3847 let mut consumer = producer.subscribe(None);
3848
3849 consumer.end_at(5);
3850
3851 producer.create_group(group::Info { sequence: 2 }).unwrap();
3853 producer.create_group(group::Info { sequence: 5 }).unwrap();
3854 producer.create_group(group::Info { sequence: 3 }).unwrap();
3855 producer.create_group(group::Info { sequence: 8 }).unwrap();
3857 producer.create_group(group::Info { sequence: 4 }).unwrap();
3858
3859 assert_eq!(
3861 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3862 2
3863 );
3864 assert_eq!(
3865 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3866 3
3867 );
3868 assert_eq!(
3869 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3870 4
3871 );
3872 assert_eq!(
3873 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3874 5
3875 );
3876 assert!(consumer.next_group().now_or_never().is_none());
3878
3879 consumer.end_at(10);
3881 assert_eq!(
3882 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3883 8
3884 );
3885 }
3886
3887 #[tokio::test]
3891 async fn end_at_parks_recv_group() {
3892 let mut producer = track_producer("test", None);
3893 let mut consumer = producer.subscribe(None);
3894
3895 for s in 0..3 {
3896 producer.create_group(group::Info { sequence: s }).unwrap();
3897 }
3898
3899 consumer.end_at(1);
3900 assert_eq!(
3901 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3902 0
3903 );
3904 assert_eq!(
3905 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3906 1
3907 );
3908 assert!(consumer.recv_group().now_or_never().is_none(), "capped at 1");
3909
3910 producer.finish().unwrap();
3912 assert!(
3913 consumer.recv_group().now_or_never().is_none(),
3914 "still parked after finish"
3915 );
3916
3917 consumer.end_at(None);
3918 assert_eq!(
3919 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3920 2
3921 );
3922 assert!(
3923 matches!(consumer.recv_group().now_or_never(), Some(Ok(None))),
3924 "finished once the parked group drains"
3925 );
3926 }
3927
3928 #[tokio::test]
3931 async fn recv_group_serves_arrivals_behind_the_cap() {
3932 let mut producer = track_producer("test", None);
3933 let mut consumer = producer.subscribe(None);
3934
3935 consumer.end_at(1);
3936
3937 producer.create_group(group::Info { sequence: 2 }).unwrap();
3939 producer.create_group(group::Info { sequence: 0 }).unwrap();
3940 producer.create_group(group::Info { sequence: 1 }).unwrap();
3941
3942 assert_eq!(
3943 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3944 0
3945 );
3946 assert_eq!(
3947 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3948 1
3949 );
3950 assert!(consumer.recv_group().now_or_never().is_none(), "capped at 1");
3951
3952 consumer.end_at(2);
3953 assert_eq!(
3954 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3955 2
3956 );
3957 }
3958
3959 #[tokio::test]
3962 async fn start_at_drops_parked_recv_groups() {
3963 let mut producer = track_producer("test", None);
3964 let mut consumer = producer.subscribe(None);
3965
3966 consumer.end_at(0);
3967 producer.create_group(group::Info { sequence: 1 }).unwrap();
3968 assert!(
3969 consumer.recv_group().now_or_never().is_none(),
3970 "group 1 parked at the cap"
3971 );
3972
3973 consumer.start_at(2);
3974 consumer.end_at(None);
3975 producer.create_group(group::Info { sequence: 2 }).unwrap();
3976 assert_eq!(
3977 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3978 2,
3979 "the overtaken parked group is dropped, not re-offered"
3980 );
3981 }
3982
3983 #[tokio::test]
3987 async fn evicted_parked_recv_groups_are_dropped() {
3988 let mut producer = track_producer("test", None);
3989 let mut consumer = producer.subscribe(None);
3990
3991 producer.create_group(group::Info { sequence: 0 }).unwrap();
3992 assert_eq!(
3993 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3994 0
3995 );
3996
3997 consumer.end_at(0);
3998 let straggler = producer.create_group(group::Info { sequence: 1 }).unwrap();
3999 assert!(
4000 consumer.recv_group().now_or_never().is_none(),
4001 "group 1 parked at the cap"
4002 );
4003
4004 straggler.abort(Error::Old).unwrap();
4006 producer.finish().unwrap();
4007
4008 consumer.end_at(None);
4009 assert!(
4010 matches!(consumer.recv_group().now_or_never(), Some(Ok(None))),
4011 "a dead parked group must not be delivered or hold the stream open"
4012 );
4013 }
4014
4015 #[tokio::test]
4019 async fn evicted_parked_group_wakes_the_clean_end() {
4020 use std::sync::atomic::{AtomicUsize, Ordering};
4021 use std::task::{Context, Wake};
4022
4023 struct CountWaker(AtomicUsize);
4026 impl Wake for CountWaker {
4027 fn wake(self: std::sync::Arc<Self>) {
4028 self.0.fetch_add(1, Ordering::SeqCst);
4029 }
4030 }
4031
4032 let mut producer = track_producer("test", None);
4033 let mut consumer = producer.subscribe(None);
4034
4035 producer.create_group(group::Info { sequence: 0 }).unwrap();
4036 assert_eq!(
4037 consumer.recv_group().now_or_never().unwrap().unwrap().unwrap().sequence,
4038 0
4039 );
4040
4041 consumer.end_at(0);
4042 let straggler = producer.create_group(group::Info { sequence: 1 }).unwrap();
4043 assert!(consumer.recv_group().now_or_never().is_none(), "parked at the cap");
4044 producer.finish().unwrap();
4045
4046 let counter = std::sync::Arc::new(CountWaker(AtomicUsize::new(0)));
4047 let waker = std::task::Waker::from(counter.clone());
4048 let mut cx = Context::from_waker(&waker);
4049 let mut fut = std::pin::pin!(consumer.recv_group());
4050 assert!(
4051 fut.as_mut().poll(&mut cx).is_pending(),
4052 "the parked group holds it open"
4053 );
4054
4055 straggler.abort(Error::Old).unwrap();
4056 assert!(counter.0.load(Ordering::SeqCst) > 0, "the eviction wakeup was lost");
4057 assert!(matches!(fut.as_mut().poll(&mut cx), Poll::Ready(Ok(None))));
4058 }
4059
4060 #[tokio::test]
4061 async fn read_frame_returns_single_frame_per_group() {
4062 let mut producer = track_producer("test", None);
4063 let mut consumer = producer.subscribe(None);
4064
4065 producer.write_frame(Timestamp::ZERO, b"hello".as_slice()).unwrap();
4066 producer.write_frame(Timestamp::ZERO, b"world".as_slice()).unwrap();
4067
4068 let frame = consumer
4069 .read_frame()
4070 .now_or_never()
4071 .expect("should not block")
4072 .expect("would have errored")
4073 .expect("track should not be closed");
4074 assert_eq!(&frame.payload[..], b"hello");
4075
4076 let frame = consumer
4077 .read_frame()
4078 .now_or_never()
4079 .expect("should not block")
4080 .expect("would have errored")
4081 .expect("track should not be closed");
4082 assert_eq!(&frame.payload[..], b"world");
4083 }
4084
4085 #[tokio::test]
4086 async fn read_frame_preserves_timestamp() {
4087 let mut producer = track_producer("test", None);
4088 let mut consumer = producer.subscribe(None);
4089
4090 producer
4091 .write_frame(Timestamp::from_micros(20_000).unwrap(), b"hello".as_slice())
4092 .unwrap();
4093
4094 let frame = consumer
4095 .read_frame()
4096 .now_or_never()
4097 .expect("should not block")
4098 .expect("would have errored")
4099 .expect("track should not be closed");
4100 assert_eq!(frame.timestamp.as_micros(), 20_000);
4101 assert_eq!(&frame.payload[..], b"hello");
4102 }
4103
4104 #[tokio::test]
4105 async fn read_frame_skips_stalled_group_for_newer_ready_frame() {
4106 let mut producer = track_producer("test", None);
4107 let mut consumer = producer.subscribe(None);
4108
4109 let _stalled = producer.create_group(group::Info { sequence: 3 }).unwrap();
4111 let mut g5 = producer.create_group(group::Info { sequence: 5 }).unwrap();
4113 g5.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"later"))
4114 .unwrap();
4115 g5.finish().unwrap();
4116
4117 let frame = consumer
4119 .read_frame()
4120 .now_or_never()
4121 .expect("should not block on stalled earlier group")
4122 .expect("would have errored")
4123 .expect("track should not be closed");
4124 assert_eq!(&frame.payload[..], b"later");
4125 }
4126
4127 #[tokio::test]
4128 async fn read_frame_discards_rest_of_multi_frame_group() {
4129 let mut producer = track_producer("test", None);
4130 let mut consumer = producer.subscribe(None);
4131
4132 let mut g0 = producer.create_group(group::Info { sequence: 0 }).unwrap();
4134 g0.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"one"))
4135 .unwrap();
4136 g0.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"two"))
4137 .unwrap();
4138 g0.finish().unwrap();
4139
4140 producer.write_frame(Timestamp::ZERO, b"next".as_slice()).unwrap();
4142
4143 let frame = consumer
4144 .read_frame()
4145 .now_or_never()
4146 .expect("should not block")
4147 .expect("would have errored")
4148 .expect("track should not be closed");
4149 assert_eq!(&frame.payload[..], b"one");
4150
4151 let frame = consumer
4153 .read_frame()
4154 .now_or_never()
4155 .expect("should not block")
4156 .expect("would have errored")
4157 .expect("track should not be closed");
4158 assert_eq!(&frame.payload[..], b"next");
4159 }
4160
4161 #[tokio::test]
4162 async fn read_frame_waits_for_pending_group_after_finish() {
4163 let mut producer = track_producer("test", None);
4166 let mut consumer = producer.subscribe(None);
4167
4168 let mut g0 = producer.create_group(group::Info { sequence: 0 }).unwrap();
4169 producer.finish().unwrap();
4170
4171 assert!(
4173 consumer.read_frame().now_or_never().is_none(),
4174 "read_frame must block on a pending group even after finish()"
4175 );
4176
4177 g0.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"late"))
4179 .unwrap();
4180 let frame = consumer
4181 .read_frame()
4182 .now_or_never()
4183 .expect("should not block once a frame is written")
4184 .expect("would have errored")
4185 .expect("track should not be closed");
4186 assert_eq!(&frame.payload[..], b"late");
4187 }
4188
4189 #[tokio::test]
4190 async fn read_frame_respects_start_at() {
4191 let mut producer = track_producer("test", None);
4194 let mut consumer = producer.subscribe(None);
4195 consumer.start_at(5);
4196
4197 let mut g3 = producer.create_group(group::Info { sequence: 3 }).unwrap();
4199 g3.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"skip-me"))
4200 .unwrap();
4201 g3.finish().unwrap();
4202
4203 let mut g5 = producer.create_group(group::Info { sequence: 5 }).unwrap();
4204 g5.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"keep"))
4205 .unwrap();
4206 g5.finish().unwrap();
4207
4208 let frame = consumer
4209 .read_frame()
4210 .now_or_never()
4211 .expect("should not block")
4212 .expect("would have errored")
4213 .expect("track should not be closed");
4214 assert_eq!(&frame.payload[..], b"keep");
4215 }
4216
4217 #[tokio::test]
4218 async fn read_frame_returns_none_when_finished() {
4219 let mut producer = track_producer("test", None);
4220 let mut consumer = producer.subscribe(None);
4221
4222 producer.write_frame(Timestamp::ZERO, b"only".as_slice()).unwrap();
4223 producer.finish().unwrap();
4224
4225 let frame = consumer
4226 .read_frame()
4227 .now_or_never()
4228 .expect("should not block")
4229 .expect("would have errored")
4230 .expect("track should not be closed");
4231 assert_eq!(&frame.payload[..], b"only");
4232
4233 let done = consumer
4234 .read_frame()
4235 .now_or_never()
4236 .expect("should not block")
4237 .expect("would have errored");
4238 assert!(done.is_none());
4239 }
4240
4241 #[test]
4242 fn append_group_returns_bounds_exceeded_on_sequence_overflow() {
4243 let mut producer = track_producer("test", None);
4244 {
4245 let mut state = producer.state.write().ok().unwrap();
4246 state.max_sequence = Some(u64::MAX);
4247 }
4248
4249 assert!(matches!(producer.append_group(), Err(Error::BoundsExceeded(_))));
4250 }
4251
4252 #[tokio::test]
4253 async fn fetch_cache_hit() {
4254 let mut producer = track_producer("test", None);
4255
4256 let mut group = producer.append_group().unwrap(); group
4259 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hello"))
4260 .unwrap();
4261 group.finish().unwrap();
4262
4263 let dynamic = producer.dynamic();
4266 let consumer = producer.consume();
4267 assert!(consumer.peek_group(0).is_some());
4268 let mut g = consumer.fetch_group(0, None).await.unwrap();
4269 assert_eq!(g.sequence, 0);
4270 assert_eq!(&g.read_frame().await.unwrap().unwrap().payload[..], b"hello");
4271
4272 assert!(dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending());
4274 }
4275
4276 #[tokio::test]
4277 async fn fetch_miss_signals_dynamic() {
4278 let producer = track_producer("test", None);
4279 let dynamic = producer.dynamic();
4280 let consumer = producer.consume();
4281
4282 assert!(consumer.peek_group(5).is_none());
4286 let pending = consumer.fetch_group(5, group::Fetch::default().with_priority(7));
4287 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
4288
4289 let req = dynamic
4290 .requested_group()
4291 .now_or_never()
4292 .expect("should not block")
4293 .unwrap();
4294 assert_eq!(req.sequence(), 5);
4295 assert_eq!(req.priority(), 7);
4296
4297 let mut group = req.accept(None).unwrap();
4299 group
4300 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hi"))
4301 .unwrap();
4302 group.finish().unwrap();
4303
4304 let mut g = pending.await.unwrap();
4305 assert_eq!(g.sequence, 5);
4306 assert_eq!(&g.read_frame().await.unwrap().unwrap().payload[..], b"hi");
4307 }
4308
4309 #[tokio::test]
4310 async fn fetch_miss_rejects() {
4311 let producer = track_producer("test", None);
4312 let dynamic = producer.dynamic();
4313 let consumer = producer.consume();
4314
4315 let pending = consumer.fetch_group(5, None);
4316 let req = dynamic
4317 .requested_group()
4318 .now_or_never()
4319 .expect("should not block")
4320 .unwrap();
4321
4322 req.reject(Error::Cancel);
4323 assert!(matches!(pending.await, Err(Error::Cancel)));
4324 let fetch = producer.state.read().fetch.clone();
4325 assert!(fetch.read().is_empty());
4326 }
4327
4328 #[tokio::test]
4329 async fn fetch_miss_drop_rejects() {
4330 let producer = track_producer("test", None);
4331 let dynamic = producer.dynamic();
4332 let consumer = producer.consume();
4333
4334 let pending = consumer.fetch_group(5, None);
4335 let req = dynamic
4336 .requested_group()
4337 .now_or_never()
4338 .expect("should not block")
4339 .unwrap();
4340
4341 drop(req);
4342 assert!(matches!(pending.await, Err(Error::Dropped)));
4343 }
4344
4345 #[tokio::test]
4346 async fn fetch_reject_does_not_poison_retry() {
4347 let producer = track_producer("test", None);
4348 let dynamic = producer.dynamic();
4349 let consumer = producer.consume();
4350
4351 let pending = consumer.fetch_group(5, None);
4352 let req = dynamic
4353 .requested_group()
4354 .now_or_never()
4355 .expect("should not block")
4356 .unwrap();
4357 req.reject(Error::Cancel);
4358 assert!(matches!(pending.await, Err(Error::Cancel)));
4359
4360 let retry = consumer.fetch_group(5, None);
4361 let req = dynamic
4362 .requested_group()
4363 .now_or_never()
4364 .expect("should not block")
4365 .unwrap();
4366 let mut group = req.accept(None).unwrap();
4367 group
4368 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"retry"))
4369 .unwrap();
4370 group.finish().unwrap();
4371
4372 let mut group = retry.await.unwrap();
4373 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"retry");
4374 }
4375
4376 #[tokio::test]
4377 async fn fetch_coalesces_concurrent() {
4378 let producer = track_producer("test", None);
4379 let dynamic = producer.dynamic();
4380 let consumer = producer.consume();
4381
4382 let first = consumer.fetch_group(5, group::Fetch::default().with_priority(1));
4385 let second = consumer.fetch_group(5, group::Fetch::default().with_priority(7));
4386 assert!(kio::Pollable::poll(&*first, &kio::Waiter::noop()).is_pending());
4387
4388 let req = dynamic
4389 .requested_group()
4390 .now_or_never()
4391 .expect("should not block")
4392 .unwrap();
4393 assert_eq!(req.sequence(), 5);
4394 assert_eq!(req.priority(), 7);
4395 assert!(
4396 dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending(),
4397 "the second fetch queued a duplicate request"
4398 );
4399
4400 let third = consumer.fetch_group(5, None);
4402
4403 let mut group = req.accept(None).unwrap();
4405 group
4406 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hi"))
4407 .unwrap();
4408 group.finish().unwrap();
4409
4410 assert_eq!(first.await.unwrap().sequence, 5);
4411 assert_eq!(second.await.unwrap().sequence, 5);
4412 assert_eq!(third.await.unwrap().sequence, 5);
4413 }
4414
4415 #[tokio::test]
4416 async fn fetch_coalesced_reject_fails_all() {
4417 let producer = track_producer("test", None);
4418 let dynamic = producer.dynamic();
4419 let consumer = producer.consume();
4420
4421 let first = consumer.fetch_group(5, None);
4422 let second = consumer.fetch_group(5, None);
4423 let req = dynamic
4424 .requested_group()
4425 .now_or_never()
4426 .expect("should not block")
4427 .unwrap();
4428 req.reject(Error::Cancel);
4429
4430 assert!(matches!(first.await, Err(Error::Cancel)));
4431 assert!(matches!(second.await, Err(Error::Cancel)));
4432
4433 let retry = consumer.fetch_group(5, None);
4435 assert!(kio::Pollable::poll(&*retry, &kio::Waiter::noop()).is_pending());
4436 let req = dynamic
4437 .requested_group()
4438 .now_or_never()
4439 .expect("should not block")
4440 .unwrap();
4441 assert_eq!(req.sequence(), 5);
4442 }
4443
4444 #[tokio::test]
4445 async fn fetch_queued_fails_when_handlers_leave() {
4446 let producer = track_producer("test", None);
4447 let dynamic = producer.dynamic();
4448 let consumer = producer.consume();
4449
4450 let pending = consumer.fetch_group(5, None);
4452 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
4453 drop(dynamic);
4454 assert!(matches!(pending.await, Err(Error::NotFound)));
4455
4456 let fetch = producer.state.read().fetch.clone();
4458 assert!(fetch.read().is_empty());
4459 }
4460
4461 #[tokio::test]
4462 async fn fetch_miss_no_dynamic_not_found() {
4463 let mut producer = track_producer("test", None);
4466 producer.append_group().unwrap(); let consumer = producer.consume();
4468 assert!(matches!(consumer.fetch_group(5, None).await, Err(Error::NotFound)));
4469 }
4470
4471 #[tokio::test]
4472 async fn fetch_past_final_not_found() {
4473 let mut producer = track_producer("test", None);
4474 producer.append_group().unwrap(); producer.finish().unwrap(); let dynamic = producer.dynamic();
4480 let consumer = producer.consume();
4481 assert!(matches!(consumer.fetch_group(5, None).await, Err(Error::NotFound)));
4482
4483 assert!(dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending());
4485 }
4486
4487 fn pooled_producer(capacity: u64) -> (Producer, cache::Pool) {
4489 let pool = cache::Pool::new(capacity);
4490 let broadcast = broadcast::Info {
4491 origin: crate::origin::Info::default().with_pool(pool.clone()),
4492 ..Default::default()
4493 };
4494 let producer = Producer::new(Arc::new(broadcast), "test", None);
4495 (producer, pool)
4496 }
4497
4498 fn finished_group(producer: &mut Producer, size: usize) -> u64 {
4499 let mut group = producer.append_group().unwrap();
4500 group
4501 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; size]))
4502 .unwrap();
4503 group.finish().unwrap();
4504 group.sequence
4505 }
4506
4507 #[tokio::test]
4510 async fn debt_evicts_oldest_group() {
4511 tokio::time::pause();
4512
4513 let (mut producer, pool) = pooled_producer(10_000);
4515
4516 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let consumer = producer.consume();
4521 assert!(consumer.peek_group(0).is_none(), "oldest group is evicted");
4522 assert!(consumer.peek_group(2).is_some(), "latest group survives");
4523 assert!(pool.used() <= 21_000, "usage hovers near capacity: {}", pool.used());
4526
4527 let mut subscriber = producer.subscribe(None);
4529 assert!(subscriber.assert_group().sequence > 0, "evicted group is not delivered");
4530 }
4531
4532 #[tokio::test]
4534 async fn latest_group_never_evicted() {
4535 tokio::time::pause();
4536
4537 let (mut producer, pool) = pooled_producer(100);
4539 finished_group(&mut producer, 1000); assert!(pool.used() > 100, "the latest may exceed the budget");
4541
4542 finished_group(&mut producer, 1000); finished_group(&mut producer, 1000); let consumer = producer.consume();
4547 assert!(consumer.peek_group(0).is_none());
4548 let mut group = consumer.peek_group(2).expect("latest survives");
4549 assert_eq!(group.read_frame().await.unwrap().unwrap().payload.len(), 1000);
4550 }
4551
4552 #[tokio::test]
4556 async fn fetch_refresh_survives_eviction() {
4557 tokio::time::pause();
4558
4559 let (mut producer, _pool) = pooled_producer(10_000);
4560 let consumer = producer.consume();
4561
4562 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_secs(1)).await;
4564 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_secs(1)).await;
4566 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_millis(500)).await;
4568
4569 let mut fetched = consumer.fetch_group(0, None).await.unwrap();
4571 assert_eq!(fetched.read_frame().await.unwrap().unwrap().payload.len(), 3_000);
4572 tokio::time::advance(Duration::from_millis(500)).await;
4573
4574 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_secs(1)).await;
4578 finished_group(&mut producer, 3_000); assert!(consumer.peek_group(0).is_some(), "refreshed group survives");
4581 assert!(consumer.peek_group(1).is_none(), "unread group is evicted instead");
4582 }
4583
4584 #[tokio::test]
4587 async fn eviction_aborts_readers() {
4588 tokio::time::pause();
4589
4590 let (mut producer, _pool) = pooled_producer(10_000);
4591 let mut subscriber = producer.subscribe(None);
4592
4593 finished_group(&mut producer, 10_000); let mut group0 = subscriber.assert_group();
4595
4596 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let read = group0.read_frame().await;
4600 assert!(matches!(read, Err(Error::Evicted)), "expected Evicted, got {read:?}");
4601 }
4602
4603 #[tokio::test]
4607 async fn small_writes_carry_debt() {
4608 tokio::time::pause();
4609
4610 let (mut producer, pool) = pooled_producer(22_000);
4611 let consumer = producer.consume();
4612
4613 finished_group(&mut producer, 20_000); for _ in 0..3 {
4618 finished_group(&mut producer, 1_000);
4619 }
4620 assert!(consumer.peek_group(0).is_some(), "debt smaller than the victim carries");
4621
4622 for _ in 0..20 {
4624 finished_group(&mut producer, 1_000);
4625 }
4626 assert!(
4627 consumer.peek_group(0).is_none(),
4628 "accumulated debt evicts the large group"
4629 );
4630 assert!(pool.used() <= 24_000, "usage hovers near capacity: {}", pool.used());
4633 }
4634
4635 #[tokio::test]
4639 async fn payment_capped_per_write() {
4640 tokio::time::pause();
4641
4642 let (mut producer, pool) = pooled_producer(1 << 40);
4643 for _ in 0..10 {
4644 finished_group(&mut producer, 1_000);
4645 }
4646
4647 pool.resize(100);
4649 let before = pool.used();
4650
4651 finished_group(&mut producer, 1_000);
4653
4654 let consumer = producer.consume();
4655 assert!(consumer.peek_group(0).is_none(), "the oldest groups are evicted");
4656 assert!(consumer.peek_group(1).is_none());
4657 assert!(consumer.peek_group(2).is_some(), "the backlog drains gradually");
4658 assert!(pool.used() > before - 4_000, "one write must not dump the backlog");
4659 }
4660
4661 #[tokio::test]
4665 async fn accept_preserves_write_accounting() {
4666 tokio::time::pause();
4667
4668 let pool = cache::Pool::new(12_000);
4669 let broadcast = broadcast::Info {
4670 origin: crate::origin::Info::default().with_pool(pool.clone()),
4671 ..Default::default()
4672 };
4673 let request = Request::new(Arc::new(broadcast), "test");
4674 let dynamic = request.dynamic();
4675 let consumer = request.consume();
4676
4677 let pending = consumer.fetch_group(0, None);
4679 let req = dynamic
4680 .requested_group()
4681 .now_or_never()
4682 .expect("should not block")
4683 .unwrap();
4684 let mut backfill = req.accept(None).unwrap();
4685 pending.await.unwrap();
4686 backfill
4687 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 30_000]))
4688 .unwrap();
4689
4690 let mut producer = request.accept(None);
4693 producer.append_group().unwrap().finish().unwrap();
4694 producer.append_group().unwrap().finish().unwrap();
4695
4696 assert!(
4697 producer.consume().peek_group(0).is_none(),
4698 "pre-accept backfill growth is reclaimed after accept"
4699 );
4700 assert!(pool.used() <= 13_000, "usage converges: {}", pool.used());
4701 }
4702
4703 #[tokio::test]
4706 async fn recreated_sequence_bounds_eviction_hints() {
4707 let (mut producer, _pool) = pooled_producer(1 << 40);
4708 producer.create_group(5u64.into()).unwrap().finish().unwrap();
4709
4710 for _ in 0..200 {
4711 let group = producer.create_group(1u64.into()).unwrap();
4712 group.abort(Error::Cancel).unwrap();
4713 }
4714
4715 let state = producer.state.read();
4716 assert!(
4717 state.evict.len() <= 2 * state.lookup.len() + EVICT_SLACK,
4718 "stale hints are compacted: {} entries for {} slots",
4719 state.evict.len(),
4720 state.lookup.len()
4721 );
4722 }
4723
4724 #[tokio::test]
4727 async fn same_tick_write_outranks_inserted() {
4728 tokio::time::pause();
4729
4730 let (mut producer, _pool) = pooled_producer(10_000);
4732
4733 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();
4740 assert!(consumer.peek_group(0).is_none(), "insert-only content pays first");
4741 assert!(consumer.peek_group(1).is_some(), "same-tick written content survives");
4742 }
4743
4744 #[tokio::test]
4747 async fn frame_only_writer_pays() {
4748 tokio::time::pause();
4749
4750 let (mut producer, pool) = pooled_producer(2_000);
4751 let mut demoted = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); demoted
4757 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 300_000]))
4758 .unwrap();
4759
4760 assert!(
4761 pool.used() <= 5_000,
4762 "the frame write settled the debt: {}",
4763 pool.used()
4764 );
4765 assert!(matches!(demoted.finish(), Err(Error::Evicted)));
4766 }
4767
4768 #[tokio::test]
4771 async fn each_track_owns_its_account() {
4772 let broadcast = Arc::new(broadcast::Info::default());
4773 let info = Info::default();
4774 let a = Producer::new(broadcast.clone(), "a", info.clone());
4775 let b = Producer::new(broadcast, "b", info);
4776
4777 let a = a.state.read().cache.clone();
4778 let b = b.state.read().cache.clone();
4779 assert!(!Arc::ptr_eq(&a, &b), "each track owns its account");
4780 }
4781
4782 #[tokio::test]
4785 async fn a_dynamic_defers_teardown() {
4786 let (mut producer, pool) = pooled_producer(1 << 40);
4787 let dynamic = producer.dynamic();
4788 finished_group(&mut producer, 100);
4789
4790 drop(producer);
4791 assert!(pool.used() > 0, "the handler still serves the cache");
4792
4793 drop(dynamic);
4794 assert_eq!(pool.used(), 0, "the last handle tears it down");
4795 }
4796
4797 #[tokio::test]
4803 async fn finished_track_frees_its_cache() {
4804 let (mut producer, pool) = pooled_producer(1 << 40);
4805 finished_group(&mut producer, 100);
4806 producer.finish().unwrap();
4807
4808 let state = producer.state.downgrade();
4809 drop(producer);
4810
4811 assert!(state.upgrade().is_none(), "the track state is freed");
4812 assert_eq!(pool.used(), 0, "so are its cached bytes");
4813 }
4814
4815 #[tokio::test]
4819 async fn teardown_ignores_a_settling_group() {
4820 let (mut producer, pool) = pooled_producer(1 << 40);
4821 finished_group(&mut producer, 100);
4822
4823 let settling = producer.state.downgrade().upgrade().expect("open");
4825 drop(producer);
4826
4827 assert_eq!(pool.used(), 0, "the abrupt teardown still released the cache");
4828 drop(settling);
4829 }
4830
4831 #[tokio::test]
4834 async fn cached_group_outlives_its_track() {
4835 let (mut producer, pool) = pooled_producer(1 << 40);
4836 let sequence = finished_group(&mut producer, 100);
4837 let group = producer.consume().peek_group(sequence).expect("cached");
4838 producer.finish().unwrap();
4839
4840 let state = producer.state.downgrade();
4841 drop(producer);
4842 assert!(state.upgrade().is_none(), "the track state is freed");
4843 assert!(pool.used() > 0, "the retained group keeps its own bytes");
4844
4845 drop(group);
4846 assert_eq!(pool.used(), 0, "which it releases when dropped");
4847 }
4848
4849 #[tokio::test]
4853 async fn pre_accept_backfill_settles_late_writes() {
4854 tokio::time::pause();
4855
4856 let pool = cache::Pool::new(2_000);
4857 let broadcast = broadcast::Info {
4858 origin: crate::origin::Info::default().with_pool(pool.clone()),
4859 ..Default::default()
4860 };
4861 let request = Request::new(Arc::new(broadcast), "test");
4862 let dynamic = request.dynamic();
4863 let consumer = request.consume();
4864
4865 let pending = consumer.fetch_group(0, None);
4867 let req = dynamic
4868 .requested_group()
4869 .now_or_never()
4870 .expect("should not block")
4871 .unwrap();
4872 let mut backfill = req.accept(None).unwrap();
4873 pending.await.unwrap();
4874
4875 let mut producer = request.accept(None);
4877 producer.append_group().unwrap().finish().unwrap();
4878
4879 backfill
4882 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 300_000]))
4883 .unwrap();
4884
4885 assert!(
4886 pool.used() <= 5_000,
4887 "the frame write settled the debt: {}",
4888 pool.used()
4889 );
4890 }
4891
4892 #[tokio::test]
4896 async fn write_restarts_retention_clock() {
4897 tokio::time::pause();
4898
4899 let (mut producer, _pool) = pooled_producer(1 << 40);
4900 let mut straggler = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
4905 straggler
4906 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 100]))
4907 .unwrap();
4908 producer.append_group().unwrap().finish().unwrap(); let consumer = producer.consume();
4911 assert!(consumer.peek_group(0).is_some(), "the write restarted the clock");
4912
4913 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
4915 producer.append_group().unwrap().finish().unwrap(); assert!(consumer.peek_group(0).is_none(), "idle content still expires");
4917 }
4918
4919 #[tokio::test]
4922 async fn refreshed_front_does_not_starve_expiry() {
4923 tokio::time::pause();
4924
4925 let (mut producer, _pool) = pooled_producer(1 << 40);
4926 let dynamic = producer.dynamic();
4927 let consumer = producer.consume();
4928
4929 producer.create_group(10u64.into()).unwrap().finish().unwrap();
4930 for sequence in 1..=5u64 {
4931 let pending = consumer.fetch_group(sequence, None);
4932 let req = dynamic
4933 .requested_group()
4934 .now_or_never()
4935 .expect("should not block")
4936 .unwrap();
4937 let mut group = req.accept(None).unwrap();
4938 group
4939 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 100]))
4940 .unwrap();
4941 group.finish().unwrap();
4942 pending.await.unwrap();
4943 }
4944
4945 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
4948 for sequence in 1..=4u64 {
4949 consumer.fetch_group(sequence, None).await.unwrap();
4950 }
4951
4952 for _ in 0..3 {
4954 producer.append_group().unwrap().finish().unwrap();
4955 }
4956 assert!(consumer.peek_group(5).is_none(), "expired backfill is reclaimed");
4957 assert!(consumer.peek_group(1).is_some(), "refreshed backfill survives");
4958 }
4959
4960 #[tokio::test]
4963 async fn recreated_sequence_delivered_once() {
4964 let (mut producer, _pool) = pooled_producer(1 << 40);
4965
4966 producer.create_group(0u64.into()).unwrap().finish().unwrap();
4967 let aborted = producer.create_group(1u64.into()).unwrap();
4968 aborted.abort(Error::Cancel).unwrap();
4969 producer.create_group(2u64.into()).unwrap().finish().unwrap();
4970 producer.create_group(1u64.into()).unwrap().finish().unwrap();
4971
4972 let mut subscriber = producer.subscribe(None);
4973 assert_eq!(subscriber.assert_group().sequence, 0);
4974 assert_eq!(subscriber.assert_group().sequence, 2);
4975 assert_eq!(
4976 subscriber.assert_group().sequence,
4977 1,
4978 "replacement arrives at its own position"
4979 );
4980 subscriber.assert_no_group();
4981 }
4982
4983 #[tokio::test]
4987 async fn datagrams_do_not_block_eviction() {
4988 tokio::time::pause();
4989
4990 let (mut producer, pool) = pooled_producer(1_000);
4991 for _ in 0..10 {
4992 finished_group(&mut producer, 1_000);
4993 producer.append_datagram(Timestamp::ZERO, &b"beat"[..]).unwrap();
4994 }
4995
4996 let consumer = producer.consume();
4997 assert!(consumer.peek_group(0).is_none(), "old groups still evict");
4998 assert!(
4999 pool.used() < 4 * 1_256,
5000 "interleaved datagrams must not bypass the budget: {}",
5001 pool.used()
5002 );
5003 }
5004
5005 #[tokio::test]
5009 async fn aborted_group_leaves_no_ghost_sample() {
5010 tokio::time::pause();
5011
5012 let (mut producer, pool) = pooled_producer(1 << 40);
5013 let group0 = producer.append_group().unwrap();
5014 producer.append_group().unwrap(); assert!(pool.average().is_some(), "demoted group is sampled");
5017 group0.abort(Error::Cancel).unwrap();
5018 assert_eq!(pool.average(), None, "the abort must remove the sample");
5019 }
5020
5021 #[tokio::test]
5024 async fn empty_groups_repay_overhead() {
5025 tokio::time::pause();
5026
5027 let (mut producer, pool) = pooled_producer(1_000);
5028 for _ in 0..100 {
5029 let mut group = producer.append_group().unwrap();
5030 group.finish().unwrap();
5031 }
5032
5033 assert!(
5034 pool.used() <= 3_000,
5035 "empty-group overhead must stay near the budget: {}",
5036 pool.used()
5037 );
5038 }
5039
5040 #[tokio::test]
5043 async fn growth_on_demoted_group_is_billed() {
5044 tokio::time::pause();
5045
5046 let (mut producer, pool) = pooled_producer(2_000);
5047 let mut straggler = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); straggler
5052 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 10_000]))
5053 .unwrap();
5054
5055 producer.append_group().unwrap().finish().unwrap(); let consumer = producer.consume();
5059 assert!(consumer.peek_group(0).is_none(), "the ballooned group is evicted");
5060 assert!(pool.used() <= 3_000, "growth is reclaimed: {}", pool.used());
5061 }
5062
5063 #[tokio::test]
5066 async fn refilled_sequence_stays_out_of_subscriptions() {
5067 let (mut producer, _pool) = pooled_producer(1 << 40);
5068 let dynamic = producer.dynamic();
5069 let consumer = producer.consume();
5070
5071 producer.create_group(0u64.into()).unwrap().finish().unwrap();
5072 let aborted = producer.create_group(1u64.into()).unwrap();
5073 aborted.abort(Error::Cancel).unwrap();
5074 producer.create_group(2u64.into()).unwrap().finish().unwrap();
5075
5076 let pending = consumer.fetch_group(1, None);
5079 let req = dynamic
5080 .requested_group()
5081 .now_or_never()
5082 .expect("should not block")
5083 .unwrap();
5084 let mut group = req.accept(None).unwrap();
5085 group
5086 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"backfill"))
5087 .unwrap();
5088 group.finish().unwrap();
5089 pending.await.unwrap();
5090
5091 assert!(consumer.peek_group(1).is_some());
5093 let mut subscriber = producer.subscribe(None);
5094 assert_eq!(subscriber.assert_group().sequence, 0);
5095 assert_eq!(subscriber.assert_group().sequence, 2);
5096 subscriber.assert_no_group();
5097 }
5098
5099 #[tokio::test]
5102 async fn expired_backfill_behind_refreshed_reclaimed() {
5103 tokio::time::pause();
5104
5105 let (mut producer, _pool) = pooled_producer(1 << 40);
5106 let dynamic = producer.dynamic();
5107 let consumer = producer.consume();
5108
5109 producer.create_group(5u64.into()).unwrap().finish().unwrap();
5110 for sequence in [2u64, 3u64] {
5111 let pending = consumer.fetch_group(sequence, None);
5112 let req = dynamic
5113 .requested_group()
5114 .now_or_never()
5115 .expect("should not block")
5116 .unwrap();
5117 let mut group = req.accept(None).unwrap();
5118 group
5119 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 100]))
5120 .unwrap();
5121 group.finish().unwrap();
5122 pending.await.unwrap();
5123 }
5124
5125 tokio::time::advance(Duration::from_secs(4)).await;
5127 consumer.fetch_group(2, None).await.unwrap();
5128 tokio::time::advance(DEFAULT_LATENCY_MAX - Duration::from_secs(2)).await;
5129 producer.create_group(6u64.into()).unwrap().finish().unwrap();
5130
5131 let consumer = producer.consume();
5132 assert!(consumer.peek_group(2).is_some(), "refreshed backfill survives");
5133 assert!(consumer.peek_group(3).is_none(), "expired backfill is reclaimed");
5134 }
5135
5136 #[tokio::test]
5139 async fn same_tick_fetch_protects() {
5140 tokio::time::pause();
5141
5142 let (mut producer, _pool) = pooled_producer(10_000);
5144 let consumer = producer.consume();
5145
5146 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();
5151
5152 finished_group(&mut producer, 3_000); finished_group(&mut producer, 3_000); assert!(consumer.peek_group(0).is_some(), "same-tick refresh protects");
5156 assert!(consumer.peek_group(1).is_none(), "the unread group dies instead");
5157 }
5158
5159 #[tokio::test]
5163 async fn refetched_latest_stays_protected() {
5164 tokio::time::pause();
5165
5166 let (mut producer, _pool) = pooled_producer(10_000);
5167 let dynamic = producer.dynamic();
5168 let consumer = producer.consume();
5169
5170 let straggler = producer.append_group().unwrap(); let latest = producer.append_group().unwrap(); latest.abort(Error::Cancel).unwrap();
5175
5176 let pending = consumer.fetch_group(1, None);
5178 let req = dynamic
5179 .requested_group()
5180 .now_or_never()
5181 .expect("should not block")
5182 .unwrap();
5183 let mut group = req.accept(None).unwrap();
5184 group
5185 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 1000]))
5186 .unwrap();
5187 group.finish().unwrap();
5188 pending.await.unwrap();
5189
5190 {
5193 let state = producer.state.read();
5194 assert!(state.lookup.contains_key(&1), "refetched group is cached");
5195 assert!(
5196 state.evict.iter().all(|(sequence, _)| *sequence != 1),
5197 "the live edge must not be an eviction candidate"
5198 );
5199 }
5200 drop(straggler);
5201 }
5202
5203 #[tokio::test]
5206 async fn eviction_allows_refetch() {
5207 tokio::time::pause();
5208
5209 let (mut producer, _pool) = pooled_producer(10_000);
5210 let dynamic = producer.dynamic();
5211
5212 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let consumer = producer.consume();
5217 assert!(consumer.peek_group(0).is_none());
5218 let pending = consumer.fetch_group(0, None);
5219
5220 let req = dynamic
5221 .requested_group()
5222 .now_or_never()
5223 .expect("should not block")
5224 .unwrap();
5225 assert_eq!(req.sequence(), 0);
5226
5227 let mut group = req.accept(None).unwrap();
5228 group
5229 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"refetched"))
5230 .unwrap();
5231 group.finish().unwrap();
5232
5233 let mut group = pending.await.unwrap();
5234 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"refetched");
5235 }
5236
5237 #[tokio::test]
5240 async fn fetched_backfill_not_subscribed() {
5241 let (mut producer, _pool) = pooled_producer(1 << 40);
5242 let dynamic = producer.dynamic();
5243 let consumer = producer.consume();
5244
5245 producer.create_group(5u64.into()).unwrap().finish().unwrap();
5247 producer.create_group(6u64.into()).unwrap().finish().unwrap();
5248
5249 let pending = consumer.fetch_group(2, None);
5251 let req = dynamic
5252 .requested_group()
5253 .now_or_never()
5254 .expect("should not block")
5255 .unwrap();
5256 let mut group = req.accept(None).unwrap();
5257 group
5258 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"backfill"))
5259 .unwrap();
5260 group.finish().unwrap();
5261 let mut fetched = pending.await.unwrap();
5262 assert_eq!(&fetched.read_frame().await.unwrap().unwrap().payload[..], b"backfill");
5263 assert!(consumer.peek_group(2).is_some(), "backfill is cached for later fetches");
5264
5265 let mut subscriber = producer.subscribe(None);
5267 assert_eq!(subscriber.assert_group().sequence, 5);
5268 assert_eq!(subscriber.assert_group().sequence, 6);
5269 subscriber.assert_no_group();
5270 }
5271
5272 #[tokio::test]
5275 async fn expired_backfill_reclaimed() {
5276 tokio::time::pause();
5277
5278 let (mut producer, pool) = pooled_producer(1 << 40);
5279 let dynamic = producer.dynamic();
5280 let consumer = producer.consume();
5281
5282 producer.create_group(5u64.into()).unwrap().finish().unwrap();
5283
5284 let pending = consumer.fetch_group(2, None);
5286 let req = dynamic
5287 .requested_group()
5288 .now_or_never()
5289 .expect("should not block")
5290 .unwrap();
5291 let mut group = req.accept(None).unwrap();
5292 group
5293 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 1000]))
5294 .unwrap();
5295 group.finish().unwrap();
5296 pending.await.unwrap();
5297 let used = pool.used();
5298
5299 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
5301 producer.create_group(6u64.into()).unwrap().finish().unwrap();
5302
5303 assert!(consumer.peek_group(2).is_none(), "expired backfill is reclaimed");
5304 assert!(pool.used() < used, "its bytes are released");
5305 }
5306
5307 #[tokio::test]
5308 async fn fetch_aborts_with_track() {
5309 let producer = track_producer("test", None);
5310 let dynamic = producer.dynamic();
5311 let consumer = producer.consume();
5312
5313 let pending = consumer.fetch_group(3, None);
5314 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
5315
5316 producer.abort(Error::Cancel).unwrap();
5317 assert!(pending.await.is_err());
5318 drop(dynamic);
5319 }
5320}