1use crate::{Error, Result, Timescale, Timestamp, coding};
17use crate::{broadcast, cache, frame, group, stats};
18
19use super::{Datagram, Requests};
20
21pub use super::subscription::Subscription;
22
23use std::{
24 collections::{HashMap, VecDeque},
25 sync::Arc,
26 sync::OnceLock,
27 sync::atomic::{AtomicBool, Ordering},
28 task::{Poll, ready},
29 time::Duration,
30};
31
32pub const DEFAULT_LATENCY_MAX: Duration = Duration::from_secs(5);
34
35const MAX_DATAGRAM_AGE: Duration = Duration::from_millis(50);
41
42const EVICT_SLACK: usize = 64;
45
46const EVICT_SCAN: usize = 4;
50
51#[derive(Clone, Debug)]
62#[non_exhaustive]
63pub struct Info {
64 pub timescale: Timescale,
71 pub latency_max: Duration,
80 pub priority: u8,
83 pub ordered: bool,
87}
88
89impl Default for Info {
90 fn default() -> Self {
91 Self {
92 timescale: Timescale::default(),
93 latency_max: DEFAULT_LATENCY_MAX,
94 priority: 0,
95 ordered: false,
96 }
97 }
98}
99
100impl Info {
101 pub fn with_timescale(mut self, timescale: Timescale) -> Self {
106 self.timescale = timescale;
107 self
108 }
109
110 pub fn with_latency_max(mut self, latency_max: Duration) -> Self {
112 self.latency_max = latency_max;
113 self
114 }
115
116 pub fn with_priority(mut self, priority: u8) -> Self {
118 self.priority = priority;
119 self
120 }
121
122 pub fn with_ordered(mut self, ordered: bool) -> Self {
126 self.ordered = ordered;
127 self
128 }
129}
130
131#[derive(Default)]
132pub(crate) struct TrackState {
133 info: Option<Info>,
136
137 broadcast: Arc<broadcast::Info>,
140
141 cache: Arc<cache::Track>,
145
146 lookup: HashMap<u64, Slot>,
150
151 arrival: VecDeque<(u64, u32)>,
156
157 evict: VecDeque<(u64, u32)>,
165
166 debt: u64,
171
172 datagrams: VecDeque<(Datagram, web_async::time::Instant)>,
176
177 datagram_offset: usize,
180
181 offset: usize,
184
185 max_sequence: Option<u64>,
188
189 latest_group: Option<u64>,
194
195 next_stamp: u32,
197
198 expire_cursor: usize,
201
202 final_sequence: Option<u64>,
204
205 abort: Option<Error>,
207
208 subscriptions: kio::Shared<Subscriptions>,
212
213 fetch: kio::Shared<FetchState>,
216}
217
218struct Slot {
224 group: group::Producer,
225
226 stamp: u32,
231}
232
233type Subscriptions = Vec<kio::Consumer<Subscription>>;
235
236type FetchState = Requests<u64, PendingFetch>;
241
242struct PendingFetch {
244 priority: u8,
246
247 result: kio::Producer<FetchOutcome>,
252}
253
254#[derive(Default)]
257struct FetchOutcome {
258 rejected: Option<Error>,
259}
260
261impl TrackState {
262 fn poll_info(&self) -> Poll<Result<Info>> {
263 if let Some(info) = &self.info {
264 Poll::Ready(Ok(info.clone()))
265 } else {
266 Poll::Pending
267 }
268 }
269
270 fn poll_recv_group(&self, index: usize, min_sequence: u64) -> Poll<Result<Option<(group::Consumer, usize)>>> {
274 let start = index.saturating_sub(self.offset);
275 for (i, (sequence, stamp)) in self.arrival.iter().enumerate().skip(start) {
276 if *sequence >= min_sequence
277 && let Some(slot) = self.lookup.get(sequence)
278 && slot.stamp == *stamp
279 && !slot.group.is_aborted()
280 {
281 return Poll::Ready(Ok(Some((slot.group.consume(), self.offset + i))));
282 }
283 }
284
285 if self.is_complete() {
287 Poll::Ready(Ok(None))
288 } else if let Some(err) = &self.abort {
289 Poll::Ready(Err(err.clone()))
290 } else {
291 Poll::Pending
292 }
293 }
294
295 fn poll_recv_datagram(&self, index: usize) -> Poll<Result<Option<(Datagram, usize)>>> {
301 let start = index.saturating_sub(self.datagram_offset);
302 if let Some((datagram, _)) = self.datagrams.get(start) {
303 return Poll::Ready(Ok(Some((datagram.clone(), self.datagram_offset + start))));
304 }
305
306 if self.is_complete() {
308 Poll::Ready(Ok(None))
309 } else if let Some(err) = &self.abort {
310 Poll::Ready(Err(err.clone()))
311 } else {
312 Poll::Pending
313 }
314 }
315
316 fn push_datagram(&mut self, datagram: Datagram) {
318 let now = web_async::time::Instant::now();
319 self.datagrams.push_back((datagram, now));
320 while let Some((_, at)) = self.datagrams.front() {
321 if now.duration_since(*at) <= MAX_DATAGRAM_AGE {
322 break;
323 }
324 self.datagrams.pop_front();
325 self.datagram_offset += 1;
326 }
327 }
328
329 fn poll_read_frame(
333 &self,
334 index: usize,
335 next_sequence: u64,
336 waiter: &kio::Waiter,
337 ) -> Poll<Result<Option<(frame::Frame, usize, u64)>>> {
338 let start = index.saturating_sub(self.offset);
339 let mut pending_seen = false;
340 for (i, (sequence, stamp)) in self.arrival.iter().enumerate().skip(start) {
341 if *sequence < next_sequence {
342 continue;
343 }
344 let Some(slot) = self.lookup.get(sequence) else {
345 continue;
346 };
347 if slot.stamp != *stamp {
348 continue;
351 }
352
353 let mut consumer = slot.group.consume();
354 match consumer.poll_read_frame(waiter) {
355 Poll::Ready(Ok(Some(frame))) => {
356 return Poll::Ready(Ok(Some((frame, self.offset + i, *sequence))));
357 }
358 Poll::Ready(Ok(None)) => continue,
359 Poll::Ready(Err(_)) => continue,
362 Poll::Pending => {
363 pending_seen = true;
364 continue;
365 }
366 }
367 }
368
369 if pending_seen {
372 Poll::Pending
373 } else if self.is_complete() {
374 Poll::Ready(Ok(None))
375 } else if let Some(err) = &self.abort {
376 Poll::Ready(Err(err.clone()))
377 } else {
378 Poll::Pending
379 }
380 }
381
382 fn poll_next_in_range(
392 &self,
393 next_sequence: u64,
394 end_sequence: Option<u64>,
395 ) -> Poll<Result<Option<group::Consumer>>> {
396 if let Some(end) = end_sequence
400 && end < next_sequence
401 {
402 if let Some(err) = &self.abort {
403 return Poll::Ready(Err(err.clone()));
404 }
405 return Poll::Pending;
406 }
407
408 let mut best: Option<&group::Producer> = None;
409 for slot in self.lookup.values() {
410 let group = &slot.group;
411 if group.sequence < next_sequence {
412 continue;
413 }
414 if let Some(end) = end_sequence
415 && group.sequence > end
416 {
417 continue;
418 }
419 if group.is_aborted() {
420 continue;
421 }
422 if best.is_none_or(|b| group.sequence < b.sequence) {
423 best = Some(group);
424 }
425 }
426
427 if let Some(group) = best {
428 return Poll::Ready(Ok(Some(group.consume())));
429 }
430
431 if let Some(err) = &self.abort {
433 return Poll::Ready(Err(err.clone()));
434 }
435 if let Some(fin) = self.final_sequence
438 && next_sequence >= fin
439 {
440 return Poll::Ready(Ok(None));
441 }
442 Poll::Pending
443 }
444
445 #[cfg(test)]
450 fn cached_group(&self, sequence: u64) -> Option<group::Consumer> {
451 let slot = self.lookup.get(&sequence)?;
452 if slot.group.is_aborted() {
453 return None;
454 }
455 Some(slot.group.consume())
456 }
457
458 fn latency_bound(&self) -> Option<Duration> {
461 self.info.as_ref().map(|info| info.latency_max)
462 }
463
464 fn poll_fetch_cached(&self, sequence: u64) -> Poll<Result<group::Consumer>> {
469 if let Some(slot) = self.lookup.get(&sequence)
470 && !slot.group.is_aborted()
471 {
472 slot.group.cache_refresh();
476 return Poll::Ready(Ok(slot.group.consume()));
477 }
478
479 if let Some(err) = &self.abort {
480 return Poll::Ready(Err(err.clone()));
481 }
482
483 if self.final_sequence.is_some_and(|fin| sequence >= fin) {
485 return Poll::Ready(Err(Error::NotFound));
486 }
487
488 Poll::Pending
489 }
490
491 fn evict_expired(&mut self, max_age: Duration) {
500 let now = self.cache.pool().now();
501 let max_ticks = cache::Pool::ticks(max_age);
502
503 let len = self.evict.len();
504 if len > 0 {
505 let start = self.expire_cursor % len;
506 for step in 0..len.min(EVICT_SCAN) {
507 let (sequence, stamp) = self.evict[(start + step) % len];
508 let Some(slot) = self.lookup.get(&sequence) else {
509 continue;
510 };
511 if slot.stamp != stamp {
512 continue;
514 }
515 if slot.group.is_aborted() {
518 self.lookup.remove(&sequence);
519 continue;
520 }
521 if Some(sequence) == self.latest_group || now.saturating_sub(slot.group.cache_accessed()) <= max_ticks {
522 continue;
523 }
524 let slot = self.lookup.remove(&sequence).unwrap();
528 let _ = slot.group.abort(Error::Old);
529 }
530 self.expire_cursor = (start + EVICT_SCAN) % len;
531 }
532
533 while let Some((sequence, stamp)) = self.arrival.front() {
536 if self.lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp) {
537 break;
538 }
539 self.arrival.pop_front();
540 self.offset += 1;
541 }
542
543 while let Some((sequence, stamp)) = self.evict.front() {
545 if self.lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp) {
546 break;
547 }
548 self.evict.pop_front();
549 }
550
551 if self.evict.len() > 2 * self.lookup.len() + EVICT_SLACK {
554 let lookup = &self.lookup;
555 self.evict
556 .retain(|(sequence, stamp)| lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp));
557 }
558 }
559
560 fn clear_cache(&mut self) {
563 self.lookup.clear();
564 self.arrival.clear();
565 self.evict.clear();
566 self.latest_group = None;
567 self.debt = 0;
568 }
569
570 fn install(&mut self, mut info: Info) {
576 info.latency_max = info.latency_max.min(self.broadcast.origin.cache_duration);
577 self.info = Some(info);
578 }
579
580 fn spawn(broadcast: Arc<broadcast::Info>) -> kio::Producer<Self> {
587 let state = kio::Producer::new(Self {
588 broadcast: broadcast.clone(),
589 ..Default::default()
590 });
591 let cache = cache::Track::new(broadcast.origin.pool.clone(), state.downgrade());
592 state.write().ok().expect("a new track is open").cache = cache;
593 state
594 }
595
596 fn claim_sequence(&mut self, sequence: u64) -> Result<()> {
602 if let Some(slot) = self.lookup.get(&sequence) {
603 if !slot.group.is_aborted() {
604 return Err(Error::Duplicate);
605 }
606 self.lookup.remove(&sequence);
607 }
608 Ok(())
609 }
610
611 fn insert_group(&mut self, group: &group::Producer, visible: bool) {
618 let sequence = group.sequence;
619 self.next_stamp = self.next_stamp.wrapping_add(1);
620 let stamp = self.next_stamp;
621
622 if self.latest_group.is_none_or(|latest| sequence >= latest) {
626 if let Some(latest) = self.latest_group
629 && sequence > latest
630 && let Some(prev) = self.lookup.get(&latest)
631 {
632 prev.group.cache_demote();
633 self.evict.push_back((latest, prev.stamp));
634 }
635 self.latest_group = Some(sequence);
636 } else {
637 group.cache_demote();
638 self.evict.push_back((sequence, stamp));
639 }
640
641 self.max_sequence = Some(self.max_sequence.map_or(sequence, |max| max.max(sequence)));
642 self.lookup.insert(
643 sequence,
644 Slot {
645 group: group.clone(),
646 stamp,
647 },
648 );
649 if visible {
650 self.arrival.push_back((sequence, stamp));
651 }
652 }
653
654 fn commit_group(&mut self, group: &group::Producer, visible: bool, latency_max: Duration) {
658 self.charge_debt();
659 self.insert_group(group, visible);
660 self.evict_expired(latency_max);
661 }
662
663 pub(super) fn charge_debt(&mut self) {
676 let written = self.cache.take_written();
677 let pool = self.cache.pool().clone();
678 match pool.accrue(written) {
679 Some(mut accrued) => {
680 if self.oldest_is_stale(&pool) {
681 accrued = accrued.saturating_mul(2);
682 }
683 self.debt = self.debt.saturating_add(accrued).min(pool.used());
686 self.pay_debt(&pool, written.saturating_mul(2));
689 }
690 None => self.debt = 0,
693 }
694 }
695
696 fn oldest_is_stale(&self, pool: &cache::Pool) -> bool {
700 let Some(average) = pool.average() else {
701 return false;
702 };
703 let Some((sequence, stamp)) = self.evict.front() else {
704 return false;
705 };
706 let Some(slot) = self.lookup.get(sequence) else {
707 return false;
708 };
709 slot.stamp == *stamp && !slot.group.is_aborted() && slot.group.cache_accessed() <= average
710 }
711
712 fn pay_debt(&mut self, pool: &cache::Pool, cap: u64) {
725 let average = pool.average().unwrap_or(0);
726 let mut paid = 0u64;
727 let mut scanned = 0usize;
728 for _ in 0..self.evict.len() {
729 if self.debt == 0 || paid >= cap || scanned >= EVICT_SCAN {
730 return;
731 }
732 let Some((sequence, stamp)) = self.evict.pop_front() else {
733 return;
734 };
735 let Some(slot) = self.lookup.get(&sequence) else {
736 continue;
738 };
739 if slot.stamp != stamp {
740 continue;
742 }
743 if slot.group.is_aborted() {
744 self.lookup.remove(&sequence);
746 continue;
747 }
748 if Some(sequence) == self.latest_group {
749 self.evict.push_back((sequence, stamp));
751 continue;
752 }
753
754 scanned += 1;
755 if slot.group.cache_accessed() > average {
759 self.evict.push_back((sequence, stamp));
760 continue;
761 }
762 let size = slot.group.cache_size();
765 if size > self.debt {
766 self.evict.push_front((sequence, stamp));
767 return;
768 }
769
770 self.debt -= size;
771 paid = paid.saturating_add(size);
772 let slot = self.lookup.remove(&sequence).unwrap();
773 let _ = slot.group.abort(Error::Evicted);
774 }
775 }
776
777 fn set_final(&mut self, final_sequence: u64) -> Result<()> {
780 if self.final_sequence.is_some() {
781 return Err(Error::Closed);
782 }
783 if let Some(max) = self.max_sequence
784 && final_sequence <= max
785 {
786 return Err(Error::ProtocolViolation);
787 }
788 self.final_sequence = Some(final_sequence);
789 Ok(())
790 }
791
792 fn is_complete(&self) -> bool {
798 self.final_sequence
799 .is_some_and(|fin| self.max_sequence.map_or(0, |max| max.saturating_add(1)) >= fin)
800 }
801
802 fn poll_finished(&self) -> Poll<Result<u64>> {
803 if let Some(fin) = self.final_sequence {
804 Poll::Ready(Ok(fin))
805 } else if let Some(err) = &self.abort {
806 Poll::Ready(Err(err.clone()))
807 } else {
808 Poll::Pending
809 }
810 }
811
812 fn modify(producer: &kio::Producer<Self>) -> Result<kio::Mut<'_, Self>> {
813 producer.write().map_err(|r| r.abort.clone().unwrap_or(Error::Dropped))
814 }
815
816 fn insert_group_request(&mut self, sequence: u64, info: Option<Info>) -> Result<group::Producer> {
822 if let Some(err) = &self.abort {
823 return Err(err.clone());
824 }
825 if let Some(fin) = self.final_sequence
826 && sequence >= fin
827 {
828 return Err(Error::Closed);
829 }
830
831 if self.info.is_none() {
835 self.install(info.unwrap_or_default());
836 }
837 let info = self.info.clone().unwrap();
838
839 self.claim_sequence(sequence)?;
841
842 let latency_max = info.latency_max;
843 let group = group::Producer::new(group::Info { sequence }, info, self.cache.clone());
844 group.cache_refresh();
849 self.commit_group(&group, false, latency_max);
850 Ok(group)
851 }
852}
853
854#[derive(Clone)]
856pub struct Producer {
857 name: Arc<str>,
858 broadcast: Arc<broadcast::Info>,
861 state: kio::Producer<TrackState>,
862 prev_subscription: Option<Subscription>,
863 alive: Arc<Alive>,
865 stats: stats::Scope,
869}
870
871impl Producer {
872 pub(crate) fn new(
880 broadcast: Arc<broadcast::Info>,
881 name: impl Into<Arc<str>>,
882 info: impl Into<Option<Info>>,
883 ) -> Self {
884 let name = name.into();
885 let state = TrackState::spawn(broadcast.clone());
886 state
887 .write()
888 .ok()
889 .expect("a new track is open")
890 .install(info.into().unwrap_or_default());
891 let alive = Alive::new(name.clone(), state.clone());
892 alive.publish(None);
893 Self {
894 name,
895 state,
896 broadcast,
897 prev_subscription: None,
898 alive,
899 stats: stats::Scope::default(),
900 }
901 }
902
903 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
907 self.alive.publish(Some(&scope));
908 self.stats = scope;
909 self
910 }
911
912 pub fn name(&self) -> &str {
914 &self.name
915 }
916
917 pub fn broadcast(&self) -> &broadcast::Info {
919 &self.broadcast
920 }
921
922 pub fn create_group(&mut self, group: group::Info) -> Result<group::Producer> {
924 let mut state = self.modify()?;
925 if let Some(fin) = state.final_sequence
926 && group.sequence >= fin
927 {
928 return Err(Error::Closed);
929 }
930 let track = state.info.clone().unwrap();
931 let latency_max = track.latency_max;
932
933 state.claim_sequence(group.sequence)?;
935
936 let group = group::Producer::new(group, track, state.cache.clone()).with_meter(self.stats.meter());
937 state.commit_group(&group, true, latency_max);
938
939 Ok(group)
940 }
941
942 pub fn append_group(&mut self) -> Result<group::Producer> {
944 let mut state = self.modify()?;
945 let sequence = match state.max_sequence {
946 Some(s) => s.checked_add(1).ok_or(coding::BoundsExceeded)?,
947 None => 0,
948 };
949 if let Some(fin) = state.final_sequence
950 && sequence >= fin
951 {
952 return Err(Error::Closed);
953 }
954
955 let track = state.info.clone().unwrap();
956 let latency_max = track.latency_max;
957
958 let group =
959 group::Producer::new(group::Info { sequence }, track, state.cache.clone()).with_meter(self.stats.meter());
960 state.commit_group(&group, true, latency_max);
961
962 Ok(group)
963 }
964
965 pub fn append_datagram<B: crate::IntoBytes>(&mut self, timestamp: Timestamp, payload: B) -> Result<u64> {
977 let payload = payload.into_bytes();
978 if payload.len() > super::datagram::MAX_DATAGRAM_PAYLOAD {
979 return Err(Error::FrameTooLarge);
980 }
981 let meter = self.stats.meter();
983 let mut state = self.modify()?;
984 let timescale = state.info.as_ref().unwrap().timescale;
986 let timestamp = timestamp.convert(timescale).map_err(|_| Error::TimestampMismatch)?;
987 let sequence = match state.max_sequence {
988 Some(s) => s.checked_add(1).ok_or(coding::BoundsExceeded)?,
989 None => 0,
990 };
991 if let Some(fin) = state.final_sequence
992 && sequence >= fin
993 {
994 return Err(Error::Closed);
995 }
996 state.max_sequence = Some(sequence);
997 meter.datagram(payload.len() as u64);
998 state.push_datagram(Datagram {
999 sequence,
1000 timestamp,
1001 payload,
1002 });
1003 Ok(sequence)
1004 }
1005
1006 pub fn write_datagram(&mut self, mut datagram: Datagram) -> Result<()> {
1012 if datagram.payload.len() > super::datagram::MAX_DATAGRAM_PAYLOAD {
1013 return Err(Error::FrameTooLarge);
1014 }
1015 let meter = self.stats.meter();
1017 let mut state = self.modify()?;
1018 let timescale = state.info.as_ref().unwrap().timescale;
1020 datagram.timestamp = datagram
1021 .timestamp
1022 .convert(timescale)
1023 .map_err(|_| Error::TimestampMismatch)?;
1024 if let Some(fin) = state.final_sequence
1025 && datagram.sequence >= fin
1026 {
1027 return Err(Error::Closed);
1028 }
1029 state.max_sequence = Some(state.max_sequence.unwrap_or(0).max(datagram.sequence));
1030 meter.datagram(datagram.payload.len() as u64);
1031 state.push_datagram(datagram);
1032 Ok(())
1033 }
1034
1035 pub fn write_frame<B: crate::IntoBytes>(&mut self, timestamp: Timestamp, frame: B) -> Result<()> {
1040 let mut group = self.append_group()?;
1041 group.write_frame(timestamp, frame)?;
1042 group.finish()?;
1043 Ok(())
1044 }
1045
1046 pub fn finish(&mut self) -> Result<()> {
1052 let mut state = self.modify()?;
1053 let final_sequence = match state.max_sequence {
1054 Some(max) => max.checked_add(1).ok_or(coding::BoundsExceeded)?,
1055 None => 0,
1056 };
1057 state.set_final(final_sequence)
1058 }
1059
1060 pub fn finish_at(&mut self, final_sequence: u64) -> Result<()> {
1073 self.modify()?.set_final(final_sequence)
1074 }
1075
1076 pub fn final_sequence(&self) -> Option<u64> {
1081 self.state.read().final_sequence
1082 }
1083
1084 pub fn abort(self, err: Error) -> Result<()> {
1095 let mut guard = self.modify()?;
1096 guard.abort = Some(err);
1097 guard.clear_cache();
1098 guard.datagrams.clear();
1099 guard.close();
1100 Ok(())
1101 }
1102
1103 pub async fn unused(&self) -> Result<()> {
1105 self.state.unused().await.map_err(|_| self.abort_reason())
1106 }
1107
1108 pub async fn used(&self) -> Result<()> {
1110 self.state.used().await.map_err(|_| self.abort_reason())
1111 }
1112
1113 pub async fn closed(&self) -> Error {
1115 kio::wait(|waiter| self.poll_closed(waiter)).await
1116 }
1117
1118 pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<Error> {
1120 self.state.poll_closed(waiter).map(|()| self.abort_reason())
1121 }
1122
1123 fn abort_reason(&self) -> Error {
1125 self.state.read().abort.clone().unwrap_or(Error::Dropped)
1126 }
1127
1128 pub fn is_closed(&self) -> bool {
1130 self.state.read().is_closed()
1131 }
1132
1133 pub fn latest(&self) -> Option<u64> {
1135 self.state.read().max_sequence
1136 }
1137
1138 pub fn is_clone(&self, other: &Self) -> bool {
1140 self.state.same_channel(&other.state)
1141 }
1142
1143 pub(crate) fn weak(&self) -> TrackWeak {
1145 TrackWeak {
1146 name: self.name.clone(),
1147 state: self.state.weak(),
1148 }
1149 }
1150
1151 pub fn demand(&self) -> Demand {
1159 Demand {
1160 name: self.name.clone(),
1161 state: self.state.weak(),
1162 }
1163 }
1164
1165 pub fn consume(&self) -> Consumer {
1170 Consumer::plain(self.name.clone(), self.state.consume())
1171 }
1172
1173 pub fn subscribe(&self, subscription: impl Into<Option<Subscription>>) -> Subscriber {
1178 let preferences = subscription.into().unwrap_or_default();
1179
1180 let info = self.state.read().info.clone().expect("producer always has info");
1185 let subscription = kio::Producer::new(preferences);
1186 register_subscription(self.state.read(), &subscription);
1187
1188 Subscriber {
1189 name: self.name.clone(),
1190 info,
1191 inner: SubscriberKind::Plain(PlainSubscriber {
1192 state: self.state.consume(),
1193 subscription,
1194 index: 0,
1195 datagram_index: 0,
1196 min_sequence: 0,
1197 next_sequence: 0,
1198 end_sequence: None,
1199 }),
1200 stats: stats::Scope::default(),
1202 _stats_sub: stats::Subscription::default(),
1203 }
1204 }
1205
1206 pub async fn subscription_changed(&mut self) -> Result<Option<Subscription>> {
1212 kio::wait(|waiter| self.poll_subscription_changed(waiter)).await
1213 }
1214
1215 pub fn subscription(&self) -> Option<Subscription> {
1223 let state = self.state.read();
1224 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
1225 drop(state);
1226 snapshot_subscription(&subs, bound)
1227 }
1228
1229 pub fn poll_subscription_changed(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Subscription>>> {
1233 if self.state.poll_closed(waiter).is_ready() {
1236 let abort = self.state.read().abort.clone();
1237 return Poll::Ready(Err(abort.unwrap_or(Error::Dropped)));
1238 }
1239
1240 let state = self.state.read();
1242 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
1243 drop(state);
1244
1245 let prev = &self.prev_subscription;
1246 let mut combined = None;
1247 let mut guard = match subs.poll(waiter, |subs| {
1248 let next = combined_subscription(subs, bound, waiter);
1249 if &next == prev {
1250 Poll::Pending
1251 } else {
1252 combined = next;
1253 Poll::Ready(())
1254 }
1255 }) {
1256 Poll::Ready(guard) => guard,
1257 Poll::Pending => return Poll::Pending,
1258 };
1259 guard.retain(|sub| !sub.is_closed());
1261 drop(guard);
1262 self.prev_subscription = combined.clone();
1263 Poll::Ready(Ok(combined))
1264 }
1265
1266 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
1268 self.state.poll_unused(waiter).map(|_| ())
1269 }
1270
1271 pub fn dynamic(&self) -> Dynamic {
1275 Dynamic::new(self.name.clone(), self.state.clone(), self.alive.clone())
1276 }
1277
1278 fn modify(&self) -> Result<kio::Mut<'_, TrackState>> {
1279 TrackState::modify(&self.state)
1280 }
1281}
1282
1283fn poll_requested_group(
1287 state: &kio::Producer<TrackState>,
1288 fetch: &kio::Shared<FetchState>,
1289 waiter: &kio::Waiter,
1290) -> Poll<Result<GroupRequest>> {
1291 if let Poll::Ready(mut guard) = fetch.poll(waiter, |fetch| {
1293 if fetch.has_queued() {
1294 Poll::Ready(())
1295 } else {
1296 Poll::Pending
1297 }
1298 }) {
1299 let sequence = guard.pop().expect("predicate guaranteed a request");
1300 let pending = guard.get(&sequence).expect("popped key must be pending");
1304 let priority = pending.priority;
1305 let result = pending.result.clone();
1306 drop(guard);
1307 return Poll::Ready(Ok(GroupRequest {
1308 state: state.clone(),
1309 fetch: fetch.clone(),
1310 sequence,
1311 priority,
1312 result,
1313 done: false,
1314 }));
1315 }
1316
1317 match state.poll_ref(waiter, |state| match &state.abort {
1319 Some(err) => Poll::Ready(err.clone()),
1320 None => Poll::Pending,
1321 }) {
1322 Poll::Ready(Ok(err)) => Poll::Ready(Err(err)),
1323 Poll::Ready(Err(closed)) => Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped))),
1324 Poll::Pending => Poll::Pending,
1325 }
1326}
1327
1328pub struct Dynamic {
1338 name: Arc<str>,
1339 state: kio::Producer<TrackState>,
1341 fetch: kio::Shared<FetchState>,
1343 alive: Arc<Alive>,
1346}
1347
1348impl Dynamic {
1349 fn new(name: Arc<str>, state: kio::Producer<TrackState>, alive: Arc<Alive>) -> Self {
1350 let fetch = state.read().fetch.clone();
1351 fetch.lock().add_handler();
1352 Self {
1353 name,
1354 state,
1355 fetch,
1356 alive,
1357 }
1358 }
1359
1360 pub fn name(&self) -> &str {
1362 &self.name
1363 }
1364
1365 pub async fn requested_group(&self) -> Result<GroupRequest> {
1371 kio::wait(|waiter| self.poll_requested_group(waiter)).await
1372 }
1373
1374 pub fn poll_requested_group(&self, waiter: &kio::Waiter) -> Poll<Result<GroupRequest>> {
1376 poll_requested_group(&self.state, &self.fetch, waiter)
1377 }
1378
1379 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
1381 self.state.poll_unused(waiter).map(|_| ())
1382 }
1383}
1384
1385impl Clone for Dynamic {
1386 fn clone(&self) -> Self {
1387 self.fetch.lock().add_handler();
1389 Self {
1390 name: self.name.clone(),
1391 state: self.state.clone(),
1392 fetch: self.fetch.clone(),
1393 alive: self.alive.clone(),
1394 }
1395 }
1396}
1397
1398impl Drop for Dynamic {
1399 fn drop(&mut self) {
1400 let mut fetch = self.fetch.lock();
1406 if fetch.remove_handler() {
1407 fetch.drain_queued();
1408 }
1409 }
1410}
1411
1412struct Alive {
1421 name: Arc<str>,
1422 state: kio::Producer<TrackState>,
1423
1424 published: AtomicBool,
1427
1428 stats: OnceLock<stats::Subscription>,
1431}
1432
1433impl Alive {
1434 fn new(name: Arc<str>, state: kio::Producer<TrackState>) -> Arc<Self> {
1435 Arc::new(Self {
1436 name,
1437 state,
1438 published: Default::default(),
1439 stats: Default::default(),
1440 })
1441 }
1442
1443 fn publish(&self, stats: Option<&stats::Scope>) {
1447 self.published.store(true, Ordering::Relaxed);
1448 if let Some(scope) = stats {
1449 let _ = self.stats.set(scope.subscribe());
1452 }
1453 }
1454}
1455
1456impl Drop for Alive {
1457 fn drop(&mut self) {
1458 if !self.published.load(Ordering::Relaxed) {
1460 return;
1461 }
1462 match self.state.write() {
1470 Ok(mut state) => {
1471 if state.final_sequence.is_some() || state.abort.is_some() {
1472 return;
1473 }
1474 tracing::warn!(
1475 track = %self.name,
1476 "track::Producer dropped without finish() or abort()"
1477 );
1478 state.clear_cache();
1479 state.datagrams.clear();
1480 }
1481 Err(state) => {
1482 if state.final_sequence.is_some() || state.abort.is_some() {
1483 return;
1484 }
1485 tracing::warn!(
1486 track = %self.name,
1487 "track::Producer dropped without finish() or abort()"
1488 );
1489 }
1490 }
1491 }
1492}
1493
1494fn combined_subscription(subs: &Subscriptions, bound: Option<Duration>, waiter: &kio::Waiter) -> Option<Subscription> {
1500 let mut combined = None;
1501 for sub in subs.iter() {
1502 if sub.is_closed() {
1507 continue;
1508 }
1509 let _ = sub.poll_closed(waiter);
1514 if let Poll::Ready(Ok(sub)) = sub.poll(waiter, |sub| sub.poll_combined(&combined)) {
1515 combined = Some(sub);
1516 }
1517 }
1518 clamp_combined(combined, bound)
1519}
1520
1521fn snapshot_subscription(subs: &kio::Shared<Subscriptions>, bound: Option<Duration>) -> Option<Subscription> {
1523 let mut combined: Option<Subscription> = None;
1524 for sub in subs.read().iter() {
1525 if sub.is_closed() {
1527 continue;
1528 }
1529 if let Poll::Ready(merged) = sub.read().poll_combined(&combined) {
1530 combined = Some(merged);
1531 }
1532 }
1533 clamp_combined(combined, bound)
1534}
1535
1536fn clamp_combined(combined: Option<Subscription>, bound: Option<Duration>) -> Option<Subscription> {
1544 let mut combined = combined?;
1545 if let Some(bound) = bound {
1546 combined.latency_max = combined.latency_max.min(bound);
1547 }
1548 Some(combined)
1549}
1550
1551fn register_subscription(state: kio::Ref<'_, TrackState>, subscription: &kio::Producer<Subscription>) {
1555 if state.is_closed() {
1556 return;
1557 }
1558 let subs = state.subscriptions.clone();
1559 drop(state);
1560 subs.lock().push(subscription.consume());
1561}
1562
1563#[derive(Clone)]
1565pub(crate) struct TrackWeak {
1566 name: Arc<str>,
1567 state: kio::ProducerWeak<TrackState>,
1568}
1569
1570impl TrackWeak {
1571 pub fn consume(&self) -> Consumer {
1572 Consumer::plain(self.name.clone(), self.state.consume())
1573 }
1574
1575 pub(crate) fn name(&self) -> &Arc<str> {
1578 &self.name
1579 }
1580
1581 pub(crate) fn is_used(&self) -> bool {
1584 !self.state.is_closed() && self.state.is_used()
1585 }
1586
1587 pub(crate) fn poll_used(&self, waiter: &kio::Waiter) {
1590 let _ = self.state.poll_used(waiter);
1591 }
1592
1593 pub(crate) fn poll_unused(&self, waiter: &kio::Waiter) {
1596 let _ = self.state.poll_unused(waiter);
1597 }
1598}
1599
1600impl super::WeakEntry for TrackWeak {
1601 fn is_closed(&self) -> bool {
1602 self.state.is_closed()
1603 }
1604
1605 fn same_channel(&self, other: &Self) -> bool {
1606 self.state.same_channel(&other.state)
1607 }
1608}
1609
1610#[derive(Clone)]
1619pub struct Demand {
1620 name: Arc<str>,
1621 state: kio::ProducerWeak<TrackState>,
1622}
1623
1624impl Demand {
1625 pub fn name(&self) -> &str {
1627 &self.name
1628 }
1629
1630 pub async fn used(&self) -> Result<()> {
1632 self.state.used().await.map_err(|_| self.abort_reason())
1633 }
1634
1635 pub async fn unused(&self) -> Result<()> {
1637 self.state.unused().await.map_err(|_| self.abort_reason())
1638 }
1639
1640 pub async fn closed(&self) -> Error {
1642 self.state.closed().await;
1643 self.abort_reason()
1644 }
1645
1646 fn abort_reason(&self) -> Error {
1648 self.state.read().abort.clone().unwrap_or(Error::Dropped)
1649 }
1650}
1651
1652#[derive(Clone)]
1663pub struct Consumer {
1664 name: Arc<str>,
1665 inner: ConsumerKind,
1666 stats: stats::Scope,
1669}
1670
1671#[derive(Clone)]
1672enum ConsumerKind {
1673 Plain(kio::Consumer<TrackState>),
1674 Spliced(super::resume::Consumer),
1675}
1676
1677impl Consumer {
1678 fn plain(name: Arc<str>, state: kio::Consumer<TrackState>) -> Self {
1679 Self {
1680 name,
1681 inner: ConsumerKind::Plain(state),
1682 stats: stats::Scope::default(),
1683 }
1684 }
1685
1686 pub(crate) fn spliced(name: Arc<str>, resume: super::resume::Consumer) -> Self {
1688 Self {
1689 name,
1690 inner: ConsumerKind::Spliced(resume),
1691 stats: stats::Scope::default(),
1692 }
1693 }
1694
1695 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
1698 self.stats = scope;
1699 self
1700 }
1701
1702 pub fn name(&self) -> &str {
1704 &self.name
1705 }
1706
1707 pub fn subscribe(&self, subscription: impl Into<Option<Subscription>>) -> kio::Pending<Subscribing> {
1713 let subscription = kio::Producer::new(subscription.into().unwrap_or_default());
1714
1715 let inner = match &self.inner {
1716 ConsumerKind::Plain(state) => {
1717 register_subscription(state.read(), &subscription);
1720 SubscribingKind::Plain(state.clone())
1721 }
1722 ConsumerKind::Spliced(resume) => SubscribingKind::Spliced(resume.clone()),
1724 };
1725
1726 kio::Pending::new(Subscribing {
1727 name: self.name.clone(),
1728 inner,
1729 subscription,
1730 stats: self.stats.clone(),
1731 })
1732 }
1733
1734 #[cfg(test)]
1738 pub(crate) fn peek_group(&self, sequence: u64) -> Option<group::Consumer> {
1739 match &self.inner {
1740 ConsumerKind::Plain(state) => state.read().cached_group(sequence),
1741 ConsumerKind::Spliced(_) => None,
1744 }
1745 }
1746
1747 pub fn fetch_group(&self, sequence: u64, options: impl Into<Option<group::Fetch>>) -> kio::Pending<Fetching> {
1759 let options = options.into().unwrap_or_default();
1760
1761 self.stats.fetch();
1765
1766 let state = match &self.inner {
1767 ConsumerKind::Plain(state) => state,
1768 ConsumerKind::Spliced(resume) => {
1771 return kio::Pending::new(Fetching {
1772 inner: FetchingKind::Spliced(resume.fetch_group(sequence, options)),
1773 stats: self.stats.clone(),
1774 });
1775 }
1776 };
1777
1778 let mut result = None;
1779
1780 let (fetch, unresolved) = {
1784 let state = state.read();
1785 (state.fetch.clone(), state.poll_fetch_cached(sequence).is_pending())
1786 };
1787
1788 if unresolved {
1789 let mut fetch = fetch.lock();
1790 if let Some(pending) = fetch.join(&sequence) {
1791 pending.priority = pending.priority.max(options.priority);
1794 result = Some(pending.result.consume());
1795 } else {
1796 let producer = kio::Producer::<FetchOutcome>::default();
1800 let consumer = producer.consume();
1801 let attempt = PendingFetch {
1802 priority: options.priority,
1803 result: producer,
1804 };
1805 if fetch.insert(sequence, attempt).is_ok() {
1806 result = Some(consumer);
1807 }
1808 }
1809 }
1810
1811 kio::Pending::new(Fetching {
1812 inner: FetchingKind::Plain {
1813 state: state.clone(),
1814 fetch,
1815 sequence,
1816 result,
1817 },
1818 stats: self.stats.clone(),
1819 })
1820 }
1821
1822 pub fn info(&self) -> kio::Pending<Querying> {
1829 kio::Pending::new(Querying {
1830 inner: match &self.inner {
1831 ConsumerKind::Plain(state) => QueryingKind::Plain(state.clone()),
1832 ConsumerKind::Spliced(resume) => QueryingKind::Spliced(resume.clone()),
1833 },
1834 })
1835 }
1836
1837 pub fn latest(&self) -> Option<u64> {
1839 match &self.inner {
1840 ConsumerKind::Plain(state) => state.read().max_sequence,
1841 ConsumerKind::Spliced(resume) => resume.latest(),
1842 }
1843 }
1844
1845 pub(crate) fn poll_complete(&self, waiter: &kio::Waiter) -> Poll<Result<()>> {
1850 let ConsumerKind::Plain(state) = &self.inner else {
1851 return Poll::Pending;
1853 };
1854 match ready!(state.poll(waiter, |state| {
1855 if state.is_complete() {
1856 Poll::Ready(())
1857 } else {
1858 Poll::Pending
1859 }
1860 })) {
1861 Ok(_) => Poll::Ready(Ok(())),
1862 Err(closed) => Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped))),
1865 }
1866 }
1867}
1868
1869pub struct Subscribing {
1872 name: Arc<str>,
1873 inner: SubscribingKind,
1874 subscription: kio::Producer<Subscription>,
1875 stats: stats::Scope,
1876}
1877
1878enum SubscribingKind {
1879 Plain(kio::Consumer<TrackState>),
1880 Spliced(super::resume::Consumer),
1881}
1882
1883impl Subscribing {
1884 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<Subscriber>> {
1887 match &self.inner {
1888 SubscribingKind::Plain(state) => {
1889 let info = ready!(state.poll(waiter, |state| state.poll_info()))
1891 .map_err(|e| e.abort.clone().unwrap_or(Error::Dropped))??;
1892
1893 Poll::Ready(Ok(Subscriber {
1894 name: self.name.clone(),
1895 info,
1896 inner: SubscriberKind::Plain(PlainSubscriber {
1897 state: state.clone(),
1898 subscription: self.subscription.clone(),
1899 index: 0,
1900 datagram_index: 0,
1901 min_sequence: 0,
1902 next_sequence: 0,
1903 end_sequence: None,
1904 }),
1905 stats: self.stats.clone(),
1906 _stats_sub: self.stats.subscribe(),
1907 }))
1908 }
1909 SubscribingKind::Spliced(resume) => {
1910 let info = ready!(resume.poll_info(waiter))?;
1913
1914 Poll::Ready(Ok(Subscriber {
1915 name: self.name.clone(),
1916 info,
1917 inner: SubscriberKind::Spliced(Box::new(resume.subscribe_shared(self.subscription.clone()))),
1918 stats: self.stats.clone(),
1919 _stats_sub: self.stats.subscribe(),
1920 }))
1921 }
1922 }
1923 }
1924
1925 pub fn update(&mut self, subscription: Subscription) -> Result<()> {
1930 let mut state = self.subscription.write().map_err(|_| Error::Closed)?;
1931 *state = subscription;
1932 Ok(())
1933 }
1934}
1935
1936impl kio::Pollable for Subscribing {
1937 type Output = Result<Subscriber>;
1938
1939 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
1940 self.poll_ok(waiter)
1941 }
1942}
1943
1944pub struct Querying {
1947 inner: QueryingKind,
1948}
1949
1950enum QueryingKind {
1951 Plain(kio::Consumer<TrackState>),
1952 Spliced(super::resume::Consumer),
1953}
1954
1955impl Querying {
1956 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<Info>> {
1958 match &self.inner {
1959 QueryingKind::Plain(state) => {
1960 let info = ready!(state.poll(waiter, |state| state.poll_info()))
1962 .map_err(|e| e.abort.clone().unwrap_or(Error::Dropped))??;
1963 Poll::Ready(Ok(info))
1964 }
1965 QueryingKind::Spliced(resume) => resume.poll_info(waiter),
1966 }
1967 }
1968}
1969
1970impl kio::Pollable for Querying {
1971 type Output = Result<Info>;
1972
1973 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
1974 self.poll_ok(waiter)
1975 }
1976}
1977
1978pub struct GroupRequest {
1987 state: kio::Producer<TrackState>,
1988 fetch: kio::Shared<FetchState>,
1990 sequence: u64,
1991 priority: u8,
1992 result: kio::Producer<FetchOutcome>,
1994 done: bool,
1995}
1996
1997impl GroupRequest {
1998 pub fn sequence(&self) -> u64 {
2000 self.sequence
2001 }
2002
2003 pub fn priority(&self) -> u8 {
2005 self.priority
2006 }
2007
2008 pub fn accept(mut self, info: impl Into<Option<Info>>) -> Result<group::Producer> {
2016 self.done = true;
2017 let res = TrackState::modify(&self.state)
2021 .and_then(|mut state| state.insert_group_request(self.sequence, info.into()));
2022 self.remove();
2023 res
2024 }
2025
2026 pub fn reject(mut self, err: Error) {
2028 self.done = true;
2029 self.remove();
2032 if let Ok(mut outcome) = self.result.write() {
2033 outcome.rejected = Some(err);
2034 }
2035 }
2036
2037 fn remove(&self) {
2040 self.fetch
2041 .lock()
2042 .remove_if(&self.sequence, |pending| pending.result.same_channel(&self.result));
2043 }
2044}
2045
2046impl Drop for GroupRequest {
2047 fn drop(&mut self) {
2048 if self.done {
2049 return;
2050 }
2051 self.remove();
2052 if let Ok(mut outcome) = self.result.write() {
2053 outcome.rejected = Some(Error::Dropped);
2054 }
2055 }
2056}
2057
2058pub struct Fetching {
2064 inner: FetchingKind,
2065 stats: stats::Scope,
2068}
2069
2070enum FetchingKind {
2071 Plain {
2072 state: kio::Consumer<TrackState>,
2073 fetch: kio::Shared<FetchState>,
2074 sequence: u64,
2075 result: Option<kio::Consumer<FetchOutcome>>,
2077 },
2078 Spliced(kio::Pending<super::resume::Fetching>),
2080}
2081
2082impl kio::Pollable for Fetching {
2083 type Output = Result<group::Consumer>;
2084
2085 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
2086 let (state, fetch, sequence, result) = match &self.inner {
2087 FetchingKind::Plain {
2088 state,
2089 fetch,
2090 sequence,
2091 result,
2092 } => (state, fetch, *sequence, result.as_ref()),
2093 FetchingKind::Spliced(spliced) => {
2094 return kio::Pollable::poll(&**spliced, waiter)
2097 .map(|res| res.map(|group| group.with_meter(self.stats.meter())));
2098 }
2099 };
2100
2101 match state.poll(waiter, |state| state.poll_fetch_cached(sequence)) {
2104 Poll::Ready(Ok(res)) => return Poll::Ready(res.map(|group| group.with_meter(self.stats.meter()))),
2105 Poll::Ready(Err(closed)) => {
2106 return Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped)));
2107 }
2108 Poll::Pending => {}
2109 }
2110
2111 let Some(result) = result else {
2113 return match fetch.poll(waiter, |fetch| match fetch.has_handlers() {
2116 false => Poll::Ready(()),
2117 true => Poll::Pending,
2118 }) {
2119 Poll::Ready(_guard) => Poll::Ready(Err(Error::NotFound)),
2120 Poll::Pending => Poll::Pending,
2121 };
2122 };
2123
2124 match result.poll(waiter, |outcome| match &outcome.rejected {
2127 Some(err) => Poll::Ready(err.clone()),
2128 None => Poll::Pending,
2129 }) {
2130 Poll::Ready(Ok(err)) => Poll::Ready(Err(err)),
2131 Poll::Ready(Err(_closed)) => Poll::Ready(Err(Error::NotFound)),
2132 Poll::Pending => Poll::Pending,
2133 }
2134 }
2135}
2136
2137pub struct Subscriber {
2160 name: Arc<str>,
2161 info: Info,
2162 inner: SubscriberKind,
2163 stats: stats::Scope,
2166 _stats_sub: stats::Subscription,
2169}
2170
2171enum SubscriberKind {
2172 Plain(PlainSubscriber),
2173 Spliced(Box<super::resume::Subscriber>),
2175}
2176
2177struct PlainSubscriber {
2179 state: kio::Consumer<TrackState>,
2180
2181 subscription: kio::Producer<Subscription>,
2182 index: usize,
2184 datagram_index: usize,
2186 min_sequence: u64,
2188 next_sequence: u64,
2191 end_sequence: Option<u64>,
2196}
2197
2198impl PlainSubscriber {
2199 fn poll<F, R>(&self, waiter: &kio::Waiter, f: F) -> Poll<Result<R>>
2201 where
2202 F: Fn(&kio::Ref<'_, TrackState>) -> Poll<Result<R>>,
2203 {
2204 Poll::Ready(match ready!(self.state.poll(waiter, f)) {
2205 Ok(res) => res,
2206 Err(state) => Err(state.abort.clone().unwrap_or(Error::Dropped)),
2208 })
2209 }
2210
2211 fn poll_recv_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2212 let Some((consumer, found_index)) =
2213 ready!(self.poll(waiter, |state| state.poll_recv_group(self.index, self.min_sequence))?)
2214 else {
2215 return Poll::Ready(Ok(None));
2216 };
2217
2218 self.index = found_index + 1;
2219 Poll::Ready(Ok(Some(consumer)))
2220 }
2221
2222 fn poll_recv_datagram(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Datagram>>> {
2223 let Some((datagram, found_index)) =
2224 ready!(self.poll(waiter, |state| state.poll_recv_datagram(self.datagram_index))?)
2225 else {
2226 return Poll::Ready(Ok(None));
2227 };
2228
2229 self.datagram_index = found_index + 1;
2230 self.next_sequence = self.next_sequence.max(datagram.sequence.saturating_add(1));
2231 Poll::Ready(Ok(Some(datagram)))
2232 }
2233
2234 fn poll_next_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2235 let floor = self.next_sequence.max(self.min_sequence);
2236 let Some(group) = ready!(self.poll(waiter, |state| state.poll_next_in_range(floor, self.end_sequence))?) else {
2237 return Poll::Ready(Ok(None));
2238 };
2239 self.next_sequence = group.sequence.saturating_add(1);
2240 Poll::Ready(Ok(Some(group)))
2241 }
2242
2243 fn poll_read_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Frame>>> {
2244 let lower = self.min_sequence.max(self.next_sequence);
2245 let Some((frame, found_index, sequence)) =
2246 ready!(self.poll(waiter, |state| { state.poll_read_frame(self.index, lower, waiter) })?)
2247 else {
2248 return Poll::Ready(Ok(None));
2249 };
2250
2251 self.index = found_index + 1;
2252 self.next_sequence = sequence.saturating_add(1);
2253 Poll::Ready(Ok(Some(frame)))
2254 }
2255}
2256
2257#[derive(Clone)]
2263pub struct SubscriberControl {
2264 subscription: kio::Producer<Subscription>,
2265}
2266
2267impl SubscriberControl {
2268 pub fn subscription(&self) -> Subscription {
2270 self.subscription.read().clone()
2271 }
2272
2273 pub fn update(&self, subscription: Subscription) -> Result<()> {
2278 let mut state = self.subscription.write().map_err(|_| Error::Closed)?;
2279 *state = subscription;
2280 Ok(())
2281 }
2282}
2283
2284impl Subscriber {
2285 pub fn info(&self) -> &Info {
2290 &self.info
2291 }
2292
2293 pub fn name(&self) -> &str {
2295 &self.name
2296 }
2297
2298 pub fn control(&self) -> SubscriberControl {
2300 SubscriberControl {
2301 subscription: match &self.inner {
2302 SubscriberKind::Plain(plain) => plain.subscription.clone(),
2303 SubscriberKind::Spliced(spliced) => spliced.prefs(),
2304 },
2305 }
2306 }
2307
2308 pub fn poll_recv_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2319 let meter = self.stats.meter();
2320 let res = match &mut self.inner {
2321 SubscriberKind::Plain(plain) => plain.poll_recv_group(waiter),
2322 SubscriberKind::Spliced(spliced) => spliced.poll_recv_group(waiter),
2323 };
2324 res.map(|res| res.map(|group| group.map(|group| group.with_meter(meter))))
2325 }
2326
2327 pub async fn recv_group(&mut self) -> Result<Option<group::Consumer>> {
2333 kio::wait(|waiter| self.poll_recv_group(waiter)).await
2334 }
2335
2336 pub fn poll_recv_datagram(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Datagram>>> {
2347 let meter = self.stats.meter();
2348 let res = match &mut self.inner {
2349 SubscriberKind::Plain(plain) => plain.poll_recv_datagram(waiter),
2350 SubscriberKind::Spliced(spliced) => spliced.poll_recv_datagram(waiter),
2351 };
2352 if let Poll::Ready(Ok(Some(datagram))) = &res {
2355 meter.datagram(datagram.payload.len() as u64);
2356 }
2357 res
2358 }
2359
2360 pub async fn recv_datagram(&mut self) -> Result<Option<Datagram>> {
2367 kio::wait(|waiter| self.poll_recv_datagram(waiter)).await
2368 }
2369
2370 pub fn poll_next_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2379 let meter = self.stats.meter();
2380 let res = match &mut self.inner {
2381 SubscriberKind::Plain(plain) => plain.poll_next_group(waiter),
2382 SubscriberKind::Spliced(spliced) => spliced.poll_next_group(waiter),
2383 };
2384 res.map(|res| res.map(|group| group.map(|group| group.with_meter(meter))))
2385 }
2386
2387 pub async fn next_group(&mut self) -> Result<Option<group::Consumer>> {
2393 kio::wait(|waiter| self.poll_next_group(waiter)).await
2394 }
2395
2396 pub fn poll_read_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Frame>>> {
2400 let meter = self.stats.meter();
2401 let res = match &mut self.inner {
2402 SubscriberKind::Plain(plain) => plain.poll_read_frame(waiter),
2403 SubscriberKind::Spliced(spliced) => spliced.poll_read_frame(waiter),
2404 };
2405 if let Poll::Ready(Ok(Some(frame))) = &res {
2408 meter.group();
2409 meter.frames(1);
2410 meter.bytes(frame.payload.len() as u64);
2411 }
2412 res
2413 }
2414
2415 pub async fn read_frame(&mut self) -> Result<Option<frame::Frame>> {
2420 kio::wait(|waiter| self.poll_read_frame(waiter)).await
2421 }
2422
2423 pub fn is_clone(&self, other: &Self) -> bool {
2425 match (&self.inner, &other.inner) {
2426 (SubscriberKind::Plain(a), SubscriberKind::Plain(b)) => a.state.same_channel(&b.state),
2427 (SubscriberKind::Spliced(a), SubscriberKind::Spliced(b)) => a.is_clone(b),
2428 _ => false,
2429 }
2430 }
2431
2432 pub fn poll_finished(&mut self, waiter: &kio::Waiter) -> Poll<Result<u64>> {
2434 match &mut self.inner {
2435 SubscriberKind::Plain(plain) => plain.poll(waiter, |state| state.poll_finished()),
2436 SubscriberKind::Spliced(spliced) => spliced.poll_finished(waiter),
2437 }
2438 }
2439
2440 pub async fn finished(&mut self) -> Result<u64> {
2448 kio::wait(|waiter| self.poll_finished(waiter)).await
2449 }
2450
2451 pub fn start_at(&mut self, sequence: u64) {
2458 match &mut self.inner {
2459 SubscriberKind::Plain(plain) => plain.min_sequence = sequence,
2460 SubscriberKind::Spliced(spliced) => spliced.start_at(sequence),
2461 }
2462 }
2463
2464 pub fn end_at(&mut self, sequence: impl Into<Option<u64>>) {
2477 match &mut self.inner {
2478 SubscriberKind::Plain(plain) => plain.end_sequence = sequence.into(),
2479 SubscriberKind::Spliced(spliced) => spliced.end_at(sequence),
2480 }
2481 }
2482
2483 pub fn subscription(&self) -> Subscription {
2485 self.control().subscription()
2486 }
2487
2488 pub fn update(&mut self, subscription: Subscription) -> Result<()> {
2494 match &mut self.inner {
2495 SubscriberKind::Plain(plain) => {
2496 let mut state = plain.subscription.write().map_err(|_| Error::Closed)?;
2497 *state = subscription;
2498 }
2499 SubscriberKind::Spliced(spliced) => spliced.update(subscription),
2500 }
2501 Ok(())
2502 }
2503
2504 pub fn latest(&self) -> Option<u64> {
2506 match &self.inner {
2507 SubscriberKind::Plain(plain) => plain.state.read().max_sequence,
2508 SubscriberKind::Spliced(spliced) => spliced.latest(),
2509 }
2510 }
2511}
2512
2513pub struct Request {
2525 name: Arc<str>,
2526 broadcast: Arc<broadcast::Info>,
2528 state: kio::Producer<TrackState>,
2529
2530 prev_subscription: Option<Subscription>,
2532
2533 alive: Arc<Alive>,
2536
2537 _dynamic: Dynamic,
2542
2543 stats: stats::Scope,
2546}
2547
2548impl Request {
2549 pub(crate) fn new(broadcast: Arc<broadcast::Info>, name: impl Into<Arc<str>>) -> Self {
2550 let name = name.into();
2551 let state = TrackState::spawn(broadcast.clone());
2552 let alive = Alive::new(name.clone(), state.clone());
2553 let dynamic = Dynamic::new(name.clone(), state.clone(), alive.clone());
2554 Self {
2555 name,
2556 broadcast,
2557 state,
2558 prev_subscription: None,
2559 alive,
2560 _dynamic: dynamic,
2561 stats: stats::Scope::default(),
2562 }
2563 }
2564
2565 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
2568 self.stats = scope;
2569 self
2570 }
2571
2572 pub fn name(&self) -> &str {
2574 &self.name
2575 }
2576
2577 pub fn consume(&self) -> Consumer {
2579 Consumer::plain(self.name.clone(), self.state.consume())
2580 }
2581
2582 pub fn dynamic(&self) -> Dynamic {
2586 Dynamic::new(self.name.clone(), self.state.clone(), self.alive.clone())
2587 }
2588
2589 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
2592 self.state.poll_unused(waiter).map(|_| ())
2593 }
2594
2595 pub fn accept(self, info: impl Into<Option<Info>>) -> Producer {
2602 if let Ok(mut state) = self.state.write() {
2605 state.install(info.into().unwrap_or_default());
2606 }
2607 self.alive.publish(Some(&self.stats));
2610 Producer {
2611 name: self.name,
2612 broadcast: self.broadcast,
2613 state: self.state,
2614 prev_subscription: None,
2615 alive: self.alive,
2616 stats: self.stats,
2617 }
2618 }
2619
2620 pub fn reject(self, err: Error) {
2622 if let Ok(mut state) = self.state.write() {
2623 state.abort = Some(err);
2624 }
2625 }
2626
2627 pub fn subscription(&self) -> Option<Subscription> {
2630 let state = self.state.read();
2631 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
2632 drop(state);
2633 snapshot_subscription(&subs, bound)
2634 }
2635
2636 pub async fn subscription_changed(&mut self) -> Option<Subscription> {
2639 kio::wait(|waiter| self.poll_subscription_changed(waiter)).await
2640 }
2641
2642 pub fn poll_subscription_changed(&mut self, waiter: &kio::Waiter) -> Poll<Option<Subscription>> {
2644 let state = self.state.read();
2645 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
2646 drop(state);
2647
2648 let prev = &self.prev_subscription;
2649 let mut combined = None;
2650 let mut guard = ready!(subs.poll(waiter, |subs| {
2651 let next = combined_subscription(subs, bound, waiter);
2652 if &next == prev {
2653 Poll::Pending
2654 } else {
2655 combined = next;
2656 Poll::Ready(())
2657 }
2658 }));
2659 guard.retain(|sub| !sub.is_closed());
2661 drop(guard);
2662 self.prev_subscription = combined.clone();
2663 Poll::Ready(combined)
2664 }
2665
2666 pub(super) fn weak(&self) -> TrackWeak {
2667 TrackWeak {
2668 name: self.name.clone(),
2669 state: self.state.weak(),
2670 }
2671 }
2672}
2673
2674#[cfg(test)]
2675use futures::FutureExt;
2676
2677#[cfg(test)]
2678#[allow(missing_docs)] impl Subscriber {
2680 pub fn assert_group(&mut self) -> group::Consumer {
2681 self.recv_group()
2682 .now_or_never()
2683 .expect("group would have blocked")
2684 .expect("would have errored")
2685 .expect("track was closed")
2686 }
2687
2688 pub fn assert_no_group(&mut self) {
2689 assert!(
2690 self.recv_group().now_or_never().is_none(),
2691 "recv_group would not have blocked"
2692 );
2693 }
2694
2695 pub fn assert_not_closed(&mut self) {
2696 assert!(self.finished().now_or_never().is_none(), "should not be closed");
2697 }
2698
2699 pub fn assert_closed(&mut self) {
2700 assert!(self.finished().now_or_never().is_some(), "should be closed");
2701 }
2702
2703 pub fn assert_error(&mut self) {
2705 assert!(
2706 self.finished().now_or_never().expect("should not block").is_err(),
2707 "should be error"
2708 );
2709 }
2710
2711 pub fn assert_is_clone(&self, other: &Self) {
2712 assert!(self.is_clone(other), "should be clone");
2713 }
2714
2715 pub fn assert_not_clone(&self, other: &Self) {
2716 assert!(!self.is_clone(other), "should not be clone");
2717 }
2718}
2719
2720#[cfg(test)]
2721mod test {
2722 use super::*;
2723 use crate::model::test_tracing::count_drop_warnings;
2724
2725 fn track_producer(name: impl Into<Arc<str>>, info: impl Into<Option<Info>>) -> Producer {
2728 Producer::new(Arc::new(broadcast::Info::default()), name, info)
2729 }
2730
2731 fn live_groups(state: &TrackState) -> usize {
2733 state.lookup.len()
2734 }
2735
2736 fn first_live_sequence(state: &TrackState) -> u64 {
2738 state
2739 .arrival
2740 .iter()
2741 .find(|(sequence, stamp)| state.lookup.get(sequence).is_some_and(|slot| slot.stamp == *stamp))
2742 .map(|(sequence, _)| *sequence)
2743 .unwrap()
2744 }
2745
2746 fn recv_datagram(dg: &mut Subscriber) -> Datagram {
2748 dg.recv_datagram()
2749 .now_or_never()
2750 .expect("datagram would have blocked")
2751 .expect("would have errored")
2752 .expect("track was closed")
2753 }
2754
2755 #[tokio::test]
2756 async fn append_datagram_shares_group_sequence() {
2757 let mut producer = track_producer("test", None);
2758 let ts = Timestamp::from_millis(10).unwrap();
2759
2760 assert_eq!(producer.append_group().unwrap().sequence, 0);
2762 assert_eq!(producer.append_datagram(ts, &b"a"[..]).unwrap(), 1);
2763 assert_eq!(producer.append_group().unwrap().sequence, 2);
2764 assert_eq!(producer.append_datagram(ts, &b"b"[..]).unwrap(), 3);
2765 assert_eq!(producer.latest(), Some(3));
2766 }
2767
2768 #[tokio::test]
2769 async fn append_datagram_roundtrip() {
2770 let mut producer = track_producer("test", None);
2771 let mut dg = producer.subscribe(None);
2772
2773 let ts = Timestamp::from_millis(42).unwrap();
2774 let seq = producer.append_datagram(ts, &b"hello"[..]).unwrap();
2775
2776 let got = recv_datagram(&mut dg);
2777 assert_eq!(got.sequence, seq);
2778 assert_eq!(got.timestamp, ts);
2779 assert_eq!(&got.payload[..], b"hello");
2780 }
2781
2782 #[tokio::test]
2783 async fn write_datagram_preserves_sequence() {
2784 let mut producer = track_producer("test", None);
2785 let mut dg = producer.subscribe(None);
2786
2787 let ts = Timestamp::from_millis(5).unwrap();
2788 producer
2790 .write_datagram(Datagram {
2791 sequence: 100,
2792 timestamp: ts,
2793 payload: bytes::Bytes::from_static(b"x"),
2794 })
2795 .unwrap();
2796
2797 assert_eq!(recv_datagram(&mut dg).sequence, 100);
2798 assert_eq!(producer.append_group().unwrap().sequence, 101);
2800 }
2801
2802 #[tokio::test]
2803 async fn recv_datagram_advances_ordered_group_cursor() {
2804 let mut producer = track_producer("test", None);
2805 let mut subscriber = producer.subscribe(None);
2806 let ts = Timestamp::from_millis(5).unwrap();
2807
2808 producer
2809 .write_datagram(Datagram {
2810 sequence: 5,
2811 timestamp: ts,
2812 payload: bytes::Bytes::from_static(b"x"),
2813 })
2814 .unwrap();
2815 assert_eq!(recv_datagram(&mut subscriber).sequence, 5);
2816
2817 producer.create_group(group::Info { sequence: 3 }).unwrap();
2818 producer.create_group(group::Info { sequence: 6 }).unwrap();
2819
2820 let group = subscriber
2821 .next_group()
2822 .now_or_never()
2823 .expect("group would have blocked")
2824 .expect("would have errored")
2825 .expect("track was closed");
2826 assert_eq!(group.sequence, 6);
2827 }
2828
2829 #[tokio::test]
2830 async fn datagram_normalized_to_track_timescale() {
2831 let info = Info::default().with_timescale(Timescale::MICRO);
2832 let mut producer = track_producer("test", info);
2833 let mut dg = producer.subscribe(None);
2834
2835 producer
2837 .append_datagram(Timestamp::from_millis(2).unwrap(), &b"z"[..])
2838 .unwrap();
2839 let got = recv_datagram(&mut dg);
2840 assert_eq!(got.timestamp.scale(), Timescale::MICRO);
2841 assert_eq!(got.timestamp.value(), 2_000);
2842 }
2843
2844 #[tokio::test]
2845 async fn datagram_rejects_oversized() {
2846 let mut producer = track_producer("test", None);
2847 let big = bytes::Bytes::from(vec![0u8; crate::model::datagram::MAX_DATAGRAM_PAYLOAD + 1]);
2848 let ts = Timestamp::from_millis(0).unwrap();
2849 assert!(matches!(
2850 producer.append_datagram(ts, big.clone()),
2851 Err(Error::FrameTooLarge)
2852 ));
2853 assert!(matches!(
2854 producer.write_datagram(Datagram {
2855 sequence: 0,
2856 timestamp: ts,
2857 payload: big,
2858 }),
2859 Err(Error::FrameTooLarge)
2860 ));
2861 }
2862
2863 #[tokio::test]
2864 async fn datagram_fanout_to_subscribers() {
2865 let mut producer = track_producer("test", None);
2866 let mut a = producer.subscribe(None);
2868 let mut b = producer.subscribe(None);
2869 let ts = Timestamp::from_millis(1).unwrap();
2870
2871 producer.append_datagram(ts, &b"first"[..]).unwrap();
2872 producer.append_datagram(ts, &b"second"[..]).unwrap();
2873
2874 assert_eq!(&recv_datagram(&mut a).payload[..], b"first");
2876 assert_eq!(&recv_datagram(&mut a).payload[..], b"second");
2877 assert_eq!(&recv_datagram(&mut b).payload[..], b"first");
2878 assert_eq!(&recv_datagram(&mut b).payload[..], b"second");
2879 }
2880
2881 #[tokio::test]
2882 async fn datagram_evicts_stale() {
2883 tokio::time::pause();
2884
2885 let mut producer = track_producer("test", None);
2886 let mut dg = producer.subscribe(None);
2887 let ts = Timestamp::from_millis(0).unwrap();
2888
2889 producer.append_datagram(ts, &b"old"[..]).unwrap(); tokio::time::advance(MAX_DATAGRAM_AGE + Duration::from_millis(10)).await;
2893 producer.append_datagram(ts, &b"new"[..]).unwrap(); let got = recv_datagram(&mut dg);
2897 assert_eq!(got.sequence, 1);
2898 assert_eq!(&got.payload[..], b"new");
2899 }
2900
2901 #[tokio::test]
2902 async fn datagram_recv_pends_until_written() {
2903 let mut producer = track_producer("test", None);
2904 let mut dg = producer.subscribe(None);
2905
2906 assert!(
2907 dg.recv_datagram().now_or_never().is_none(),
2908 "should block with no datagrams"
2909 );
2910
2911 producer
2912 .append_datagram(Timestamp::from_millis(0).unwrap(), &b"go"[..])
2913 .unwrap();
2914 assert_eq!(&recv_datagram(&mut dg).payload[..], b"go");
2915 }
2916
2917 #[tokio::test]
2921 async fn datagram_wire_roundtrip_between_tracks() {
2922 use crate::coding::{Decode, Encode};
2923 use crate::lite;
2924
2925 let version = lite::Version::Lite05;
2926
2927 let mut origin = track_producer("test", None);
2929 let mut origin_dg = origin.subscribe(None);
2930 let ts = Timestamp::from_millis(7).unwrap();
2931 let seq = origin.append_datagram(ts, &b"payload"[..]).unwrap();
2932
2933 let d = recv_datagram(&mut origin_dg);
2934 let body = lite::Datagram {
2935 subscribe: 5,
2936 sequence: d.sequence,
2937 timestamp: d.timestamp.value(),
2938 payload: d.payload.clone(),
2939 }
2940 .encode_bytes(version)
2941 .unwrap();
2942
2943 let mut slice = &body[..];
2945 let wire = lite::Datagram::decode(&mut slice, version).unwrap();
2946 let mut downstream = track_producer("test", None);
2947 let mut downstream_dg = downstream.subscribe(None);
2948 downstream
2949 .write_datagram(Datagram {
2950 sequence: wire.sequence,
2951 timestamp: Timestamp::new(wire.timestamp, Timescale::MILLI).unwrap(),
2952 payload: wire.payload,
2953 })
2954 .unwrap();
2955
2956 let got = recv_datagram(&mut downstream_dg);
2957 assert_eq!(got.sequence, seq);
2958 assert_eq!(got.timestamp, ts);
2959 assert_eq!(&got.payload[..], b"payload");
2960 }
2961
2962 #[tokio::test]
2963 async fn evict_expired_groups() {
2964 tokio::time::pause();
2965
2966 let mut producer = track_producer("test", None);
2967
2968 producer.append_group().unwrap(); producer.append_group().unwrap(); producer.append_group().unwrap(); {
2974 let state = producer.state.read();
2975 assert_eq!(live_groups(&state), 3);
2976 assert_eq!(state.offset, 0);
2977 }
2978
2979 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
2981
2982 producer.append_group().unwrap(); {
2989 let state = producer.state.read();
2990 assert_eq!(live_groups(&state), 1);
2991 assert_eq!(first_live_sequence(&state), 3);
2992 assert_eq!(state.offset, 3);
2993 assert!(!state.lookup.contains_key(&0));
2994 assert!(!state.lookup.contains_key(&1));
2995 assert!(!state.lookup.contains_key(&2));
2996 assert!(state.lookup.contains_key(&3));
2997 }
2998 }
2999
3000 #[tokio::test]
3004 async fn aging_out_a_finished_group_keeps_the_clean_end() {
3005 tokio::time::pause();
3006
3007 let mut producer = track_producer("test", None);
3008 let mut group = producer.create_group(group::Info { sequence: 0 }).unwrap();
3009 let mut consumer = group.consume();
3010
3011 group
3012 .write_frame(Timestamp::from_millis(0).unwrap(), b"hello".as_slice())
3013 .unwrap();
3014 assert_eq!(consumer.next_frame().await.unwrap().unwrap().size, 5);
3015
3016 tokio::time::advance(DEFAULT_LATENCY_MAX * 12).await;
3018 group.finish().unwrap();
3019 let _next = producer.create_group(group::Info { sequence: 1 }).unwrap();
3020
3021 assert!(consumer.next_frame().await.unwrap().is_none());
3022 }
3023
3024 #[tokio::test]
3025 async fn evict_keeps_max_sequence() {
3026 tokio::time::pause();
3027
3028 let mut producer = track_producer("test", None);
3029 producer.append_group().unwrap(); tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3033
3034 producer.append_group().unwrap(); {
3038 let state = producer.state.read();
3039 assert_eq!(live_groups(&state), 1);
3040 assert_eq!(first_live_sequence(&state), 1);
3041 assert_eq!(state.offset, 1);
3042 }
3043 }
3044
3045 #[tokio::test]
3046 async fn no_eviction_when_fresh() {
3047 tokio::time::pause();
3048
3049 let mut producer = track_producer("test", None);
3050 producer.append_group().unwrap(); producer.append_group().unwrap(); producer.append_group().unwrap(); {
3055 let state = producer.state.read();
3056 assert_eq!(live_groups(&state), 3);
3057 assert_eq!(state.offset, 0);
3058 }
3059 }
3060
3061 #[tokio::test]
3062 async fn consumer_skips_evicted_groups() {
3063 tokio::time::pause();
3064
3065 let mut producer = track_producer("test", None);
3066 producer.append_group().unwrap(); let mut consumer = producer.subscribe(None);
3069
3070 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3071 producer.append_group().unwrap(); let group = consumer.assert_group();
3075 assert_eq!(group.sequence, 1);
3076 }
3077
3078 #[tokio::test]
3079 async fn cache_age_controls_eviction() {
3080 tokio::time::pause();
3081
3082 let mut producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(1)));
3084 producer.append_group().unwrap(); tokio::time::advance(Duration::from_secs(2)).await;
3088 producer.append_group().unwrap(); let state = producer.state.read();
3092 assert_eq!(live_groups(&state), 1);
3093 assert_eq!(first_live_sequence(&state), 1);
3094 }
3095
3096 #[test]
3097 fn latency_max_clamped_to_cache() {
3098 let producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(2)));
3099
3100 let mut subscriber = producer.subscribe(Subscription::default().with_latency_max(Duration::from_secs(10)));
3104 assert_eq!(subscriber.subscription().latency_max, Duration::from_secs(10));
3105 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3106
3107 subscriber
3109 .update(Subscription::default().with_latency_max(Duration::from_millis(500)))
3110 .unwrap();
3111 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_millis(500));
3112
3113 subscriber
3114 .update(Subscription::default().with_latency_max(Duration::ZERO))
3115 .unwrap();
3116 assert_eq!(producer.subscription().unwrap().latency_max, Duration::ZERO);
3117 }
3118
3119 fn track_producer_capped(name: impl Into<Arc<str>>, info: Info, cap: Duration) -> Producer {
3122 let origin = crate::origin::Info::default().with_cache_duration(cap);
3123 Producer::new(Arc::new(broadcast::Info { origin }), name, info)
3124 }
3125
3126 #[test]
3127 fn origin_cache_duration_clamps_latency_max() {
3128 let capped = track_producer_capped(
3131 "test",
3132 Info::default().with_latency_max(Duration::from_secs(60)),
3133 Duration::from_secs(1),
3134 );
3135 assert_eq!(capped.state.read().latency_bound(), Some(Duration::from_secs(1)));
3136
3137 let under = track_producer_capped(
3138 "test",
3139 Info::default().with_latency_max(Duration::from_millis(500)),
3140 Duration::from_secs(1),
3141 );
3142 assert_eq!(under.state.read().latency_bound(), Some(Duration::from_millis(500)));
3143 }
3144
3145 #[tokio::test]
3146 async fn origin_cache_duration_caps_eviction() {
3147 tokio::time::pause();
3148
3149 let mut producer = track_producer_capped(
3151 "test",
3152 Info::default().with_latency_max(Duration::from_secs(60)),
3153 Duration::from_secs(1),
3154 );
3155 producer.append_group().unwrap(); tokio::time::advance(Duration::from_secs(2)).await;
3159 producer.append_group().unwrap(); let state = producer.state.read();
3163 assert_eq!(live_groups(&state), 1);
3164 assert_eq!(first_live_sequence(&state), 1);
3165 }
3166
3167 #[test]
3168 fn latency_max_clamped_via_every_update_path() {
3169 let producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(2)));
3170 let over = Subscription::default().with_latency_max(Duration::from_secs(10));
3171
3172 let mut subscriber = producer.subscribe(over.clone());
3175 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3176
3177 subscriber.control().update(over.clone()).unwrap();
3178 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3179
3180 subscriber.update(over).unwrap();
3181 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3182 }
3183
3184 #[test]
3185 fn latency_max_aggregate_clamps_the_max_across_subscribers() {
3186 let producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(2)));
3187
3188 let _a = producer.subscribe(Subscription::default().with_latency_max(Duration::from_millis(500)));
3191 let _b = producer.subscribe(Subscription::default().with_latency_max(Duration::from_secs(10)));
3192
3193 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
3194 }
3195
3196 #[test]
3197 fn subscriber_control_updates_while_read_future_is_pending() {
3198 let producer = track_producer("test", None);
3199 let mut subscriber = producer.subscribe(None);
3200 let control = subscriber.control();
3201
3202 let mut recv = Box::pin(subscriber.recv_group());
3203 assert!(recv.as_mut().now_or_never().is_none());
3204
3205 control
3206 .update(Subscription::default().with_priority(7).with_ordered(false))
3207 .unwrap();
3208
3209 let aggregate = producer.subscription().expect("expected an active subscription");
3210 assert_eq!(aggregate.priority, 7);
3211 assert!(!aggregate.ordered);
3212 }
3213
3214 #[test]
3215 fn dropped_subscriber_leaves_no_ghost_in_aggregate() {
3216 let mut producer = track_producer("test", None);
3221 let a = producer.subscribe(Subscription::default().with_priority(5));
3222
3223 let waiter = kio::Waiter::noop();
3225 assert!(
3226 matches!(producer.poll_subscription_changed(&waiter), Poll::Ready(Ok(Some(_)))),
3227 "one live subscriber should aggregate to Some",
3228 );
3229
3230 drop(a);
3232
3233 assert!(
3235 matches!(producer.poll_subscription_changed(&waiter), Poll::Ready(Ok(None))),
3236 "a dropped subscriber must not linger in the aggregate",
3237 );
3238
3239 assert!(
3241 producer.subscription().is_none(),
3242 "snapshot must exclude a dropped subscriber",
3243 );
3244 }
3245
3246 #[test]
3247 fn dropped_subscriber_wakes_the_aggregate() {
3248 use std::sync::atomic::{AtomicBool, Ordering};
3255
3256 let mut producer = track_producer("test", None);
3257 let a = producer.subscribe(Subscription::default().with_priority(5));
3258
3259 let woken = Arc::new(AtomicBool::new(false));
3260 let waiter = kio::Waiter::new(futures::task::waker(Arc::new(FlagWake(woken.clone()))));
3261
3262 assert!(matches!(
3264 producer.poll_subscription_changed(&waiter),
3265 Poll::Ready(Ok(Some(_)))
3266 ));
3267 assert!(
3268 producer.poll_subscription_changed(&waiter).is_pending(),
3269 "the aggregate is unchanged, so this poll must park",
3270 );
3271 assert!(!woken.load(Ordering::SeqCst), "nothing happened yet");
3272
3273 drop(a);
3274 assert!(
3275 woken.load(Ordering::SeqCst),
3276 "the last subscriber leaving must wake the aggregate watcher",
3277 );
3278 }
3279
3280 struct FlagWake(Arc<std::sync::atomic::AtomicBool>);
3282
3283 impl futures::task::ArcWake for FlagWake {
3284 fn wake_by_ref(arc_self: &Arc<Self>) {
3285 arc_self.0.store(true, std::sync::atomic::Ordering::SeqCst);
3286 }
3287 }
3288
3289 #[tokio::test]
3290 async fn out_of_order_max_sequence_at_front() {
3291 tokio::time::pause();
3292
3293 let mut producer = track_producer("test", None);
3294
3295 producer.create_group(group::Info { sequence: 5 }).unwrap();
3297 producer.create_group(group::Info { sequence: 3 }).unwrap();
3298 producer.create_group(group::Info { sequence: 4 }).unwrap();
3299
3300 {
3302 let state = producer.state.read();
3303 assert_eq!(state.max_sequence, Some(5));
3304 }
3305
3306 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3308
3309 producer.append_group().unwrap(); {
3315 let state = producer.state.read();
3316 assert_eq!(live_groups(&state), 1);
3317 assert_eq!(first_live_sequence(&state), 6);
3318 assert!(!state.lookup.contains_key(&3));
3319 assert!(!state.lookup.contains_key(&4));
3320 assert!(!state.lookup.contains_key(&5));
3321 assert!(state.lookup.contains_key(&6));
3322 }
3323 }
3324
3325 #[tokio::test]
3326 async fn max_sequence_at_front_blocks_trim() {
3327 tokio::time::pause();
3328
3329 let mut producer = track_producer("test", None);
3330
3331 producer.create_group(group::Info { sequence: 5 }).unwrap();
3333
3334 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3335
3336 producer.create_group(group::Info { sequence: 3 }).unwrap();
3338
3339 {
3342 let state = producer.state.read();
3343 assert_eq!(live_groups(&state), 2);
3344 assert_eq!(state.offset, 0);
3345 }
3346
3347 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
3349
3350 producer.create_group(group::Info { sequence: 2 }).unwrap();
3352
3353 {
3358 let state = producer.state.read();
3359 assert_eq!(live_groups(&state), 2);
3360 assert_eq!(state.offset, 0);
3361 assert!(state.lookup.contains_key(&5));
3362 assert!(!state.lookup.contains_key(&3));
3363 assert!(state.lookup.contains_key(&2));
3364 }
3365
3366 let mut consumer = producer.subscribe(None);
3368 let group = consumer.assert_group();
3369 assert_eq!(group.sequence, 5);
3371 }
3372
3373 #[tokio::test]
3374 async fn abort_clears_cached_groups() {
3375 let mut producer = track_producer("test", None);
3376 producer.append_group().unwrap();
3377 producer.append_group().unwrap();
3378
3379 let mut consumer = producer.subscribe(None);
3381 assert_eq!(live_groups(&producer.state.read()), 2);
3382
3383 producer.clone().abort(Error::Cancel).unwrap();
3384
3385 {
3386 let state = producer.state.read();
3387 assert!(state.lookup.is_empty(), "cached groups should be dropped on abort");
3388 assert!(state.arrival.is_empty());
3389 assert!(state.evict.is_empty());
3390 }
3391
3392 let result = consumer.recv_group().now_or_never().expect("should not block");
3394 assert!(matches!(result, Err(Error::Cancel)));
3395 }
3396
3397 #[tokio::test]
3398 async fn drop_unfinished_clears_cached_groups() {
3399 let producer = track_producer("test", None);
3400 let mut writer = producer.clone();
3401 writer.append_group().unwrap();
3402
3403 let mut consumer = producer.subscribe(None);
3405 assert_eq!(live_groups(&producer.state.read()), 1);
3406
3407 drop(writer);
3409 drop(producer);
3410
3411 let result = consumer.recv_group().now_or_never().expect("should not block");
3412 assert!(matches!(result, Err(Error::Dropped)));
3413 }
3414
3415 #[tokio::test]
3416 async fn drop_after_abort_does_not_warn() {
3417 let warns = count_drop_warnings("track::Producer dropped without finish", || {
3420 let producer = track_producer("test", None);
3421 let keep = producer.clone();
3422 let mut writer = producer.clone();
3423 let mut group = writer.append_group().unwrap();
3424 group.finish().unwrap();
3425 let _consumer = producer.subscribe(None);
3426 writer.abort(Error::Cancel).unwrap();
3427 drop(keep);
3428 });
3429 assert_eq!(warns, 0, "abort-then-drop must not emit unfinished-producer WARN");
3430 }
3431
3432 #[tokio::test]
3433 async fn drop_unfinished_warns() {
3434 let warns = count_drop_warnings("track::Producer dropped without finish", || {
3435 let producer = track_producer("test", None);
3436 let mut writer = producer.clone();
3437 writer.append_group().unwrap();
3438 let _consumer = producer.subscribe(None);
3439 drop(writer);
3440 drop(producer);
3441 });
3442 assert!(warns >= 1, "unfinished drop must emit unfinished-producer WARN");
3443 }
3444
3445 #[tokio::test]
3446 async fn drop_finished_keeps_cached_groups() {
3447 let mut producer = track_producer("test", None);
3448 producer.append_group().unwrap();
3449 producer.finish().unwrap();
3450
3451 let mut consumer = producer.subscribe(None);
3452 drop(producer);
3453
3454 assert_eq!(consumer.assert_group().sequence, 0);
3456 let done = consumer.recv_group().now_or_never().expect("should not block").unwrap();
3457 assert!(done.is_none(), "consumer should drain then see clean finish");
3458 }
3459
3460 #[test]
3461 fn append_finish_cannot_be_rewritten() {
3462 let mut producer = track_producer("test", None);
3463
3464 assert!(producer.finish().is_ok());
3466 assert!(producer.finish().is_err());
3467 assert!(producer.append_group().is_err());
3468 }
3469
3470 #[test]
3471 fn finish_after_groups() {
3472 let mut producer = track_producer("test", None);
3473
3474 producer.append_group().unwrap();
3475 assert!(producer.finish().is_ok());
3476 assert!(producer.finish().is_err());
3477 assert!(producer.append_group().is_err());
3478 }
3479
3480 #[test]
3481 fn finish_at_rejects_a_boundary_at_or_below_the_live_edge() {
3482 let mut producer = track_producer("test", None);
3483 producer.create_group(group::Info { sequence: 5 }).unwrap();
3484
3485 assert!(producer.finish_at(4).is_err());
3488 assert!(producer.finish_at(5).is_err());
3489 assert!(producer.finish_at(6).is_ok());
3490
3491 {
3492 let state = producer.state.read();
3493 assert_eq!(state.final_sequence, Some(6));
3494 }
3495
3496 assert!(producer.finish_at(6).is_err());
3498 assert!(producer.create_group(group::Info { sequence: 4 }).is_ok());
3499 assert!(producer.create_group(group::Info { sequence: 6 }).is_err());
3500 }
3501
3502 #[test]
3503 fn final_sequence_reports_the_declared_boundary() {
3504 let mut producer = track_producer("test", None);
3505 assert_eq!(producer.final_sequence(), None);
3506
3507 producer.create_group(group::Info { sequence: 5 }).unwrap();
3508 assert_eq!(producer.final_sequence(), None, "a group does not declare a boundary");
3509
3510 producer.finish_at(9).unwrap();
3511 assert_eq!(producer.final_sequence(), Some(9));
3512
3513 assert!(producer.finish().is_err());
3515 }
3516
3517 #[test]
3518 fn final_sequence_reports_the_live_edge_after_finish() {
3519 let mut producer = track_producer("test", None);
3520 producer.create_group(group::Info { sequence: 5 }).unwrap();
3521 producer.finish().unwrap();
3522 assert_eq!(producer.final_sequence(), Some(6));
3523 }
3524
3525 #[tokio::test]
3526 async fn finish_at_declares_a_future_boundary() {
3527 let mut producer = track_producer("test", None);
3528 producer.create_group(group::Info { sequence: 5 }).unwrap();
3529
3530 producer.finish_at(7).unwrap();
3532
3533 let mut consumer = producer.subscribe(None);
3534 assert_eq!(consumer.assert_group().sequence, 5);
3535
3536 let boundary = consumer
3539 .finished()
3540 .now_or_never()
3541 .expect("boundary is known immediately")
3542 .expect("would have errored");
3543 assert_eq!(boundary, 7);
3544 assert!(
3545 consumer.recv_group().now_or_never().is_none(),
3546 "should wait for the outstanding group"
3547 );
3548
3549 producer.create_group(group::Info { sequence: 6 }).unwrap();
3551 assert_eq!(consumer.assert_group().sequence, 6);
3552 let done = consumer
3553 .recv_group()
3554 .now_or_never()
3555 .expect("should not block")
3556 .expect("would have errored");
3557 assert!(done.is_none(), "track completes once the boundary is reached");
3558 }
3559
3560 #[tokio::test]
3561 async fn recv_group_finishes_without_waiting_for_gaps() {
3562 let mut producer = track_producer("test", None);
3563 producer.create_group(group::Info { sequence: 1 }).unwrap();
3564 producer.finish().unwrap();
3565
3566 let mut consumer = producer.subscribe(None);
3567 assert_eq!(consumer.assert_group().sequence, 1);
3568
3569 let done = consumer
3570 .recv_group()
3571 .now_or_never()
3572 .expect("should not block")
3573 .expect("would have errored");
3574 assert!(done.is_none(), "track should finish without waiting for gaps");
3575 }
3576
3577 #[tokio::test]
3578 async fn next_group_skips_late_arrivals() {
3579 let mut producer = track_producer("test", None);
3580 let mut consumer = producer.subscribe(None);
3581
3582 producer.create_group(group::Info { sequence: 5 }).unwrap();
3584 let group = consumer
3585 .next_group()
3586 .now_or_never()
3587 .expect("should not block")
3588 .expect("would have errored")
3589 .expect("track should not be closed");
3590 assert_eq!(group.sequence, 5);
3591
3592 producer.create_group(group::Info { sequence: 3 }).unwrap();
3594 producer.create_group(group::Info { sequence: 4 }).unwrap();
3596 producer.create_group(group::Info { sequence: 7 }).unwrap();
3598
3599 let group = consumer
3600 .next_group()
3601 .now_or_never()
3602 .expect("should not block")
3603 .expect("would have errored")
3604 .expect("track should not be closed");
3605 assert_eq!(group.sequence, 7);
3606
3607 assert!(
3609 consumer.next_group().now_or_never().is_none(),
3610 "should block waiting for a higher sequence"
3611 );
3612 }
3613
3614 #[tokio::test]
3615 async fn next_group_returns_arrivals_in_order() {
3616 let mut producer = track_producer("test", None);
3617 let mut consumer = producer.subscribe(None);
3618
3619 producer.create_group(group::Info { sequence: 3 }).unwrap();
3621 producer.create_group(group::Info { sequence: 5 }).unwrap();
3622
3623 let group = consumer
3624 .next_group()
3625 .now_or_never()
3626 .expect("should not block")
3627 .expect("would have errored")
3628 .expect("track should not be closed");
3629 assert_eq!(group.sequence, 3);
3630
3631 let group = consumer
3632 .next_group()
3633 .now_or_never()
3634 .expect("should not block")
3635 .expect("would have errored")
3636 .expect("track should not be closed");
3637 assert_eq!(group.sequence, 5);
3638 }
3639
3640 #[tokio::test]
3641 async fn next_group_and_recv_group_use_independent_cursors() {
3642 let mut producer = track_producer("test", None);
3643 let mut consumer = producer.subscribe(None);
3644
3645 producer.create_group(group::Info { sequence: 5 }).unwrap();
3647 producer.create_group(group::Info { sequence: 3 }).unwrap();
3648
3649 let group = consumer
3652 .next_group()
3653 .now_or_never()
3654 .expect("should not block")
3655 .expect("would have errored")
3656 .expect("track should not be closed");
3657 assert_eq!(group.sequence, 3);
3658
3659 assert_eq!(consumer.assert_group().sequence, 5);
3662 }
3663
3664 #[tokio::test]
3665 async fn end_at_caps_next_group() {
3666 let mut producer = track_producer("test", None);
3667 let mut consumer = producer.subscribe(None);
3668
3669 for s in 0..6 {
3670 producer.create_group(group::Info { sequence: s }).unwrap();
3671 }
3672
3673 consumer.end_at(2);
3674
3675 assert_eq!(
3677 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3678 0
3679 );
3680 assert_eq!(
3681 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3682 1
3683 );
3684 assert_eq!(
3685 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3686 2
3687 );
3688
3689 assert!(
3691 consumer.next_group().now_or_never().is_none(),
3692 "capped consumer must block instead of returning out-of-range groups"
3693 );
3694 }
3695
3696 #[tokio::test]
3697 async fn end_at_release_drains_cached_groups() {
3698 let mut producer = track_producer("test", None);
3699 let mut consumer = producer.subscribe(None);
3700
3701 for s in 0..6 {
3702 producer.create_group(group::Info { sequence: s }).unwrap();
3703 }
3704
3705 consumer.end_at(1);
3706 assert_eq!(
3707 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3708 0
3709 );
3710 assert_eq!(
3711 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3712 1
3713 );
3714 assert!(consumer.next_group().now_or_never().is_none(), "capped at 1");
3715
3716 consumer.end_at(4);
3718 assert_eq!(
3719 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3720 2
3721 );
3722 assert_eq!(
3723 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3724 3
3725 );
3726 assert_eq!(
3727 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3728 4
3729 );
3730 assert!(consumer.next_group().now_or_never().is_none(), "capped at 4");
3731
3732 consumer.end_at(None);
3734 assert_eq!(
3735 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3736 5
3737 );
3738 assert!(consumer.next_group().now_or_never().is_none(), "no more groups");
3739 }
3740
3741 #[tokio::test]
3742 async fn end_at_lower_than_cursor_parks_consumer() {
3743 let mut producer = track_producer("test", None);
3744 let mut consumer = producer.subscribe(None);
3745
3746 for s in 0..3 {
3747 producer.create_group(group::Info { sequence: s }).unwrap();
3748 }
3749
3750 assert_eq!(
3752 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3753 0
3754 );
3755 assert_eq!(
3756 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3757 1
3758 );
3759 assert_eq!(
3760 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3761 2
3762 );
3763
3764 consumer.end_at(1);
3766 producer.create_group(group::Info { sequence: 3 }).unwrap();
3767 producer.create_group(group::Info { sequence: 4 }).unwrap();
3768 assert!(
3769 consumer.next_group().now_or_never().is_none(),
3770 "cap is below cursor; nothing returnable until cap rises"
3771 );
3772
3773 consumer.end_at(None);
3775 assert_eq!(
3776 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3777 3
3778 );
3779 assert_eq!(
3780 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3781 4
3782 );
3783 }
3784
3785 #[tokio::test]
3786 async fn end_at_toggling_around_late_arrivals() {
3787 let mut producer = track_producer("test", None);
3788 let mut consumer = producer.subscribe(None);
3789
3790 consumer.end_at(5);
3791
3792 producer.create_group(group::Info { sequence: 2 }).unwrap();
3794 producer.create_group(group::Info { sequence: 5 }).unwrap();
3795 producer.create_group(group::Info { sequence: 3 }).unwrap();
3796 producer.create_group(group::Info { sequence: 8 }).unwrap();
3798 producer.create_group(group::Info { sequence: 4 }).unwrap();
3799
3800 assert_eq!(
3802 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3803 2
3804 );
3805 assert_eq!(
3806 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3807 3
3808 );
3809 assert_eq!(
3810 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3811 4
3812 );
3813 assert_eq!(
3814 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3815 5
3816 );
3817 assert!(consumer.next_group().now_or_never().is_none());
3819
3820 consumer.end_at(10);
3822 assert_eq!(
3823 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3824 8
3825 );
3826 }
3827
3828 #[tokio::test]
3829 async fn read_frame_returns_single_frame_per_group() {
3830 let mut producer = track_producer("test", None);
3831 let mut consumer = producer.subscribe(None);
3832
3833 producer.write_frame(Timestamp::ZERO, b"hello".as_slice()).unwrap();
3834 producer.write_frame(Timestamp::ZERO, b"world".as_slice()).unwrap();
3835
3836 let frame = consumer
3837 .read_frame()
3838 .now_or_never()
3839 .expect("should not block")
3840 .expect("would have errored")
3841 .expect("track should not be closed");
3842 assert_eq!(&frame.payload[..], b"hello");
3843
3844 let frame = consumer
3845 .read_frame()
3846 .now_or_never()
3847 .expect("should not block")
3848 .expect("would have errored")
3849 .expect("track should not be closed");
3850 assert_eq!(&frame.payload[..], b"world");
3851 }
3852
3853 #[tokio::test]
3854 async fn read_frame_preserves_timestamp() {
3855 let mut producer = track_producer("test", None);
3856 let mut consumer = producer.subscribe(None);
3857
3858 producer
3859 .write_frame(Timestamp::from_micros(20_000).unwrap(), b"hello".as_slice())
3860 .unwrap();
3861
3862 let frame = consumer
3863 .read_frame()
3864 .now_or_never()
3865 .expect("should not block")
3866 .expect("would have errored")
3867 .expect("track should not be closed");
3868 assert_eq!(frame.timestamp.as_micros(), 20_000);
3869 assert_eq!(&frame.payload[..], b"hello");
3870 }
3871
3872 #[tokio::test]
3873 async fn read_frame_skips_stalled_group_for_newer_ready_frame() {
3874 let mut producer = track_producer("test", None);
3875 let mut consumer = producer.subscribe(None);
3876
3877 let _stalled = producer.create_group(group::Info { sequence: 3 }).unwrap();
3879 let mut g5 = producer.create_group(group::Info { sequence: 5 }).unwrap();
3881 g5.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"later"))
3882 .unwrap();
3883 g5.finish().unwrap();
3884
3885 let frame = consumer
3887 .read_frame()
3888 .now_or_never()
3889 .expect("should not block on stalled earlier group")
3890 .expect("would have errored")
3891 .expect("track should not be closed");
3892 assert_eq!(&frame.payload[..], b"later");
3893 }
3894
3895 #[tokio::test]
3896 async fn read_frame_discards_rest_of_multi_frame_group() {
3897 let mut producer = track_producer("test", None);
3898 let mut consumer = producer.subscribe(None);
3899
3900 let mut g0 = producer.create_group(group::Info { sequence: 0 }).unwrap();
3902 g0.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"one"))
3903 .unwrap();
3904 g0.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"two"))
3905 .unwrap();
3906 g0.finish().unwrap();
3907
3908 producer.write_frame(Timestamp::ZERO, b"next".as_slice()).unwrap();
3910
3911 let frame = consumer
3912 .read_frame()
3913 .now_or_never()
3914 .expect("should not block")
3915 .expect("would have errored")
3916 .expect("track should not be closed");
3917 assert_eq!(&frame.payload[..], b"one");
3918
3919 let frame = consumer
3921 .read_frame()
3922 .now_or_never()
3923 .expect("should not block")
3924 .expect("would have errored")
3925 .expect("track should not be closed");
3926 assert_eq!(&frame.payload[..], b"next");
3927 }
3928
3929 #[tokio::test]
3930 async fn read_frame_waits_for_pending_group_after_finish() {
3931 let mut producer = track_producer("test", None);
3934 let mut consumer = producer.subscribe(None);
3935
3936 let mut g0 = producer.create_group(group::Info { sequence: 0 }).unwrap();
3937 producer.finish().unwrap();
3938
3939 assert!(
3941 consumer.read_frame().now_or_never().is_none(),
3942 "read_frame must block on a pending group even after finish()"
3943 );
3944
3945 g0.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"late"))
3947 .unwrap();
3948 let frame = consumer
3949 .read_frame()
3950 .now_or_never()
3951 .expect("should not block once a frame is written")
3952 .expect("would have errored")
3953 .expect("track should not be closed");
3954 assert_eq!(&frame.payload[..], b"late");
3955 }
3956
3957 #[tokio::test]
3958 async fn read_frame_respects_start_at() {
3959 let mut producer = track_producer("test", None);
3962 let mut consumer = producer.subscribe(None);
3963 consumer.start_at(5);
3964
3965 let mut g3 = producer.create_group(group::Info { sequence: 3 }).unwrap();
3967 g3.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"skip-me"))
3968 .unwrap();
3969 g3.finish().unwrap();
3970
3971 let mut g5 = producer.create_group(group::Info { sequence: 5 }).unwrap();
3972 g5.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"keep"))
3973 .unwrap();
3974 g5.finish().unwrap();
3975
3976 let frame = consumer
3977 .read_frame()
3978 .now_or_never()
3979 .expect("should not block")
3980 .expect("would have errored")
3981 .expect("track should not be closed");
3982 assert_eq!(&frame.payload[..], b"keep");
3983 }
3984
3985 #[tokio::test]
3986 async fn read_frame_returns_none_when_finished() {
3987 let mut producer = track_producer("test", None);
3988 let mut consumer = producer.subscribe(None);
3989
3990 producer.write_frame(Timestamp::ZERO, b"only".as_slice()).unwrap();
3991 producer.finish().unwrap();
3992
3993 let frame = consumer
3994 .read_frame()
3995 .now_or_never()
3996 .expect("should not block")
3997 .expect("would have errored")
3998 .expect("track should not be closed");
3999 assert_eq!(&frame.payload[..], b"only");
4000
4001 let done = consumer
4002 .read_frame()
4003 .now_or_never()
4004 .expect("should not block")
4005 .expect("would have errored");
4006 assert!(done.is_none());
4007 }
4008
4009 #[test]
4010 fn append_group_returns_bounds_exceeded_on_sequence_overflow() {
4011 let mut producer = track_producer("test", None);
4012 {
4013 let mut state = producer.state.write().ok().unwrap();
4014 state.max_sequence = Some(u64::MAX);
4015 }
4016
4017 assert!(matches!(producer.append_group(), Err(Error::BoundsExceeded(_))));
4018 }
4019
4020 #[tokio::test]
4021 async fn fetch_cache_hit() {
4022 let mut producer = track_producer("test", None);
4023
4024 let mut group = producer.append_group().unwrap(); group
4027 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hello"))
4028 .unwrap();
4029 group.finish().unwrap();
4030
4031 let dynamic = producer.dynamic();
4034 let consumer = producer.consume();
4035 assert!(consumer.peek_group(0).is_some());
4036 let mut g = consumer.fetch_group(0, None).await.unwrap();
4037 assert_eq!(g.sequence, 0);
4038 assert_eq!(&g.read_frame().await.unwrap().unwrap().payload[..], b"hello");
4039
4040 assert!(dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending());
4042 }
4043
4044 #[tokio::test]
4045 async fn fetch_miss_signals_dynamic() {
4046 let producer = track_producer("test", None);
4047 let dynamic = producer.dynamic();
4048 let consumer = producer.consume();
4049
4050 assert!(consumer.peek_group(5).is_none());
4054 let pending = consumer.fetch_group(5, group::Fetch::default().with_priority(7));
4055 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
4056
4057 let req = dynamic
4058 .requested_group()
4059 .now_or_never()
4060 .expect("should not block")
4061 .unwrap();
4062 assert_eq!(req.sequence(), 5);
4063 assert_eq!(req.priority(), 7);
4064
4065 let mut group = req.accept(None).unwrap();
4067 group
4068 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hi"))
4069 .unwrap();
4070 group.finish().unwrap();
4071
4072 let mut g = pending.await.unwrap();
4073 assert_eq!(g.sequence, 5);
4074 assert_eq!(&g.read_frame().await.unwrap().unwrap().payload[..], b"hi");
4075 }
4076
4077 #[tokio::test]
4078 async fn fetch_miss_rejects() {
4079 let producer = track_producer("test", None);
4080 let dynamic = producer.dynamic();
4081 let consumer = producer.consume();
4082
4083 let pending = consumer.fetch_group(5, None);
4084 let req = dynamic
4085 .requested_group()
4086 .now_or_never()
4087 .expect("should not block")
4088 .unwrap();
4089
4090 req.reject(Error::Cancel);
4091 assert!(matches!(pending.await, Err(Error::Cancel)));
4092 let fetch = producer.state.read().fetch.clone();
4093 assert!(fetch.read().is_empty());
4094 }
4095
4096 #[tokio::test]
4097 async fn fetch_miss_drop_rejects() {
4098 let producer = track_producer("test", None);
4099 let dynamic = producer.dynamic();
4100 let consumer = producer.consume();
4101
4102 let pending = consumer.fetch_group(5, None);
4103 let req = dynamic
4104 .requested_group()
4105 .now_or_never()
4106 .expect("should not block")
4107 .unwrap();
4108
4109 drop(req);
4110 assert!(matches!(pending.await, Err(Error::Dropped)));
4111 }
4112
4113 #[tokio::test]
4114 async fn fetch_reject_does_not_poison_retry() {
4115 let producer = track_producer("test", None);
4116 let dynamic = producer.dynamic();
4117 let consumer = producer.consume();
4118
4119 let pending = consumer.fetch_group(5, None);
4120 let req = dynamic
4121 .requested_group()
4122 .now_or_never()
4123 .expect("should not block")
4124 .unwrap();
4125 req.reject(Error::Cancel);
4126 assert!(matches!(pending.await, Err(Error::Cancel)));
4127
4128 let retry = consumer.fetch_group(5, None);
4129 let req = dynamic
4130 .requested_group()
4131 .now_or_never()
4132 .expect("should not block")
4133 .unwrap();
4134 let mut group = req.accept(None).unwrap();
4135 group
4136 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"retry"))
4137 .unwrap();
4138 group.finish().unwrap();
4139
4140 let mut group = retry.await.unwrap();
4141 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"retry");
4142 }
4143
4144 #[tokio::test]
4145 async fn fetch_coalesces_concurrent() {
4146 let producer = track_producer("test", None);
4147 let dynamic = producer.dynamic();
4148 let consumer = producer.consume();
4149
4150 let first = consumer.fetch_group(5, group::Fetch::default().with_priority(1));
4153 let second = consumer.fetch_group(5, group::Fetch::default().with_priority(7));
4154 assert!(kio::Pollable::poll(&*first, &kio::Waiter::noop()).is_pending());
4155
4156 let req = dynamic
4157 .requested_group()
4158 .now_or_never()
4159 .expect("should not block")
4160 .unwrap();
4161 assert_eq!(req.sequence(), 5);
4162 assert_eq!(req.priority(), 7);
4163 assert!(
4164 dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending(),
4165 "the second fetch queued a duplicate request"
4166 );
4167
4168 let third = consumer.fetch_group(5, None);
4170
4171 let mut group = req.accept(None).unwrap();
4173 group
4174 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hi"))
4175 .unwrap();
4176 group.finish().unwrap();
4177
4178 assert_eq!(first.await.unwrap().sequence, 5);
4179 assert_eq!(second.await.unwrap().sequence, 5);
4180 assert_eq!(third.await.unwrap().sequence, 5);
4181 }
4182
4183 #[tokio::test]
4184 async fn fetch_coalesced_reject_fails_all() {
4185 let producer = track_producer("test", None);
4186 let dynamic = producer.dynamic();
4187 let consumer = producer.consume();
4188
4189 let first = consumer.fetch_group(5, None);
4190 let second = consumer.fetch_group(5, None);
4191 let req = dynamic
4192 .requested_group()
4193 .now_or_never()
4194 .expect("should not block")
4195 .unwrap();
4196 req.reject(Error::Cancel);
4197
4198 assert!(matches!(first.await, Err(Error::Cancel)));
4199 assert!(matches!(second.await, Err(Error::Cancel)));
4200
4201 let retry = consumer.fetch_group(5, None);
4203 assert!(kio::Pollable::poll(&*retry, &kio::Waiter::noop()).is_pending());
4204 let req = dynamic
4205 .requested_group()
4206 .now_or_never()
4207 .expect("should not block")
4208 .unwrap();
4209 assert_eq!(req.sequence(), 5);
4210 }
4211
4212 #[tokio::test]
4213 async fn fetch_queued_fails_when_handlers_leave() {
4214 let producer = track_producer("test", None);
4215 let dynamic = producer.dynamic();
4216 let consumer = producer.consume();
4217
4218 let pending = consumer.fetch_group(5, None);
4220 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
4221 drop(dynamic);
4222 assert!(matches!(pending.await, Err(Error::NotFound)));
4223
4224 let fetch = producer.state.read().fetch.clone();
4226 assert!(fetch.read().is_empty());
4227 }
4228
4229 #[tokio::test]
4230 async fn fetch_miss_no_dynamic_not_found() {
4231 let mut producer = track_producer("test", None);
4234 producer.append_group().unwrap(); let consumer = producer.consume();
4236 assert!(matches!(consumer.fetch_group(5, None).await, Err(Error::NotFound)));
4237 }
4238
4239 #[tokio::test]
4240 async fn fetch_past_final_not_found() {
4241 let mut producer = track_producer("test", None);
4242 producer.append_group().unwrap(); producer.finish().unwrap(); let dynamic = producer.dynamic();
4248 let consumer = producer.consume();
4249 assert!(matches!(consumer.fetch_group(5, None).await, Err(Error::NotFound)));
4250
4251 assert!(dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending());
4253 }
4254
4255 fn pooled_producer(capacity: u64) -> (Producer, cache::Pool) {
4257 let pool = cache::Pool::new(capacity);
4258 let broadcast = broadcast::Info {
4259 origin: crate::origin::Info::default().with_pool(pool.clone()),
4260 ..Default::default()
4261 };
4262 let producer = Producer::new(Arc::new(broadcast), "test", None);
4263 (producer, pool)
4264 }
4265
4266 fn finished_group(producer: &mut Producer, size: usize) -> u64 {
4267 let mut group = producer.append_group().unwrap();
4268 group
4269 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; size]))
4270 .unwrap();
4271 group.finish().unwrap();
4272 group.sequence
4273 }
4274
4275 #[tokio::test]
4278 async fn debt_evicts_oldest_group() {
4279 tokio::time::pause();
4280
4281 let (mut producer, pool) = pooled_producer(10_000);
4283
4284 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let consumer = producer.consume();
4289 assert!(consumer.peek_group(0).is_none(), "oldest group is evicted");
4290 assert!(consumer.peek_group(2).is_some(), "latest group survives");
4291 assert!(pool.used() <= 21_000, "usage hovers near capacity: {}", pool.used());
4294
4295 let mut subscriber = producer.subscribe(None);
4297 assert!(subscriber.assert_group().sequence > 0, "evicted group is not delivered");
4298 }
4299
4300 #[tokio::test]
4302 async fn latest_group_never_evicted() {
4303 tokio::time::pause();
4304
4305 let (mut producer, pool) = pooled_producer(100);
4307 finished_group(&mut producer, 1000); assert!(pool.used() > 100, "the latest may exceed the budget");
4309
4310 finished_group(&mut producer, 1000); finished_group(&mut producer, 1000); let consumer = producer.consume();
4315 assert!(consumer.peek_group(0).is_none());
4316 let mut group = consumer.peek_group(2).expect("latest survives");
4317 assert_eq!(group.read_frame().await.unwrap().unwrap().payload.len(), 1000);
4318 }
4319
4320 #[tokio::test]
4324 async fn fetch_refresh_survives_eviction() {
4325 tokio::time::pause();
4326
4327 let (mut producer, _pool) = pooled_producer(10_000);
4328 let consumer = producer.consume();
4329
4330 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_secs(1)).await;
4332 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_secs(1)).await;
4334 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_millis(500)).await;
4336
4337 let mut fetched = consumer.fetch_group(0, None).await.unwrap();
4339 assert_eq!(fetched.read_frame().await.unwrap().unwrap().payload.len(), 3_000);
4340 tokio::time::advance(Duration::from_millis(500)).await;
4341
4342 finished_group(&mut producer, 3_000); tokio::time::advance(Duration::from_secs(1)).await;
4346 finished_group(&mut producer, 3_000); assert!(consumer.peek_group(0).is_some(), "refreshed group survives");
4349 assert!(consumer.peek_group(1).is_none(), "unread group is evicted instead");
4350 }
4351
4352 #[tokio::test]
4355 async fn eviction_aborts_readers() {
4356 tokio::time::pause();
4357
4358 let (mut producer, _pool) = pooled_producer(10_000);
4359 let mut subscriber = producer.subscribe(None);
4360
4361 finished_group(&mut producer, 10_000); let mut group0 = subscriber.assert_group();
4363
4364 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let read = group0.read_frame().await;
4368 assert!(matches!(read, Err(Error::Evicted)), "expected Evicted, got {read:?}");
4369 }
4370
4371 #[tokio::test]
4375 async fn small_writes_carry_debt() {
4376 tokio::time::pause();
4377
4378 let (mut producer, pool) = pooled_producer(22_000);
4379 let consumer = producer.consume();
4380
4381 finished_group(&mut producer, 20_000); for _ in 0..3 {
4386 finished_group(&mut producer, 1_000);
4387 }
4388 assert!(consumer.peek_group(0).is_some(), "debt smaller than the victim carries");
4389
4390 for _ in 0..20 {
4392 finished_group(&mut producer, 1_000);
4393 }
4394 assert!(
4395 consumer.peek_group(0).is_none(),
4396 "accumulated debt evicts the large group"
4397 );
4398 assert!(pool.used() <= 24_000, "usage hovers near capacity: {}", pool.used());
4401 }
4402
4403 #[tokio::test]
4407 async fn payment_capped_per_write() {
4408 tokio::time::pause();
4409
4410 let (mut producer, pool) = pooled_producer(1 << 40);
4411 for _ in 0..10 {
4412 finished_group(&mut producer, 1_000);
4413 }
4414
4415 pool.resize(100);
4417 let before = pool.used();
4418
4419 finished_group(&mut producer, 1_000);
4421
4422 let consumer = producer.consume();
4423 assert!(consumer.peek_group(0).is_none(), "the oldest groups are evicted");
4424 assert!(consumer.peek_group(1).is_none());
4425 assert!(consumer.peek_group(2).is_some(), "the backlog drains gradually");
4426 assert!(pool.used() > before - 4_000, "one write must not dump the backlog");
4427 }
4428
4429 #[tokio::test]
4433 async fn accept_preserves_write_accounting() {
4434 tokio::time::pause();
4435
4436 let pool = cache::Pool::new(12_000);
4437 let broadcast = broadcast::Info {
4438 origin: crate::origin::Info::default().with_pool(pool.clone()),
4439 ..Default::default()
4440 };
4441 let request = Request::new(Arc::new(broadcast), "test");
4442 let dynamic = request.dynamic();
4443 let consumer = request.consume();
4444
4445 let pending = consumer.fetch_group(0, None);
4447 let req = dynamic
4448 .requested_group()
4449 .now_or_never()
4450 .expect("should not block")
4451 .unwrap();
4452 let mut backfill = req.accept(None).unwrap();
4453 pending.await.unwrap();
4454 backfill
4455 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 30_000]))
4456 .unwrap();
4457
4458 let mut producer = request.accept(None);
4461 producer.append_group().unwrap().finish().unwrap();
4462 producer.append_group().unwrap().finish().unwrap();
4463
4464 assert!(
4465 producer.consume().peek_group(0).is_none(),
4466 "pre-accept backfill growth is reclaimed after accept"
4467 );
4468 assert!(pool.used() <= 13_000, "usage converges: {}", pool.used());
4469 }
4470
4471 #[tokio::test]
4474 async fn recreated_sequence_bounds_eviction_hints() {
4475 let (mut producer, _pool) = pooled_producer(1 << 40);
4476 producer.create_group(5u64.into()).unwrap().finish().unwrap();
4477
4478 for _ in 0..200 {
4479 let group = producer.create_group(1u64.into()).unwrap();
4480 group.abort(Error::Cancel).unwrap();
4481 }
4482
4483 let state = producer.state.read();
4484 assert!(
4485 state.evict.len() <= 2 * state.lookup.len() + EVICT_SLACK,
4486 "stale hints are compacted: {} entries for {} slots",
4487 state.evict.len(),
4488 state.lookup.len()
4489 );
4490 }
4491
4492 #[tokio::test]
4495 async fn same_tick_write_outranks_inserted() {
4496 tokio::time::pause();
4497
4498 let (mut producer, _pool) = pooled_producer(10_000);
4500
4501 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();
4508 assert!(consumer.peek_group(0).is_none(), "insert-only content pays first");
4509 assert!(consumer.peek_group(1).is_some(), "same-tick written content survives");
4510 }
4511
4512 #[tokio::test]
4515 async fn frame_only_writer_pays() {
4516 tokio::time::pause();
4517
4518 let (mut producer, pool) = pooled_producer(2_000);
4519 let mut demoted = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); demoted
4525 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 300_000]))
4526 .unwrap();
4527
4528 assert!(
4529 pool.used() <= 5_000,
4530 "the frame write settled the debt: {}",
4531 pool.used()
4532 );
4533 assert!(matches!(demoted.finish(), Err(Error::Evicted)));
4534 }
4535
4536 #[tokio::test]
4539 async fn each_track_owns_its_account() {
4540 let broadcast = Arc::new(broadcast::Info::default());
4541 let info = Info::default();
4542 let a = Producer::new(broadcast.clone(), "a", info.clone());
4543 let b = Producer::new(broadcast, "b", info);
4544
4545 let a = a.state.read().cache.clone();
4546 let b = b.state.read().cache.clone();
4547 assert!(!Arc::ptr_eq(&a, &b), "each track owns its account");
4548 }
4549
4550 #[tokio::test]
4553 async fn a_dynamic_defers_teardown() {
4554 let (mut producer, pool) = pooled_producer(1 << 40);
4555 let dynamic = producer.dynamic();
4556 finished_group(&mut producer, 100);
4557
4558 drop(producer);
4559 assert!(pool.used() > 0, "the handler still serves the cache");
4560
4561 drop(dynamic);
4562 assert_eq!(pool.used(), 0, "the last handle tears it down");
4563 }
4564
4565 #[tokio::test]
4571 async fn finished_track_frees_its_cache() {
4572 let (mut producer, pool) = pooled_producer(1 << 40);
4573 finished_group(&mut producer, 100);
4574 producer.finish().unwrap();
4575
4576 let state = producer.state.downgrade();
4577 drop(producer);
4578
4579 assert!(state.upgrade().is_none(), "the track state is freed");
4580 assert_eq!(pool.used(), 0, "so are its cached bytes");
4581 }
4582
4583 #[tokio::test]
4587 async fn teardown_ignores_a_settling_group() {
4588 let (mut producer, pool) = pooled_producer(1 << 40);
4589 finished_group(&mut producer, 100);
4590
4591 let settling = producer.state.downgrade().upgrade().expect("open");
4593 drop(producer);
4594
4595 assert_eq!(pool.used(), 0, "the abrupt teardown still released the cache");
4596 drop(settling);
4597 }
4598
4599 #[tokio::test]
4602 async fn cached_group_outlives_its_track() {
4603 let (mut producer, pool) = pooled_producer(1 << 40);
4604 let sequence = finished_group(&mut producer, 100);
4605 let group = producer.consume().peek_group(sequence).expect("cached");
4606 producer.finish().unwrap();
4607
4608 let state = producer.state.downgrade();
4609 drop(producer);
4610 assert!(state.upgrade().is_none(), "the track state is freed");
4611 assert!(pool.used() > 0, "the retained group keeps its own bytes");
4612
4613 drop(group);
4614 assert_eq!(pool.used(), 0, "which it releases when dropped");
4615 }
4616
4617 #[tokio::test]
4621 async fn pre_accept_backfill_settles_late_writes() {
4622 tokio::time::pause();
4623
4624 let pool = cache::Pool::new(2_000);
4625 let broadcast = broadcast::Info {
4626 origin: crate::origin::Info::default().with_pool(pool.clone()),
4627 ..Default::default()
4628 };
4629 let request = Request::new(Arc::new(broadcast), "test");
4630 let dynamic = request.dynamic();
4631 let consumer = request.consume();
4632
4633 let pending = consumer.fetch_group(0, None);
4635 let req = dynamic
4636 .requested_group()
4637 .now_or_never()
4638 .expect("should not block")
4639 .unwrap();
4640 let mut backfill = req.accept(None).unwrap();
4641 pending.await.unwrap();
4642
4643 let mut producer = request.accept(None);
4645 producer.append_group().unwrap().finish().unwrap();
4646
4647 backfill
4650 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 300_000]))
4651 .unwrap();
4652
4653 assert!(
4654 pool.used() <= 5_000,
4655 "the frame write settled the debt: {}",
4656 pool.used()
4657 );
4658 }
4659
4660 #[tokio::test]
4664 async fn write_restarts_retention_clock() {
4665 tokio::time::pause();
4666
4667 let (mut producer, _pool) = pooled_producer(1 << 40);
4668 let mut straggler = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
4673 straggler
4674 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 100]))
4675 .unwrap();
4676 producer.append_group().unwrap().finish().unwrap(); let consumer = producer.consume();
4679 assert!(consumer.peek_group(0).is_some(), "the write restarted the clock");
4680
4681 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
4683 producer.append_group().unwrap().finish().unwrap(); assert!(consumer.peek_group(0).is_none(), "idle content still expires");
4685 }
4686
4687 #[tokio::test]
4690 async fn refreshed_front_does_not_starve_expiry() {
4691 tokio::time::pause();
4692
4693 let (mut producer, _pool) = pooled_producer(1 << 40);
4694 let dynamic = producer.dynamic();
4695 let consumer = producer.consume();
4696
4697 producer.create_group(10u64.into()).unwrap().finish().unwrap();
4698 for sequence in 1..=5u64 {
4699 let pending = consumer.fetch_group(sequence, None);
4700 let req = dynamic
4701 .requested_group()
4702 .now_or_never()
4703 .expect("should not block")
4704 .unwrap();
4705 let mut group = req.accept(None).unwrap();
4706 group
4707 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 100]))
4708 .unwrap();
4709 group.finish().unwrap();
4710 pending.await.unwrap();
4711 }
4712
4713 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
4716 for sequence in 1..=4u64 {
4717 consumer.fetch_group(sequence, None).await.unwrap();
4718 }
4719
4720 for _ in 0..3 {
4722 producer.append_group().unwrap().finish().unwrap();
4723 }
4724 assert!(consumer.peek_group(5).is_none(), "expired backfill is reclaimed");
4725 assert!(consumer.peek_group(1).is_some(), "refreshed backfill survives");
4726 }
4727
4728 #[tokio::test]
4731 async fn recreated_sequence_delivered_once() {
4732 let (mut producer, _pool) = pooled_producer(1 << 40);
4733
4734 producer.create_group(0u64.into()).unwrap().finish().unwrap();
4735 let aborted = producer.create_group(1u64.into()).unwrap();
4736 aborted.abort(Error::Cancel).unwrap();
4737 producer.create_group(2u64.into()).unwrap().finish().unwrap();
4738 producer.create_group(1u64.into()).unwrap().finish().unwrap();
4739
4740 let mut subscriber = producer.subscribe(None);
4741 assert_eq!(subscriber.assert_group().sequence, 0);
4742 assert_eq!(subscriber.assert_group().sequence, 2);
4743 assert_eq!(
4744 subscriber.assert_group().sequence,
4745 1,
4746 "replacement arrives at its own position"
4747 );
4748 subscriber.assert_no_group();
4749 }
4750
4751 #[tokio::test]
4755 async fn datagrams_do_not_block_eviction() {
4756 tokio::time::pause();
4757
4758 let (mut producer, pool) = pooled_producer(1_000);
4759 for _ in 0..10 {
4760 finished_group(&mut producer, 1_000);
4761 producer.append_datagram(Timestamp::ZERO, &b"beat"[..]).unwrap();
4762 }
4763
4764 let consumer = producer.consume();
4765 assert!(consumer.peek_group(0).is_none(), "old groups still evict");
4766 assert!(
4767 pool.used() < 4 * 1_256,
4768 "interleaved datagrams must not bypass the budget: {}",
4769 pool.used()
4770 );
4771 }
4772
4773 #[tokio::test]
4777 async fn aborted_group_leaves_no_ghost_sample() {
4778 tokio::time::pause();
4779
4780 let (mut producer, pool) = pooled_producer(1 << 40);
4781 let group0 = producer.append_group().unwrap();
4782 producer.append_group().unwrap(); assert!(pool.average().is_some(), "demoted group is sampled");
4785 group0.abort(Error::Cancel).unwrap();
4786 assert_eq!(pool.average(), None, "the abort must remove the sample");
4787 }
4788
4789 #[tokio::test]
4792 async fn empty_groups_repay_overhead() {
4793 tokio::time::pause();
4794
4795 let (mut producer, pool) = pooled_producer(1_000);
4796 for _ in 0..100 {
4797 let mut group = producer.append_group().unwrap();
4798 group.finish().unwrap();
4799 }
4800
4801 assert!(
4802 pool.used() <= 3_000,
4803 "empty-group overhead must stay near the budget: {}",
4804 pool.used()
4805 );
4806 }
4807
4808 #[tokio::test]
4811 async fn growth_on_demoted_group_is_billed() {
4812 tokio::time::pause();
4813
4814 let (mut producer, pool) = pooled_producer(2_000);
4815 let mut straggler = producer.append_group().unwrap(); producer.append_group().unwrap().finish().unwrap(); straggler
4820 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 10_000]))
4821 .unwrap();
4822
4823 producer.append_group().unwrap().finish().unwrap(); let consumer = producer.consume();
4827 assert!(consumer.peek_group(0).is_none(), "the ballooned group is evicted");
4828 assert!(pool.used() <= 3_000, "growth is reclaimed: {}", pool.used());
4829 }
4830
4831 #[tokio::test]
4834 async fn refilled_sequence_stays_out_of_subscriptions() {
4835 let (mut producer, _pool) = pooled_producer(1 << 40);
4836 let dynamic = producer.dynamic();
4837 let consumer = producer.consume();
4838
4839 producer.create_group(0u64.into()).unwrap().finish().unwrap();
4840 let aborted = producer.create_group(1u64.into()).unwrap();
4841 aborted.abort(Error::Cancel).unwrap();
4842 producer.create_group(2u64.into()).unwrap().finish().unwrap();
4843
4844 let pending = consumer.fetch_group(1, None);
4847 let req = dynamic
4848 .requested_group()
4849 .now_or_never()
4850 .expect("should not block")
4851 .unwrap();
4852 let mut group = req.accept(None).unwrap();
4853 group
4854 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"backfill"))
4855 .unwrap();
4856 group.finish().unwrap();
4857 pending.await.unwrap();
4858
4859 assert!(consumer.peek_group(1).is_some());
4861 let mut subscriber = producer.subscribe(None);
4862 assert_eq!(subscriber.assert_group().sequence, 0);
4863 assert_eq!(subscriber.assert_group().sequence, 2);
4864 subscriber.assert_no_group();
4865 }
4866
4867 #[tokio::test]
4870 async fn expired_backfill_behind_refreshed_reclaimed() {
4871 tokio::time::pause();
4872
4873 let (mut producer, _pool) = pooled_producer(1 << 40);
4874 let dynamic = producer.dynamic();
4875 let consumer = producer.consume();
4876
4877 producer.create_group(5u64.into()).unwrap().finish().unwrap();
4878 for sequence in [2u64, 3u64] {
4879 let pending = consumer.fetch_group(sequence, None);
4880 let req = dynamic
4881 .requested_group()
4882 .now_or_never()
4883 .expect("should not block")
4884 .unwrap();
4885 let mut group = req.accept(None).unwrap();
4886 group
4887 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 100]))
4888 .unwrap();
4889 group.finish().unwrap();
4890 pending.await.unwrap();
4891 }
4892
4893 tokio::time::advance(Duration::from_secs(4)).await;
4895 consumer.fetch_group(2, None).await.unwrap();
4896 tokio::time::advance(DEFAULT_LATENCY_MAX - Duration::from_secs(2)).await;
4897 producer.create_group(6u64.into()).unwrap().finish().unwrap();
4898
4899 let consumer = producer.consume();
4900 assert!(consumer.peek_group(2).is_some(), "refreshed backfill survives");
4901 assert!(consumer.peek_group(3).is_none(), "expired backfill is reclaimed");
4902 }
4903
4904 #[tokio::test]
4907 async fn same_tick_fetch_protects() {
4908 tokio::time::pause();
4909
4910 let (mut producer, _pool) = pooled_producer(10_000);
4912 let consumer = producer.consume();
4913
4914 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();
4919
4920 finished_group(&mut producer, 3_000); finished_group(&mut producer, 3_000); assert!(consumer.peek_group(0).is_some(), "same-tick refresh protects");
4924 assert!(consumer.peek_group(1).is_none(), "the unread group dies instead");
4925 }
4926
4927 #[tokio::test]
4931 async fn refetched_latest_stays_protected() {
4932 tokio::time::pause();
4933
4934 let (mut producer, _pool) = pooled_producer(10_000);
4935 let dynamic = producer.dynamic();
4936 let consumer = producer.consume();
4937
4938 let straggler = producer.append_group().unwrap(); let latest = producer.append_group().unwrap(); latest.abort(Error::Cancel).unwrap();
4943
4944 let pending = consumer.fetch_group(1, None);
4946 let req = dynamic
4947 .requested_group()
4948 .now_or_never()
4949 .expect("should not block")
4950 .unwrap();
4951 let mut group = req.accept(None).unwrap();
4952 group
4953 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 1000]))
4954 .unwrap();
4955 group.finish().unwrap();
4956 pending.await.unwrap();
4957
4958 {
4961 let state = producer.state.read();
4962 assert!(state.lookup.contains_key(&1), "refetched group is cached");
4963 assert!(
4964 state.evict.iter().all(|(sequence, _)| *sequence != 1),
4965 "the live edge must not be an eviction candidate"
4966 );
4967 }
4968 drop(straggler);
4969 }
4970
4971 #[tokio::test]
4974 async fn eviction_allows_refetch() {
4975 tokio::time::pause();
4976
4977 let (mut producer, _pool) = pooled_producer(10_000);
4978 let dynamic = producer.dynamic();
4979
4980 finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); finished_group(&mut producer, 10_000); let consumer = producer.consume();
4985 assert!(consumer.peek_group(0).is_none());
4986 let pending = consumer.fetch_group(0, None);
4987
4988 let req = dynamic
4989 .requested_group()
4990 .now_or_never()
4991 .expect("should not block")
4992 .unwrap();
4993 assert_eq!(req.sequence(), 0);
4994
4995 let mut group = req.accept(None).unwrap();
4996 group
4997 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"refetched"))
4998 .unwrap();
4999 group.finish().unwrap();
5000
5001 let mut group = pending.await.unwrap();
5002 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"refetched");
5003 }
5004
5005 #[tokio::test]
5008 async fn fetched_backfill_not_subscribed() {
5009 let (mut producer, _pool) = pooled_producer(1 << 40);
5010 let dynamic = producer.dynamic();
5011 let consumer = producer.consume();
5012
5013 producer.create_group(5u64.into()).unwrap().finish().unwrap();
5015 producer.create_group(6u64.into()).unwrap().finish().unwrap();
5016
5017 let pending = consumer.fetch_group(2, None);
5019 let req = dynamic
5020 .requested_group()
5021 .now_or_never()
5022 .expect("should not block")
5023 .unwrap();
5024 let mut group = req.accept(None).unwrap();
5025 group
5026 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"backfill"))
5027 .unwrap();
5028 group.finish().unwrap();
5029 let mut fetched = pending.await.unwrap();
5030 assert_eq!(&fetched.read_frame().await.unwrap().unwrap().payload[..], b"backfill");
5031 assert!(consumer.peek_group(2).is_some(), "backfill is cached for later fetches");
5032
5033 let mut subscriber = producer.subscribe(None);
5035 assert_eq!(subscriber.assert_group().sequence, 5);
5036 assert_eq!(subscriber.assert_group().sequence, 6);
5037 subscriber.assert_no_group();
5038 }
5039
5040 #[tokio::test]
5043 async fn expired_backfill_reclaimed() {
5044 tokio::time::pause();
5045
5046 let (mut producer, pool) = pooled_producer(1 << 40);
5047 let dynamic = producer.dynamic();
5048 let consumer = producer.consume();
5049
5050 producer.create_group(5u64.into()).unwrap().finish().unwrap();
5051
5052 let pending = consumer.fetch_group(2, None);
5054 let req = dynamic
5055 .requested_group()
5056 .now_or_never()
5057 .expect("should not block")
5058 .unwrap();
5059 let mut group = req.accept(None).unwrap();
5060 group
5061 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 1000]))
5062 .unwrap();
5063 group.finish().unwrap();
5064 pending.await.unwrap();
5065 let used = pool.used();
5066
5067 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
5069 producer.create_group(6u64.into()).unwrap().finish().unwrap();
5070
5071 assert!(consumer.peek_group(2).is_none(), "expired backfill is reclaimed");
5072 assert!(pool.used() < used, "its bytes are released");
5073 }
5074
5075 #[tokio::test]
5076 async fn fetch_aborts_with_track() {
5077 let producer = track_producer("test", None);
5078 let dynamic = producer.dynamic();
5079 let consumer = producer.consume();
5080
5081 let pending = consumer.fetch_group(3, None);
5082 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
5083
5084 producer.abort(Error::Cancel).unwrap();
5085 assert!(pending.await.is_err());
5086 drop(dynamic);
5087 }
5088}