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::{HashSet, VecDeque},
25 sync::Arc,
26 task::{Poll, ready},
27 time::Duration,
28};
29
30pub const DEFAULT_LATENCY_MAX: Duration = Duration::from_secs(5);
32
33const MAX_DATAGRAM_AGE: Duration = Duration::from_millis(50);
39
40#[derive(Clone, Debug)]
49#[non_exhaustive]
50pub struct Info {
51 pub timescale: Timescale,
58 pub latency_max: Duration,
67 pub priority: u8,
70 pub ordered: bool,
74
75 pub(crate) broadcast: Arc<broadcast::Info>,
82}
83
84fn default_broadcast() -> Arc<broadcast::Info> {
88 static DEFAULT: std::sync::LazyLock<Arc<broadcast::Info>> =
89 std::sync::LazyLock::new(|| Arc::new(broadcast::Info::default()));
90 DEFAULT.clone()
91}
92
93impl Default for Info {
94 fn default() -> Self {
95 Self {
96 timescale: Timescale::default(),
97 latency_max: DEFAULT_LATENCY_MAX,
98 priority: 0,
99 ordered: false,
100 broadcast: default_broadcast(),
101 }
102 }
103}
104
105impl Info {
106 pub fn with_timescale(mut self, timescale: Timescale) -> Self {
111 self.timescale = timescale;
112 self
113 }
114
115 pub fn with_latency_max(mut self, latency_max: Duration) -> Self {
117 self.latency_max = latency_max;
118 self
119 }
120
121 pub fn with_priority(mut self, priority: u8) -> Self {
123 self.priority = priority;
124 self
125 }
126
127 pub fn with_ordered(mut self, ordered: bool) -> Self {
131 self.ordered = ordered;
132 self
133 }
134}
135
136#[derive(Default)]
137struct TrackState {
138 info: Option<Info>,
141
142 broadcast: Arc<broadcast::Info>,
146
147 latest_entry: Option<Arc<cache::Entry>>,
151
152 groups: VecDeque<Option<(group::Producer, web_async::time::Instant)>>,
154
155 datagrams: VecDeque<(Datagram, web_async::time::Instant)>,
159
160 datagram_offset: usize,
163
164 duplicates: HashSet<u64>,
168
169 offset: usize,
171
172 max_sequence: Option<u64>,
174
175 final_sequence: Option<u64>,
177
178 abort: Option<Error>,
180
181 subscriptions: kio::Shared<Subscriptions>,
185
186 fetch: kio::Shared<FetchState>,
189}
190
191type Subscriptions = Vec<kio::Consumer<Subscription>>;
193
194type FetchState = Requests<u64, PendingFetch>;
199
200struct PendingFetch {
202 priority: u8,
204
205 result: kio::Producer<FetchOutcome>,
210}
211
212#[derive(Default)]
215struct FetchOutcome {
216 rejected: Option<Error>,
217}
218
219impl TrackState {
220 fn poll_info(&self) -> Poll<Result<Info>> {
221 if let Some(info) = &self.info {
222 Poll::Ready(Ok(info.clone()))
223 } else {
224 Poll::Pending
225 }
226 }
227
228 fn poll_recv_group(&self, index: usize, min_sequence: u64) -> Poll<Result<Option<(group::Consumer, usize)>>> {
232 let start = index.saturating_sub(self.offset);
233 for (i, slot) in self.groups.iter().enumerate().skip(start) {
234 if let Some((group, _)) = slot
235 && group.sequence >= min_sequence
236 && !group.is_aborted()
237 {
238 return Poll::Ready(Ok(Some((group.consume(), self.offset + i))));
239 }
240 }
241
242 if self.is_complete() {
244 Poll::Ready(Ok(None))
245 } else if let Some(err) = &self.abort {
246 Poll::Ready(Err(err.clone()))
247 } else {
248 Poll::Pending
249 }
250 }
251
252 fn poll_recv_datagram(&self, index: usize) -> Poll<Result<Option<(Datagram, usize)>>> {
258 let start = index.saturating_sub(self.datagram_offset);
259 if let Some((datagram, _)) = self.datagrams.get(start) {
260 return Poll::Ready(Ok(Some((datagram.clone(), self.datagram_offset + start))));
261 }
262
263 if self.is_complete() {
265 Poll::Ready(Ok(None))
266 } else if let Some(err) = &self.abort {
267 Poll::Ready(Err(err.clone()))
268 } else {
269 Poll::Pending
270 }
271 }
272
273 fn push_datagram(&mut self, datagram: Datagram) {
275 let now = web_async::time::Instant::now();
276 self.datagrams.push_back((datagram, now));
277 while let Some((_, at)) = self.datagrams.front() {
278 if now.duration_since(*at) <= MAX_DATAGRAM_AGE {
279 break;
280 }
281 self.datagrams.pop_front();
282 self.datagram_offset += 1;
283 }
284 }
285
286 fn poll_read_frame(
290 &self,
291 index: usize,
292 next_sequence: u64,
293 waiter: &kio::Waiter,
294 ) -> Poll<Result<Option<(frame::Frame, usize, u64)>>> {
295 let start = index.saturating_sub(self.offset);
296 let mut pending_seen = false;
297 for (i, slot) in self.groups.iter().enumerate().skip(start) {
298 let Some((group, _)) = slot else { continue };
299 if group.sequence < next_sequence {
300 continue;
301 }
302
303 let mut consumer = group.consume();
304 match consumer.poll_read_frame(waiter) {
305 Poll::Ready(Ok(Some(frame))) => {
306 return Poll::Ready(Ok(Some((frame, self.offset + i, group.sequence))));
307 }
308 Poll::Ready(Ok(None)) => continue,
309 Poll::Ready(Err(_)) => continue,
312 Poll::Pending => {
313 pending_seen = true;
314 continue;
315 }
316 }
317 }
318
319 if pending_seen {
322 Poll::Pending
323 } else if self.is_complete() {
324 Poll::Ready(Ok(None))
325 } else if let Some(err) = &self.abort {
326 Poll::Ready(Err(err.clone()))
327 } else {
328 Poll::Pending
329 }
330 }
331
332 fn poll_next_in_range(
342 &self,
343 next_sequence: u64,
344 end_sequence: Option<u64>,
345 ) -> Poll<Result<Option<group::Consumer>>> {
346 if let Some(end) = end_sequence
350 && end < next_sequence
351 {
352 if let Some(err) = &self.abort {
353 return Poll::Ready(Err(err.clone()));
354 }
355 return Poll::Pending;
356 }
357
358 let mut best: Option<&group::Producer> = None;
359 for (group, _) in self.groups.iter().flatten() {
360 if group.sequence < next_sequence {
361 continue;
362 }
363 if let Some(end) = end_sequence
364 && group.sequence > end
365 {
366 continue;
367 }
368 if group.is_aborted() {
369 continue;
370 }
371 if best.is_none_or(|b| group.sequence < b.sequence) {
372 best = Some(group);
373 }
374 }
375
376 if let Some(group) = best {
377 return Poll::Ready(Ok(Some(group.consume())));
378 }
379
380 if let Some(err) = &self.abort {
382 return Poll::Ready(Err(err.clone()));
383 }
384 if let Some(fin) = self.final_sequence
387 && next_sequence >= fin
388 {
389 return Poll::Ready(Ok(None));
390 }
391 Poll::Pending
392 }
393
394 fn cached_group(&self, sequence: u64) -> Option<group::Consumer> {
398 self.groups
399 .iter()
400 .flatten()
401 .find(|(group, _)| group.sequence == sequence && !group.is_aborted())
402 .map(|(group, _)| group.consume())
403 }
404
405 fn latency_bound(&self) -> Option<Duration> {
408 self.info.as_ref().map(|info| info.latency_max)
409 }
410
411 fn poll_fetch_cached(&self, sequence: u64) -> Poll<Result<group::Consumer>> {
416 if let Some(group) = self.cached_group(sequence) {
417 return Poll::Ready(Ok(group));
418 }
419
420 if let Some(err) = &self.abort {
421 return Poll::Ready(Err(err.clone()));
422 }
423
424 if self.final_sequence.is_some_and(|fin| sequence >= fin) {
426 return Poll::Ready(Err(Error::NotFound));
427 }
428
429 Poll::Pending
430 }
431
432 fn evict_expired(&mut self, now: web_async::time::Instant, max_age: Duration) {
443 for slot in self.groups.iter_mut() {
444 let Some((group, created_at)) = slot else { continue };
445
446 if group.is_aborted() {
449 self.duplicates.remove(&group.sequence);
450 *slot = None;
451 continue;
452 }
453
454 if Some(group.sequence) == self.max_sequence {
455 continue;
456 }
457
458 if now.duration_since(*created_at) <= max_age {
459 break;
460 }
461
462 self.duplicates.remove(&group.sequence);
463 if let Some((group, _)) = slot.take() {
468 let _ = group.abort(Error::Old);
469 }
470 }
471
472 while let Some(None) = self.groups.front() {
474 self.groups.pop_front();
475 self.offset += 1;
476 }
477 }
478
479 fn pin_latest(&mut self, group: &group::Producer) {
482 if Some(group.sequence) != self.max_sequence {
483 return;
484 }
485 if let Some(prev) = self.latest_entry.take() {
486 prev.set_pinned(false);
487 }
488 if let Some(entry) = group.cache_entry() {
489 entry.set_pinned(true);
490 self.latest_entry = Some(entry);
491 }
492 }
493
494 fn set_final(&mut self, final_sequence: u64) -> Result<()> {
497 if self.final_sequence.is_some() {
498 return Err(Error::Closed);
499 }
500 if let Some(max) = self.max_sequence
501 && final_sequence <= max
502 {
503 return Err(Error::ProtocolViolation);
504 }
505 self.final_sequence = Some(final_sequence);
506 Ok(())
507 }
508
509 fn is_complete(&self) -> bool {
515 self.final_sequence
516 .is_some_and(|fin| self.max_sequence.map_or(0, |max| max.saturating_add(1)) >= fin)
517 }
518
519 fn poll_finished(&self) -> Poll<Result<u64>> {
520 if let Some(fin) = self.final_sequence {
521 Poll::Ready(Ok(fin))
522 } else if let Some(err) = &self.abort {
523 Poll::Ready(Err(err.clone()))
524 } else {
525 Poll::Pending
526 }
527 }
528
529 fn modify(producer: &kio::Producer<Self>) -> Result<kio::Mut<'_, Self>> {
530 producer.write().map_err(|r| r.abort.clone().unwrap_or(Error::Dropped))
531 }
532
533 fn replace_evicted(
537 &mut self,
538 sequence: u64,
539 track: Info,
540 now: web_async::time::Instant,
541 ) -> Option<Result<group::Producer>> {
542 let slot = self
543 .groups
544 .iter_mut()
545 .find(|slot| matches!(slot, Some((group, _)) if group.sequence == sequence))?;
546 let (existing, _) = slot.as_ref().unwrap();
547 if !existing.is_aborted() {
548 return Some(Err(Error::Duplicate));
549 }
550 let group = group::Producer::new(group::Info { sequence }, track);
551 *slot = Some((group.clone(), now));
552 self.pin_latest(&group);
555 Some(Ok(group))
556 }
557
558 fn insert_group_request(&mut self, sequence: u64, info: Option<Info>) -> Result<group::Producer> {
564 if let Some(err) = &self.abort {
565 return Err(err.clone());
566 }
567 if let Some(fin) = self.final_sequence
568 && sequence >= fin
569 {
570 return Err(Error::Closed);
571 }
572
573 let now = web_async::time::Instant::now();
576 let broadcast = self.broadcast.clone();
577 let info = self
578 .info
579 .get_or_insert_with(|| {
580 let mut info = info.unwrap_or_default();
581 info.broadcast = broadcast;
582 info
583 })
584 .clone();
585
586 if !self.duplicates.insert(sequence) {
587 return self
589 .replace_evicted(sequence, info, now)
590 .unwrap_or(Err(Error::Duplicate));
591 }
592
593 let latency_max = info.latency_max;
594 let group = group::Producer::new(group::Info { sequence }, info);
595 self.max_sequence = Some(self.max_sequence.unwrap_or(0).max(sequence));
596 self.groups.push_back(Some((group.clone(), now)));
597 self.pin_latest(&group);
598 self.evict_expired(now, latency_max);
599 Ok(group)
600 }
601}
602
603#[derive(Clone)]
605pub struct Producer {
606 name: Arc<str>,
607 broadcast: Arc<broadcast::Info>,
610 state: kio::Producer<TrackState>,
611 prev_subscription: Option<Subscription>,
612 stats: stats::Scope,
616}
617
618impl Producer {
619 pub(crate) fn new(
627 broadcast: Arc<broadcast::Info>,
628 name: impl Into<Arc<str>>,
629 info: impl Into<Option<Info>>,
630 ) -> Self {
631 let mut info = info.into().unwrap_or_default();
632 info.broadcast = broadcast.clone();
633 Self {
634 name: name.into(),
635 state: kio::Producer::new(TrackState {
636 info: Some(info),
637 broadcast: broadcast.clone(),
638 ..Default::default()
639 }),
640 broadcast,
641 prev_subscription: None,
642 stats: stats::Scope::default(),
643 }
644 }
645
646 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
650 scope.open_subscription();
651 self.stats = scope;
652 self
653 }
654
655 pub fn name(&self) -> &str {
657 &self.name
658 }
659
660 pub fn broadcast(&self) -> &broadcast::Info {
662 &self.broadcast
663 }
664
665 pub fn create_group(&mut self, group: group::Info) -> Result<group::Producer> {
667 let mut state = self.modify()?;
668 if let Some(fin) = state.final_sequence
669 && group.sequence >= fin
670 {
671 return Err(Error::Closed);
672 }
673 let info = state.info.as_ref().unwrap();
674 let track = info.clone();
675 let latency_max = info.latency_max;
676 let now = web_async::time::Instant::now();
677
678 if !state.duplicates.insert(group.sequence) {
679 return state
681 .replace_evicted(group.sequence, track, now)
682 .unwrap_or(Err(Error::Duplicate));
683 }
684
685 let group = group::Producer::new(group, track).with_meter(self.stats.meter());
686 state.max_sequence = Some(state.max_sequence.unwrap_or(0).max(group.sequence));
687 state.groups.push_back(Some((group.clone(), now)));
688 state.pin_latest(&group);
689 state.evict_expired(now, latency_max);
690
691 Ok(group)
692 }
693
694 pub fn append_group(&mut self) -> Result<group::Producer> {
696 let mut state = self.modify()?;
697 let sequence = match state.max_sequence {
698 Some(s) => s.checked_add(1).ok_or(coding::BoundsExceeded)?,
699 None => 0,
700 };
701 if let Some(fin) = state.final_sequence
702 && sequence >= fin
703 {
704 return Err(Error::Closed);
705 }
706
707 let info = state.info.as_ref().unwrap();
708 let track = info.clone();
709 let latency_max = info.latency_max;
710
711 let group = group::Producer::new(group::Info { sequence }, track).with_meter(self.stats.meter());
712
713 let now = web_async::time::Instant::now();
714 state.duplicates.insert(sequence);
715 state.max_sequence = Some(sequence);
716 state.groups.push_back(Some((group.clone(), now)));
717 state.pin_latest(&group);
718 state.evict_expired(now, latency_max);
719
720 Ok(group)
721 }
722
723 pub fn append_datagram<B: crate::IntoBytes>(&mut self, timestamp: Timestamp, payload: B) -> Result<u64> {
735 let payload = payload.into_bytes();
736 if payload.len() > super::datagram::MAX_DATAGRAM_PAYLOAD {
737 return Err(Error::FrameTooLarge);
738 }
739 let meter = self.stats.meter();
741 let mut state = self.modify()?;
742 let timescale = state.info.as_ref().unwrap().timescale;
744 let timestamp = timestamp.convert(timescale).map_err(|_| Error::TimestampMismatch)?;
745 let sequence = match state.max_sequence {
746 Some(s) => s.checked_add(1).ok_or(coding::BoundsExceeded)?,
747 None => 0,
748 };
749 if let Some(fin) = state.final_sequence
750 && sequence >= fin
751 {
752 return Err(Error::Closed);
753 }
754 state.max_sequence = Some(sequence);
755 meter.datagram(payload.len() as u64);
756 state.push_datagram(Datagram {
757 sequence,
758 timestamp,
759 payload,
760 });
761 Ok(sequence)
762 }
763
764 pub fn write_datagram(&mut self, mut datagram: Datagram) -> Result<()> {
770 if datagram.payload.len() > super::datagram::MAX_DATAGRAM_PAYLOAD {
771 return Err(Error::FrameTooLarge);
772 }
773 let meter = self.stats.meter();
775 let mut state = self.modify()?;
776 let timescale = state.info.as_ref().unwrap().timescale;
778 datagram.timestamp = datagram
779 .timestamp
780 .convert(timescale)
781 .map_err(|_| Error::TimestampMismatch)?;
782 if let Some(fin) = state.final_sequence
783 && datagram.sequence >= fin
784 {
785 return Err(Error::Closed);
786 }
787 state.max_sequence = Some(state.max_sequence.unwrap_or(0).max(datagram.sequence));
788 meter.datagram(datagram.payload.len() as u64);
789 state.push_datagram(datagram);
790 Ok(())
791 }
792
793 pub fn write_frame<B: crate::IntoBytes>(&mut self, timestamp: Timestamp, frame: B) -> Result<()> {
798 let mut group = self.append_group()?;
799 group.write_frame(timestamp, frame)?;
800 group.finish()?;
801 Ok(())
802 }
803
804 pub fn finish(&mut self) -> Result<()> {
810 let mut state = self.modify()?;
811 let final_sequence = match state.max_sequence {
812 Some(max) => max.checked_add(1).ok_or(coding::BoundsExceeded)?,
813 None => 0,
814 };
815 state.set_final(final_sequence)
816 }
817
818 pub fn finish_at(&mut self, final_sequence: u64) -> Result<()> {
831 self.modify()?.set_final(final_sequence)
832 }
833
834 pub fn final_sequence(&self) -> Option<u64> {
839 self.state.read().final_sequence
840 }
841
842 pub fn abort(self, err: Error) -> Result<()> {
853 let mut guard = self.modify()?;
854 guard.abort = Some(err);
855 guard.groups.clear();
856 guard.datagrams.clear();
857 guard.duplicates.clear();
858 guard.latest_entry = None;
859 guard.close();
860 Ok(())
861 }
862
863 pub async fn unused(&self) -> Result<()> {
865 self.state.unused().await.map_err(|_| self.abort_reason())
866 }
867
868 pub async fn used(&self) -> Result<()> {
870 self.state.used().await.map_err(|_| self.abort_reason())
871 }
872
873 pub async fn closed(&self) -> Error {
875 kio::wait(|waiter| self.poll_closed(waiter)).await
876 }
877
878 pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<Error> {
880 self.state.poll_closed(waiter).map(|()| self.abort_reason())
881 }
882
883 fn abort_reason(&self) -> Error {
885 self.state.read().abort.clone().unwrap_or(Error::Dropped)
886 }
887
888 pub fn is_closed(&self) -> bool {
890 self.state.read().is_closed()
891 }
892
893 pub fn latest(&self) -> Option<u64> {
895 self.state.read().max_sequence
896 }
897
898 pub fn is_clone(&self, other: &Self) -> bool {
900 self.state.same_channel(&other.state)
901 }
902
903 pub(crate) fn weak(&self) -> TrackWeak {
905 TrackWeak {
906 name: self.name.clone(),
907 state: self.state.weak(),
908 }
909 }
910
911 pub fn demand(&self) -> Demand {
919 Demand {
920 name: self.name.clone(),
921 state: self.state.weak(),
922 }
923 }
924
925 pub fn consume(&self) -> Consumer {
930 Consumer::plain(self.name.clone(), self.state.consume())
931 }
932
933 pub fn subscribe(&self, subscription: impl Into<Option<Subscription>>) -> Subscriber {
938 let preferences = subscription.into().unwrap_or_default();
939
940 let info = self
945 .state
946 .read()
947 .info
948 .as_ref()
949 .expect("producer always has info")
950 .clone();
951 let subscription = kio::Producer::new(preferences);
952 register_subscription(self.state.read(), &subscription);
953
954 Subscriber {
955 name: self.name.clone(),
956 info,
957 inner: SubscriberKind::Plain(PlainSubscriber {
958 state: self.state.consume(),
959 subscription,
960 index: 0,
961 datagram_index: 0,
962 min_sequence: 0,
963 next_sequence: 0,
964 end_sequence: None,
965 }),
966 stats: stats::Scope::default(),
968 _stats_sub: stats::Subscription::default(),
969 }
970 }
971
972 pub async fn subscription_changed(&mut self) -> Result<Option<Subscription>> {
978 kio::wait(|waiter| self.poll_subscription_changed(waiter)).await
979 }
980
981 pub fn subscription(&self) -> Option<Subscription> {
989 let state = self.state.read();
990 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
991 drop(state);
992 snapshot_subscription(&subs, bound)
993 }
994
995 pub fn poll_subscription_changed(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Subscription>>> {
999 if self.state.poll_closed(waiter).is_ready() {
1002 let abort = self.state.read().abort.clone();
1003 return Poll::Ready(Err(abort.unwrap_or(Error::Dropped)));
1004 }
1005
1006 let state = self.state.read();
1008 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
1009 drop(state);
1010
1011 let prev = &self.prev_subscription;
1012 let mut combined = None;
1013 let mut guard = match subs.poll(waiter, |subs| {
1014 let next = combined_subscription(subs, bound, waiter);
1015 if &next == prev {
1016 Poll::Pending
1017 } else {
1018 combined = next;
1019 Poll::Ready(())
1020 }
1021 }) {
1022 Poll::Ready(guard) => guard,
1023 Poll::Pending => return Poll::Pending,
1024 };
1025 guard.retain(|sub| !sub.is_closed());
1027 drop(guard);
1028 self.prev_subscription = combined.clone();
1029 Poll::Ready(Ok(combined))
1030 }
1031
1032 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
1034 self.state.poll_unused(waiter).map(|_| ())
1035 }
1036
1037 pub fn dynamic(&self) -> Dynamic {
1041 Dynamic::new(self.name.clone(), self.state.clone())
1042 }
1043
1044 fn modify(&self) -> Result<kio::Mut<'_, TrackState>> {
1045 TrackState::modify(&self.state)
1046 }
1047}
1048
1049fn poll_requested_group(
1053 state: &kio::Producer<TrackState>,
1054 fetch: &kio::Shared<FetchState>,
1055 waiter: &kio::Waiter,
1056) -> Poll<Result<GroupRequest>> {
1057 if let Poll::Ready(mut guard) = fetch.poll(waiter, |fetch| {
1059 if fetch.has_queued() {
1060 Poll::Ready(())
1061 } else {
1062 Poll::Pending
1063 }
1064 }) {
1065 let sequence = guard.pop().expect("predicate guaranteed a request");
1066 let pending = guard.get(&sequence).expect("popped key must be pending");
1070 let priority = pending.priority;
1071 let result = pending.result.clone();
1072 drop(guard);
1073 return Poll::Ready(Ok(GroupRequest {
1074 state: state.clone(),
1075 fetch: fetch.clone(),
1076 sequence,
1077 priority,
1078 result,
1079 done: false,
1080 }));
1081 }
1082
1083 match state.poll_ref(waiter, |state| match &state.abort {
1085 Some(err) => Poll::Ready(err.clone()),
1086 None => Poll::Pending,
1087 }) {
1088 Poll::Ready(Ok(err)) => Poll::Ready(Err(err)),
1089 Poll::Ready(Err(closed)) => Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped))),
1090 Poll::Pending => Poll::Pending,
1091 }
1092}
1093
1094pub struct Dynamic {
1104 name: Arc<str>,
1105 state: kio::Producer<TrackState>,
1107 fetch: kio::Shared<FetchState>,
1109}
1110
1111impl Dynamic {
1112 fn new(name: Arc<str>, state: kio::Producer<TrackState>) -> Self {
1113 let fetch = state.read().fetch.clone();
1114 fetch.lock().add_handler();
1115 Self { name, state, fetch }
1116 }
1117
1118 pub fn name(&self) -> &str {
1120 &self.name
1121 }
1122
1123 pub async fn requested_group(&self) -> Result<GroupRequest> {
1129 kio::wait(|waiter| self.poll_requested_group(waiter)).await
1130 }
1131
1132 pub fn poll_requested_group(&self, waiter: &kio::Waiter) -> Poll<Result<GroupRequest>> {
1134 poll_requested_group(&self.state, &self.fetch, waiter)
1135 }
1136
1137 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
1139 self.state.poll_unused(waiter).map(|_| ())
1140 }
1141}
1142
1143impl Clone for Dynamic {
1144 fn clone(&self) -> Self {
1145 self.fetch.lock().add_handler();
1147 Self {
1148 name: self.name.clone(),
1149 state: self.state.clone(),
1150 fetch: self.fetch.clone(),
1151 }
1152 }
1153}
1154
1155impl Drop for Dynamic {
1156 fn drop(&mut self) {
1157 let mut fetch = self.fetch.lock();
1163 if fetch.remove_handler() {
1164 fetch.drain_queued();
1165 }
1166 }
1167}
1168
1169impl Drop for Producer {
1170 fn drop(&mut self) {
1171 if !self.state.is_last() {
1176 return;
1177 }
1178 self.stats.close_subscription();
1180 if let Ok(mut state) = self.state.write()
1181 && state.final_sequence.is_none()
1182 {
1183 tracing::warn!(
1187 track = %self.name(),
1188 "track::Producer dropped without finish() or abort()"
1189 );
1190 state.groups.clear();
1191 state.datagrams.clear();
1192 state.duplicates.clear();
1193 state.latest_entry = None;
1194 }
1195 }
1196}
1197
1198fn combined_subscription(subs: &Subscriptions, bound: Option<Duration>, waiter: &kio::Waiter) -> Option<Subscription> {
1204 let mut combined = None;
1205 for sub in subs.iter() {
1206 if sub.is_closed() {
1211 continue;
1212 }
1213 if let Poll::Ready(Ok(sub)) = sub.poll(waiter, |sub| sub.poll_combined(&combined)) {
1214 combined = Some(sub);
1215 }
1216 }
1217 clamp_combined(combined, bound)
1218}
1219
1220fn snapshot_subscription(subs: &kio::Shared<Subscriptions>, bound: Option<Duration>) -> Option<Subscription> {
1222 let mut combined: Option<Subscription> = None;
1223 for sub in subs.read().iter() {
1224 if sub.is_closed() {
1226 continue;
1227 }
1228 if let Poll::Ready(merged) = sub.read().poll_combined(&combined) {
1229 combined = Some(merged);
1230 }
1231 }
1232 clamp_combined(combined, bound)
1233}
1234
1235fn clamp_combined(combined: Option<Subscription>, bound: Option<Duration>) -> Option<Subscription> {
1243 let mut combined = combined?;
1244 if let Some(bound) = bound {
1245 combined.latency_max = combined.latency_max.min(bound);
1246 }
1247 Some(combined)
1248}
1249
1250fn register_subscription(state: kio::Ref<'_, TrackState>, subscription: &kio::Producer<Subscription>) {
1254 if state.is_closed() {
1255 return;
1256 }
1257 let subs = state.subscriptions.clone();
1258 drop(state);
1259 subs.lock().push(subscription.consume());
1260}
1261
1262#[derive(Clone)]
1264pub(crate) struct TrackWeak {
1265 name: Arc<str>,
1266 state: kio::ProducerWeak<TrackState>,
1267}
1268
1269impl TrackWeak {
1270 pub fn consume(&self) -> Consumer {
1271 Consumer::plain(self.name.clone(), self.state.consume())
1272 }
1273
1274 pub(crate) fn name(&self) -> &Arc<str> {
1277 &self.name
1278 }
1279
1280 pub(crate) fn is_used(&self) -> bool {
1283 !self.state.is_closed() && self.state.is_used()
1284 }
1285
1286 pub(crate) fn poll_used(&self, waiter: &kio::Waiter) {
1289 let _ = self.state.poll_used(waiter);
1290 }
1291
1292 pub(crate) fn poll_unused(&self, waiter: &kio::Waiter) {
1295 let _ = self.state.poll_unused(waiter);
1296 }
1297}
1298
1299impl super::WeakEntry for TrackWeak {
1300 fn is_closed(&self) -> bool {
1301 self.state.is_closed()
1302 }
1303
1304 fn same_channel(&self, other: &Self) -> bool {
1305 self.state.same_channel(&other.state)
1306 }
1307}
1308
1309#[derive(Clone)]
1318pub struct Demand {
1319 name: Arc<str>,
1320 state: kio::ProducerWeak<TrackState>,
1321}
1322
1323impl Demand {
1324 pub fn name(&self) -> &str {
1326 &self.name
1327 }
1328
1329 pub async fn used(&self) -> Result<()> {
1331 self.state.used().await.map_err(|_| self.abort_reason())
1332 }
1333
1334 pub async fn unused(&self) -> Result<()> {
1336 self.state.unused().await.map_err(|_| self.abort_reason())
1337 }
1338
1339 pub async fn closed(&self) -> Error {
1341 self.state.closed().await;
1342 self.abort_reason()
1343 }
1344
1345 fn abort_reason(&self) -> Error {
1347 self.state.read().abort.clone().unwrap_or(Error::Dropped)
1348 }
1349}
1350
1351#[derive(Clone)]
1362pub struct Consumer {
1363 name: Arc<str>,
1364 inner: ConsumerKind,
1365 stats: stats::Scope,
1368}
1369
1370#[derive(Clone)]
1371enum ConsumerKind {
1372 Plain(kio::Consumer<TrackState>),
1373 Spliced(super::resume::Consumer),
1374}
1375
1376impl Consumer {
1377 fn plain(name: Arc<str>, state: kio::Consumer<TrackState>) -> Self {
1378 Self {
1379 name,
1380 inner: ConsumerKind::Plain(state),
1381 stats: stats::Scope::default(),
1382 }
1383 }
1384
1385 pub(crate) fn spliced(name: Arc<str>, resume: super::resume::Consumer) -> Self {
1387 Self {
1388 name,
1389 inner: ConsumerKind::Spliced(resume),
1390 stats: stats::Scope::default(),
1391 }
1392 }
1393
1394 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
1397 self.stats = scope;
1398 self
1399 }
1400
1401 pub fn name(&self) -> &str {
1403 &self.name
1404 }
1405
1406 pub fn subscribe(&self, subscription: impl Into<Option<Subscription>>) -> kio::Pending<Subscribing> {
1412 let subscription = kio::Producer::new(subscription.into().unwrap_or_default());
1413
1414 let inner = match &self.inner {
1415 ConsumerKind::Plain(state) => {
1416 register_subscription(state.read(), &subscription);
1419 SubscribingKind::Plain(state.clone())
1420 }
1421 ConsumerKind::Spliced(resume) => SubscribingKind::Spliced(resume.clone()),
1423 };
1424
1425 kio::Pending::new(Subscribing {
1426 name: self.name.clone(),
1427 inner,
1428 subscription,
1429 stats: self.stats.clone(),
1430 })
1431 }
1432
1433 #[cfg(test)]
1437 pub(crate) fn peek_group(&self, sequence: u64) -> Option<group::Consumer> {
1438 match &self.inner {
1439 ConsumerKind::Plain(state) => state.read().cached_group(sequence),
1440 ConsumerKind::Spliced(_) => None,
1443 }
1444 }
1445
1446 pub fn fetch_group(&self, sequence: u64, options: impl Into<Option<group::Fetch>>) -> kio::Pending<Fetching> {
1458 let options = options.into().unwrap_or_default();
1459
1460 self.stats.fetch();
1464
1465 let state = match &self.inner {
1466 ConsumerKind::Plain(state) => state,
1467 ConsumerKind::Spliced(resume) => {
1470 return kio::Pending::new(Fetching {
1471 inner: FetchingKind::Spliced(resume.fetch_group(sequence, options)),
1472 stats: self.stats.clone(),
1473 });
1474 }
1475 };
1476
1477 let mut result = None;
1478
1479 let (fetch, unresolved) = {
1483 let state = state.read();
1484 (state.fetch.clone(), state.poll_fetch_cached(sequence).is_pending())
1485 };
1486
1487 if unresolved {
1488 let mut fetch = fetch.lock();
1489 if let Some(pending) = fetch.join(&sequence) {
1490 pending.priority = pending.priority.max(options.priority);
1493 result = Some(pending.result.consume());
1494 } else {
1495 let producer = kio::Producer::<FetchOutcome>::default();
1499 let consumer = producer.consume();
1500 let attempt = PendingFetch {
1501 priority: options.priority,
1502 result: producer,
1503 };
1504 if fetch.insert(sequence, attempt).is_ok() {
1505 result = Some(consumer);
1506 }
1507 }
1508 }
1509
1510 kio::Pending::new(Fetching {
1511 inner: FetchingKind::Plain {
1512 state: state.clone(),
1513 fetch,
1514 sequence,
1515 result,
1516 },
1517 stats: self.stats.clone(),
1518 })
1519 }
1520
1521 pub fn info(&self) -> kio::Pending<Querying> {
1528 kio::Pending::new(Querying {
1529 inner: match &self.inner {
1530 ConsumerKind::Plain(state) => QueryingKind::Plain(state.clone()),
1531 ConsumerKind::Spliced(resume) => QueryingKind::Spliced(resume.clone()),
1532 },
1533 })
1534 }
1535
1536 pub fn latest(&self) -> Option<u64> {
1538 match &self.inner {
1539 ConsumerKind::Plain(state) => state.read().max_sequence,
1540 ConsumerKind::Spliced(resume) => resume.latest(),
1541 }
1542 }
1543
1544 pub(crate) fn poll_complete(&self, waiter: &kio::Waiter) -> Poll<Result<()>> {
1549 let ConsumerKind::Plain(state) = &self.inner else {
1550 return Poll::Pending;
1552 };
1553 match ready!(state.poll(waiter, |state| {
1554 if state.is_complete() {
1555 Poll::Ready(())
1556 } else {
1557 Poll::Pending
1558 }
1559 })) {
1560 Ok(_) => Poll::Ready(Ok(())),
1561 Err(closed) => Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped))),
1564 }
1565 }
1566}
1567
1568pub struct Subscribing {
1571 name: Arc<str>,
1572 inner: SubscribingKind,
1573 subscription: kio::Producer<Subscription>,
1574 stats: stats::Scope,
1575}
1576
1577enum SubscribingKind {
1578 Plain(kio::Consumer<TrackState>),
1579 Spliced(super::resume::Consumer),
1580}
1581
1582impl Subscribing {
1583 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<Subscriber>> {
1586 match &self.inner {
1587 SubscribingKind::Plain(state) => {
1588 let info = ready!(state.poll(waiter, |state| state.poll_info()))
1590 .map_err(|e| e.abort.clone().unwrap_or(Error::Dropped))??;
1591
1592 Poll::Ready(Ok(Subscriber {
1593 name: self.name.clone(),
1594 info,
1595 inner: SubscriberKind::Plain(PlainSubscriber {
1596 state: state.clone(),
1597 subscription: self.subscription.clone(),
1598 index: 0,
1599 datagram_index: 0,
1600 min_sequence: 0,
1601 next_sequence: 0,
1602 end_sequence: None,
1603 }),
1604 stats: self.stats.clone(),
1605 _stats_sub: self.stats.subscribe(),
1606 }))
1607 }
1608 SubscribingKind::Spliced(resume) => {
1609 let info = ready!(resume.poll_info(waiter))?;
1612
1613 Poll::Ready(Ok(Subscriber {
1614 name: self.name.clone(),
1615 info,
1616 inner: SubscriberKind::Spliced(Box::new(resume.subscribe_shared(self.subscription.clone()))),
1617 stats: self.stats.clone(),
1618 _stats_sub: self.stats.subscribe(),
1619 }))
1620 }
1621 }
1622 }
1623
1624 pub fn update(&mut self, subscription: Subscription) -> Result<()> {
1629 let mut state = self.subscription.write().map_err(|_| Error::Closed)?;
1630 *state = subscription;
1631 Ok(())
1632 }
1633}
1634
1635impl kio::Pollable for Subscribing {
1636 type Output = Result<Subscriber>;
1637
1638 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
1639 self.poll_ok(waiter)
1640 }
1641}
1642
1643pub struct Querying {
1646 inner: QueryingKind,
1647}
1648
1649enum QueryingKind {
1650 Plain(kio::Consumer<TrackState>),
1651 Spliced(super::resume::Consumer),
1652}
1653
1654impl Querying {
1655 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<Info>> {
1657 match &self.inner {
1658 QueryingKind::Plain(state) => {
1659 let info = ready!(state.poll(waiter, |state| state.poll_info()))
1661 .map_err(|e| e.abort.clone().unwrap_or(Error::Dropped))??;
1662 Poll::Ready(Ok(info))
1663 }
1664 QueryingKind::Spliced(resume) => resume.poll_info(waiter),
1665 }
1666 }
1667}
1668
1669impl kio::Pollable for Querying {
1670 type Output = Result<Info>;
1671
1672 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
1673 self.poll_ok(waiter)
1674 }
1675}
1676
1677pub struct GroupRequest {
1686 state: kio::Producer<TrackState>,
1687 fetch: kio::Shared<FetchState>,
1689 sequence: u64,
1690 priority: u8,
1691 result: kio::Producer<FetchOutcome>,
1693 done: bool,
1694}
1695
1696impl GroupRequest {
1697 pub fn sequence(&self) -> u64 {
1699 self.sequence
1700 }
1701
1702 pub fn priority(&self) -> u8 {
1704 self.priority
1705 }
1706
1707 pub fn accept(mut self, info: impl Into<Option<Info>>) -> Result<group::Producer> {
1715 self.done = true;
1716 let res = TrackState::modify(&self.state)
1720 .and_then(|mut state| state.insert_group_request(self.sequence, info.into()));
1721 self.remove();
1722 res
1723 }
1724
1725 pub fn reject(mut self, err: Error) {
1727 self.done = true;
1728 self.remove();
1731 if let Ok(mut outcome) = self.result.write() {
1732 outcome.rejected = Some(err);
1733 }
1734 }
1735
1736 fn remove(&self) {
1739 self.fetch
1740 .lock()
1741 .remove_if(&self.sequence, |pending| pending.result.same_channel(&self.result));
1742 }
1743}
1744
1745impl Drop for GroupRequest {
1746 fn drop(&mut self) {
1747 if self.done {
1748 return;
1749 }
1750 self.remove();
1751 if let Ok(mut outcome) = self.result.write() {
1752 outcome.rejected = Some(Error::Dropped);
1753 }
1754 }
1755}
1756
1757pub struct Fetching {
1763 inner: FetchingKind,
1764 stats: stats::Scope,
1767}
1768
1769enum FetchingKind {
1770 Plain {
1771 state: kio::Consumer<TrackState>,
1772 fetch: kio::Shared<FetchState>,
1773 sequence: u64,
1774 result: Option<kio::Consumer<FetchOutcome>>,
1776 },
1777 Spliced(kio::Pending<super::resume::Fetching>),
1779}
1780
1781impl kio::Pollable for Fetching {
1782 type Output = Result<group::Consumer>;
1783
1784 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
1785 let (state, fetch, sequence, result) = match &self.inner {
1786 FetchingKind::Plain {
1787 state,
1788 fetch,
1789 sequence,
1790 result,
1791 } => (state, fetch, *sequence, result.as_ref()),
1792 FetchingKind::Spliced(spliced) => {
1793 return kio::Pollable::poll(&**spliced, waiter)
1796 .map(|res| res.map(|group| group.with_meter(self.stats.meter())));
1797 }
1798 };
1799
1800 match state.poll(waiter, |state| state.poll_fetch_cached(sequence)) {
1803 Poll::Ready(Ok(res)) => return Poll::Ready(res.map(|group| group.with_meter(self.stats.meter()))),
1804 Poll::Ready(Err(closed)) => {
1805 return Poll::Ready(Err(closed.abort.clone().unwrap_or(Error::Dropped)));
1806 }
1807 Poll::Pending => {}
1808 }
1809
1810 let Some(result) = result else {
1812 return match fetch.poll(waiter, |fetch| match fetch.has_handlers() {
1815 false => Poll::Ready(()),
1816 true => Poll::Pending,
1817 }) {
1818 Poll::Ready(_guard) => Poll::Ready(Err(Error::NotFound)),
1819 Poll::Pending => Poll::Pending,
1820 };
1821 };
1822
1823 match result.poll(waiter, |outcome| match &outcome.rejected {
1826 Some(err) => Poll::Ready(err.clone()),
1827 None => Poll::Pending,
1828 }) {
1829 Poll::Ready(Ok(err)) => Poll::Ready(Err(err)),
1830 Poll::Ready(Err(_closed)) => Poll::Ready(Err(Error::NotFound)),
1831 Poll::Pending => Poll::Pending,
1832 }
1833 }
1834}
1835
1836pub struct Subscriber {
1859 name: Arc<str>,
1860 info: Info,
1861 inner: SubscriberKind,
1862 stats: stats::Scope,
1865 _stats_sub: stats::Subscription,
1868}
1869
1870enum SubscriberKind {
1871 Plain(PlainSubscriber),
1872 Spliced(Box<super::resume::Subscriber>),
1874}
1875
1876struct PlainSubscriber {
1878 state: kio::Consumer<TrackState>,
1879
1880 subscription: kio::Producer<Subscription>,
1881 index: usize,
1883 datagram_index: usize,
1885 min_sequence: u64,
1887 next_sequence: u64,
1890 end_sequence: Option<u64>,
1895}
1896
1897impl PlainSubscriber {
1898 fn poll<F, R>(&self, waiter: &kio::Waiter, f: F) -> Poll<Result<R>>
1900 where
1901 F: Fn(&kio::Ref<'_, TrackState>) -> Poll<Result<R>>,
1902 {
1903 Poll::Ready(match ready!(self.state.poll(waiter, f)) {
1904 Ok(res) => res,
1905 Err(state) => Err(state.abort.clone().unwrap_or(Error::Dropped)),
1907 })
1908 }
1909
1910 fn poll_recv_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
1911 let Some((consumer, found_index)) =
1912 ready!(self.poll(waiter, |state| state.poll_recv_group(self.index, self.min_sequence))?)
1913 else {
1914 return Poll::Ready(Ok(None));
1915 };
1916
1917 self.index = found_index + 1;
1918 Poll::Ready(Ok(Some(consumer)))
1919 }
1920
1921 fn poll_recv_datagram(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Datagram>>> {
1922 let Some((datagram, found_index)) =
1923 ready!(self.poll(waiter, |state| state.poll_recv_datagram(self.datagram_index))?)
1924 else {
1925 return Poll::Ready(Ok(None));
1926 };
1927
1928 self.datagram_index = found_index + 1;
1929 self.next_sequence = self.next_sequence.max(datagram.sequence.saturating_add(1));
1930 Poll::Ready(Ok(Some(datagram)))
1931 }
1932
1933 fn poll_next_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
1934 let floor = self.next_sequence.max(self.min_sequence);
1935 let Some(group) = ready!(self.poll(waiter, |state| state.poll_next_in_range(floor, self.end_sequence))?) else {
1936 return Poll::Ready(Ok(None));
1937 };
1938 self.next_sequence = group.sequence.saturating_add(1);
1939 Poll::Ready(Ok(Some(group)))
1940 }
1941
1942 fn poll_read_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Frame>>> {
1943 let lower = self.min_sequence.max(self.next_sequence);
1944 let Some((frame, found_index, sequence)) =
1945 ready!(self.poll(waiter, |state| { state.poll_read_frame(self.index, lower, waiter) })?)
1946 else {
1947 return Poll::Ready(Ok(None));
1948 };
1949
1950 self.index = found_index + 1;
1951 self.next_sequence = sequence.saturating_add(1);
1952 Poll::Ready(Ok(Some(frame)))
1953 }
1954}
1955
1956#[derive(Clone)]
1962pub struct SubscriberControl {
1963 subscription: kio::Producer<Subscription>,
1964}
1965
1966impl SubscriberControl {
1967 pub fn subscription(&self) -> Subscription {
1969 self.subscription.read().clone()
1970 }
1971
1972 pub fn update(&self, subscription: Subscription) -> Result<()> {
1977 let mut state = self.subscription.write().map_err(|_| Error::Closed)?;
1978 *state = subscription;
1979 Ok(())
1980 }
1981}
1982
1983impl Subscriber {
1984 pub fn info(&self) -> &Info {
1989 &self.info
1990 }
1991
1992 pub fn name(&self) -> &str {
1994 &self.name
1995 }
1996
1997 pub fn control(&self) -> SubscriberControl {
1999 SubscriberControl {
2000 subscription: match &self.inner {
2001 SubscriberKind::Plain(plain) => plain.subscription.clone(),
2002 SubscriberKind::Spliced(spliced) => spliced.prefs(),
2003 },
2004 }
2005 }
2006
2007 pub fn poll_recv_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2018 let meter = self.stats.meter();
2019 let res = match &mut self.inner {
2020 SubscriberKind::Plain(plain) => plain.poll_recv_group(waiter),
2021 SubscriberKind::Spliced(spliced) => spliced.poll_recv_group(waiter),
2022 };
2023 res.map(|res| res.map(|group| group.map(|group| group.with_meter(meter))))
2024 }
2025
2026 pub async fn recv_group(&mut self) -> Result<Option<group::Consumer>> {
2032 kio::wait(|waiter| self.poll_recv_group(waiter)).await
2033 }
2034
2035 pub fn poll_recv_datagram(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Datagram>>> {
2046 let meter = self.stats.meter();
2047 let res = match &mut self.inner {
2048 SubscriberKind::Plain(plain) => plain.poll_recv_datagram(waiter),
2049 SubscriberKind::Spliced(spliced) => spliced.poll_recv_datagram(waiter),
2050 };
2051 if let Poll::Ready(Ok(Some(datagram))) = &res {
2054 meter.datagram(datagram.payload.len() as u64);
2055 }
2056 res
2057 }
2058
2059 pub async fn recv_datagram(&mut self) -> Result<Option<Datagram>> {
2066 kio::wait(|waiter| self.poll_recv_datagram(waiter)).await
2067 }
2068
2069 pub fn poll_next_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
2078 let meter = self.stats.meter();
2079 let res = match &mut self.inner {
2080 SubscriberKind::Plain(plain) => plain.poll_next_group(waiter),
2081 SubscriberKind::Spliced(spliced) => spliced.poll_next_group(waiter),
2082 };
2083 res.map(|res| res.map(|group| group.map(|group| group.with_meter(meter))))
2084 }
2085
2086 pub async fn next_group(&mut self) -> Result<Option<group::Consumer>> {
2092 kio::wait(|waiter| self.poll_next_group(waiter)).await
2093 }
2094
2095 pub fn poll_read_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Frame>>> {
2099 let meter = self.stats.meter();
2100 let res = match &mut self.inner {
2101 SubscriberKind::Plain(plain) => plain.poll_read_frame(waiter),
2102 SubscriberKind::Spliced(spliced) => spliced.poll_read_frame(waiter),
2103 };
2104 if let Poll::Ready(Ok(Some(frame))) = &res {
2107 meter.group();
2108 meter.frames(1);
2109 meter.bytes(frame.payload.len() as u64);
2110 }
2111 res
2112 }
2113
2114 pub async fn read_frame(&mut self) -> Result<Option<frame::Frame>> {
2119 kio::wait(|waiter| self.poll_read_frame(waiter)).await
2120 }
2121
2122 pub fn is_clone(&self, other: &Self) -> bool {
2124 match (&self.inner, &other.inner) {
2125 (SubscriberKind::Plain(a), SubscriberKind::Plain(b)) => a.state.same_channel(&b.state),
2126 (SubscriberKind::Spliced(a), SubscriberKind::Spliced(b)) => a.is_clone(b),
2127 _ => false,
2128 }
2129 }
2130
2131 pub fn poll_finished(&mut self, waiter: &kio::Waiter) -> Poll<Result<u64>> {
2133 match &mut self.inner {
2134 SubscriberKind::Plain(plain) => plain.poll(waiter, |state| state.poll_finished()),
2135 SubscriberKind::Spliced(spliced) => spliced.poll_finished(waiter),
2136 }
2137 }
2138
2139 pub async fn finished(&mut self) -> Result<u64> {
2147 kio::wait(|waiter| self.poll_finished(waiter)).await
2148 }
2149
2150 pub fn start_at(&mut self, sequence: u64) {
2157 match &mut self.inner {
2158 SubscriberKind::Plain(plain) => plain.min_sequence = sequence,
2159 SubscriberKind::Spliced(spliced) => spliced.start_at(sequence),
2160 }
2161 }
2162
2163 pub fn end_at(&mut self, sequence: impl Into<Option<u64>>) {
2176 match &mut self.inner {
2177 SubscriberKind::Plain(plain) => plain.end_sequence = sequence.into(),
2178 SubscriberKind::Spliced(spliced) => spliced.end_at(sequence),
2179 }
2180 }
2181
2182 pub fn subscription(&self) -> Subscription {
2184 self.control().subscription()
2185 }
2186
2187 pub fn update(&mut self, subscription: Subscription) -> Result<()> {
2193 match &mut self.inner {
2194 SubscriberKind::Plain(plain) => {
2195 let mut state = plain.subscription.write().map_err(|_| Error::Closed)?;
2196 *state = subscription;
2197 }
2198 SubscriberKind::Spliced(spliced) => spliced.update(subscription),
2199 }
2200 Ok(())
2201 }
2202
2203 pub fn latest(&self) -> Option<u64> {
2205 match &self.inner {
2206 SubscriberKind::Plain(plain) => plain.state.read().max_sequence,
2207 SubscriberKind::Spliced(spliced) => spliced.latest(),
2208 }
2209 }
2210}
2211
2212pub struct Request {
2224 name: Arc<str>,
2225 broadcast: Arc<broadcast::Info>,
2227 state: kio::Producer<TrackState>,
2228
2229 prev_subscription: Option<Subscription>,
2231
2232 _dynamic: Dynamic,
2237
2238 stats: stats::Scope,
2241}
2242
2243impl Request {
2244 pub(crate) fn new(broadcast: Arc<broadcast::Info>, name: impl Into<Arc<str>>) -> Self {
2245 let name = name.into();
2246 let state = kio::Producer::new(TrackState {
2247 broadcast: broadcast.clone(),
2248 ..Default::default()
2249 });
2250 let dynamic = Dynamic::new(name.clone(), state.clone());
2251 Self {
2252 name,
2253 broadcast,
2254 state,
2255 prev_subscription: None,
2256 _dynamic: dynamic,
2257 stats: stats::Scope::default(),
2258 }
2259 }
2260
2261 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
2264 self.stats = scope;
2265 self
2266 }
2267
2268 pub fn name(&self) -> &str {
2270 &self.name
2271 }
2272
2273 pub fn consume(&self) -> Consumer {
2275 Consumer::plain(self.name.clone(), self.state.consume())
2276 }
2277
2278 pub fn dynamic(&self) -> Dynamic {
2282 Dynamic::new(self.name.clone(), self.state.clone())
2283 }
2284
2285 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<()> {
2288 self.state.poll_unused(waiter).map(|_| ())
2289 }
2290
2291 pub fn accept(self, info: impl Into<Option<Info>>) -> Producer {
2298 let mut info = info.into().unwrap_or_default();
2299 info.broadcast = self.broadcast.clone();
2300 if let Ok(mut state) = self.state.write() {
2303 state.info = Some(info);
2304 }
2305 self.stats.open_subscription();
2308 Producer {
2309 name: self.name,
2310 broadcast: self.broadcast,
2311 state: self.state,
2312 prev_subscription: None,
2313 stats: self.stats,
2314 }
2315 }
2316
2317 pub fn reject(self, err: Error) {
2319 if let Ok(mut state) = self.state.write() {
2320 state.abort = Some(err);
2321 }
2322 }
2323
2324 pub fn subscription(&self) -> Option<Subscription> {
2327 let state = self.state.read();
2328 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
2329 drop(state);
2330 snapshot_subscription(&subs, bound)
2331 }
2332
2333 pub async fn subscription_changed(&mut self) -> Option<Subscription> {
2336 kio::wait(|waiter| self.poll_subscription_changed(waiter)).await
2337 }
2338
2339 pub fn poll_subscription_changed(&mut self, waiter: &kio::Waiter) -> Poll<Option<Subscription>> {
2341 let state = self.state.read();
2342 let (subs, bound) = (state.subscriptions.clone(), state.latency_bound());
2343 drop(state);
2344
2345 let prev = &self.prev_subscription;
2346 let mut combined = None;
2347 let mut guard = ready!(subs.poll(waiter, |subs| {
2348 let next = combined_subscription(subs, bound, waiter);
2349 if &next == prev {
2350 Poll::Pending
2351 } else {
2352 combined = next;
2353 Poll::Ready(())
2354 }
2355 }));
2356 guard.retain(|sub| !sub.is_closed());
2358 drop(guard);
2359 self.prev_subscription = combined.clone();
2360 Poll::Ready(combined)
2361 }
2362
2363 pub(super) fn weak(&self) -> TrackWeak {
2364 TrackWeak {
2365 name: self.name.clone(),
2366 state: self.state.weak(),
2367 }
2368 }
2369}
2370
2371#[cfg(test)]
2372use futures::FutureExt;
2373
2374#[cfg(test)]
2375#[allow(missing_docs)] impl Subscriber {
2377 pub fn assert_group(&mut self) -> group::Consumer {
2378 self.recv_group()
2379 .now_or_never()
2380 .expect("group would have blocked")
2381 .expect("would have errored")
2382 .expect("track was closed")
2383 }
2384
2385 pub fn assert_no_group(&mut self) {
2386 assert!(
2387 self.recv_group().now_or_never().is_none(),
2388 "recv_group would not have blocked"
2389 );
2390 }
2391
2392 pub fn assert_not_closed(&mut self) {
2393 assert!(self.finished().now_or_never().is_none(), "should not be closed");
2394 }
2395
2396 pub fn assert_closed(&mut self) {
2397 assert!(self.finished().now_or_never().is_some(), "should be closed");
2398 }
2399
2400 pub fn assert_error(&mut self) {
2402 assert!(
2403 self.finished().now_or_never().expect("should not block").is_err(),
2404 "should be error"
2405 );
2406 }
2407
2408 pub fn assert_is_clone(&self, other: &Self) {
2409 assert!(self.is_clone(other), "should be clone");
2410 }
2411
2412 pub fn assert_not_clone(&self, other: &Self) {
2413 assert!(!self.is_clone(other), "should not be clone");
2414 }
2415}
2416
2417#[cfg(test)]
2418mod test {
2419 use super::*;
2420
2421 fn track_producer(name: impl Into<Arc<str>>, info: impl Into<Option<Info>>) -> Producer {
2424 Producer::new(Arc::new(broadcast::Info::default()), name, info)
2425 }
2426
2427 fn live_groups(state: &TrackState) -> usize {
2429 state.groups.iter().flatten().count()
2430 }
2431
2432 fn first_live_sequence(state: &TrackState) -> u64 {
2434 state.groups.iter().flatten().next().unwrap().0.sequence
2435 }
2436
2437 fn recv_datagram(dg: &mut Subscriber) -> Datagram {
2439 dg.recv_datagram()
2440 .now_or_never()
2441 .expect("datagram would have blocked")
2442 .expect("would have errored")
2443 .expect("track was closed")
2444 }
2445
2446 #[tokio::test]
2447 async fn append_datagram_shares_group_sequence() {
2448 let mut producer = track_producer("test", None);
2449 let ts = Timestamp::from_millis(10).unwrap();
2450
2451 assert_eq!(producer.append_group().unwrap().sequence, 0);
2453 assert_eq!(producer.append_datagram(ts, &b"a"[..]).unwrap(), 1);
2454 assert_eq!(producer.append_group().unwrap().sequence, 2);
2455 assert_eq!(producer.append_datagram(ts, &b"b"[..]).unwrap(), 3);
2456 assert_eq!(producer.latest(), Some(3));
2457 }
2458
2459 #[tokio::test]
2460 async fn append_datagram_roundtrip() {
2461 let mut producer = track_producer("test", None);
2462 let mut dg = producer.subscribe(None);
2463
2464 let ts = Timestamp::from_millis(42).unwrap();
2465 let seq = producer.append_datagram(ts, &b"hello"[..]).unwrap();
2466
2467 let got = recv_datagram(&mut dg);
2468 assert_eq!(got.sequence, seq);
2469 assert_eq!(got.timestamp, ts);
2470 assert_eq!(&got.payload[..], b"hello");
2471 }
2472
2473 #[tokio::test]
2474 async fn write_datagram_preserves_sequence() {
2475 let mut producer = track_producer("test", None);
2476 let mut dg = producer.subscribe(None);
2477
2478 let ts = Timestamp::from_millis(5).unwrap();
2479 producer
2481 .write_datagram(Datagram {
2482 sequence: 100,
2483 timestamp: ts,
2484 payload: bytes::Bytes::from_static(b"x"),
2485 })
2486 .unwrap();
2487
2488 assert_eq!(recv_datagram(&mut dg).sequence, 100);
2489 assert_eq!(producer.append_group().unwrap().sequence, 101);
2491 }
2492
2493 #[tokio::test]
2494 async fn recv_datagram_advances_ordered_group_cursor() {
2495 let mut producer = track_producer("test", None);
2496 let mut subscriber = producer.subscribe(None);
2497 let ts = Timestamp::from_millis(5).unwrap();
2498
2499 producer
2500 .write_datagram(Datagram {
2501 sequence: 5,
2502 timestamp: ts,
2503 payload: bytes::Bytes::from_static(b"x"),
2504 })
2505 .unwrap();
2506 assert_eq!(recv_datagram(&mut subscriber).sequence, 5);
2507
2508 producer.create_group(group::Info { sequence: 3 }).unwrap();
2509 producer.create_group(group::Info { sequence: 6 }).unwrap();
2510
2511 let group = subscriber
2512 .next_group()
2513 .now_or_never()
2514 .expect("group would have blocked")
2515 .expect("would have errored")
2516 .expect("track was closed");
2517 assert_eq!(group.sequence, 6);
2518 }
2519
2520 #[tokio::test]
2521 async fn datagram_normalized_to_track_timescale() {
2522 let info = Info::default().with_timescale(Timescale::MICRO);
2523 let mut producer = track_producer("test", info);
2524 let mut dg = producer.subscribe(None);
2525
2526 producer
2528 .append_datagram(Timestamp::from_millis(2).unwrap(), &b"z"[..])
2529 .unwrap();
2530 let got = recv_datagram(&mut dg);
2531 assert_eq!(got.timestamp.scale(), Timescale::MICRO);
2532 assert_eq!(got.timestamp.value(), 2_000);
2533 }
2534
2535 #[tokio::test]
2536 async fn datagram_rejects_oversized() {
2537 let mut producer = track_producer("test", None);
2538 let big = bytes::Bytes::from(vec![0u8; crate::model::datagram::MAX_DATAGRAM_PAYLOAD + 1]);
2539 let ts = Timestamp::from_millis(0).unwrap();
2540 assert!(matches!(
2541 producer.append_datagram(ts, big.clone()),
2542 Err(Error::FrameTooLarge)
2543 ));
2544 assert!(matches!(
2545 producer.write_datagram(Datagram {
2546 sequence: 0,
2547 timestamp: ts,
2548 payload: big,
2549 }),
2550 Err(Error::FrameTooLarge)
2551 ));
2552 }
2553
2554 #[tokio::test]
2555 async fn datagram_fanout_to_subscribers() {
2556 let mut producer = track_producer("test", None);
2557 let mut a = producer.subscribe(None);
2559 let mut b = producer.subscribe(None);
2560 let ts = Timestamp::from_millis(1).unwrap();
2561
2562 producer.append_datagram(ts, &b"first"[..]).unwrap();
2563 producer.append_datagram(ts, &b"second"[..]).unwrap();
2564
2565 assert_eq!(&recv_datagram(&mut a).payload[..], b"first");
2567 assert_eq!(&recv_datagram(&mut a).payload[..], b"second");
2568 assert_eq!(&recv_datagram(&mut b).payload[..], b"first");
2569 assert_eq!(&recv_datagram(&mut b).payload[..], b"second");
2570 }
2571
2572 #[tokio::test]
2573 async fn datagram_evicts_stale() {
2574 tokio::time::pause();
2575
2576 let mut producer = track_producer("test", None);
2577 let mut dg = producer.subscribe(None);
2578 let ts = Timestamp::from_millis(0).unwrap();
2579
2580 producer.append_datagram(ts, &b"old"[..]).unwrap(); tokio::time::advance(MAX_DATAGRAM_AGE + Duration::from_millis(10)).await;
2584 producer.append_datagram(ts, &b"new"[..]).unwrap(); let got = recv_datagram(&mut dg);
2588 assert_eq!(got.sequence, 1);
2589 assert_eq!(&got.payload[..], b"new");
2590 }
2591
2592 #[tokio::test]
2593 async fn datagram_recv_pends_until_written() {
2594 let mut producer = track_producer("test", None);
2595 let mut dg = producer.subscribe(None);
2596
2597 assert!(
2598 dg.recv_datagram().now_or_never().is_none(),
2599 "should block with no datagrams"
2600 );
2601
2602 producer
2603 .append_datagram(Timestamp::from_millis(0).unwrap(), &b"go"[..])
2604 .unwrap();
2605 assert_eq!(&recv_datagram(&mut dg).payload[..], b"go");
2606 }
2607
2608 #[tokio::test]
2612 async fn datagram_wire_roundtrip_between_tracks() {
2613 use crate::coding::{Decode, Encode};
2614 use crate::lite;
2615
2616 let version = lite::Version::Lite05;
2617
2618 let mut origin = track_producer("test", None);
2620 let mut origin_dg = origin.subscribe(None);
2621 let ts = Timestamp::from_millis(7).unwrap();
2622 let seq = origin.append_datagram(ts, &b"payload"[..]).unwrap();
2623
2624 let d = recv_datagram(&mut origin_dg);
2625 let body = lite::Datagram {
2626 subscribe: 5,
2627 sequence: d.sequence,
2628 timestamp: d.timestamp.value(),
2629 payload: d.payload.clone(),
2630 }
2631 .encode_bytes(version)
2632 .unwrap();
2633
2634 let mut slice = &body[..];
2636 let wire = lite::Datagram::decode(&mut slice, version).unwrap();
2637 let mut downstream = track_producer("test", None);
2638 let mut downstream_dg = downstream.subscribe(None);
2639 downstream
2640 .write_datagram(Datagram {
2641 sequence: wire.sequence,
2642 timestamp: Timestamp::new(wire.timestamp, Timescale::MILLI).unwrap(),
2643 payload: wire.payload,
2644 })
2645 .unwrap();
2646
2647 let got = recv_datagram(&mut downstream_dg);
2648 assert_eq!(got.sequence, seq);
2649 assert_eq!(got.timestamp, ts);
2650 assert_eq!(&got.payload[..], b"payload");
2651 }
2652
2653 #[tokio::test]
2654 async fn evict_expired_groups() {
2655 tokio::time::pause();
2656
2657 let mut producer = track_producer("test", None);
2658
2659 producer.append_group().unwrap(); producer.append_group().unwrap(); producer.append_group().unwrap(); {
2665 let state = producer.state.read();
2666 assert_eq!(live_groups(&state), 3);
2667 assert_eq!(state.offset, 0);
2668 }
2669
2670 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
2672
2673 producer.append_group().unwrap(); {
2679 let state = producer.state.read();
2680 assert_eq!(live_groups(&state), 1);
2681 assert_eq!(first_live_sequence(&state), 3);
2682 assert_eq!(state.offset, 3);
2683 assert!(!state.duplicates.contains(&0));
2684 assert!(!state.duplicates.contains(&1));
2685 assert!(!state.duplicates.contains(&2));
2686 assert!(state.duplicates.contains(&3));
2687 }
2688 }
2689
2690 #[tokio::test]
2691 async fn evict_keeps_max_sequence() {
2692 tokio::time::pause();
2693
2694 let mut producer = track_producer("test", None);
2695 producer.append_group().unwrap(); tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
2699
2700 producer.append_group().unwrap(); {
2704 let state = producer.state.read();
2705 assert_eq!(live_groups(&state), 1);
2706 assert_eq!(first_live_sequence(&state), 1);
2707 assert_eq!(state.offset, 1);
2708 }
2709 }
2710
2711 #[tokio::test]
2712 async fn no_eviction_when_fresh() {
2713 tokio::time::pause();
2714
2715 let mut producer = track_producer("test", None);
2716 producer.append_group().unwrap(); producer.append_group().unwrap(); producer.append_group().unwrap(); {
2721 let state = producer.state.read();
2722 assert_eq!(live_groups(&state), 3);
2723 assert_eq!(state.offset, 0);
2724 }
2725 }
2726
2727 #[tokio::test]
2728 async fn consumer_skips_evicted_groups() {
2729 tokio::time::pause();
2730
2731 let mut producer = track_producer("test", None);
2732 producer.append_group().unwrap(); let mut consumer = producer.subscribe(None);
2735
2736 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
2737 producer.append_group().unwrap(); let group = consumer.assert_group();
2741 assert_eq!(group.sequence, 1);
2742 }
2743
2744 #[tokio::test]
2745 async fn cache_age_controls_eviction() {
2746 tokio::time::pause();
2747
2748 let mut producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(1)));
2750 producer.append_group().unwrap(); tokio::time::advance(Duration::from_secs(2)).await;
2754 producer.append_group().unwrap(); let state = producer.state.read();
2758 assert_eq!(live_groups(&state), 1);
2759 assert_eq!(first_live_sequence(&state), 1);
2760 }
2761
2762 #[test]
2763 fn latency_max_clamped_to_cache() {
2764 let producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(2)));
2765
2766 let mut subscriber = producer.subscribe(Subscription::default().with_latency_max(Duration::from_secs(10)));
2770 assert_eq!(subscriber.subscription().latency_max, Duration::from_secs(10));
2771 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
2772
2773 subscriber
2775 .update(Subscription::default().with_latency_max(Duration::from_millis(500)))
2776 .unwrap();
2777 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_millis(500));
2778
2779 subscriber
2780 .update(Subscription::default().with_latency_max(Duration::ZERO))
2781 .unwrap();
2782 assert_eq!(producer.subscription().unwrap().latency_max, Duration::ZERO);
2783 }
2784
2785 #[test]
2786 fn latency_max_clamped_via_every_update_path() {
2787 let producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(2)));
2788 let over = Subscription::default().with_latency_max(Duration::from_secs(10));
2789
2790 let mut subscriber = producer.subscribe(over.clone());
2793 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
2794
2795 subscriber.control().update(over.clone()).unwrap();
2796 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
2797
2798 subscriber.update(over).unwrap();
2799 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
2800 }
2801
2802 #[test]
2803 fn latency_max_aggregate_clamps_the_max_across_subscribers() {
2804 let producer = track_producer("test", Info::default().with_latency_max(Duration::from_secs(2)));
2805
2806 let _a = producer.subscribe(Subscription::default().with_latency_max(Duration::from_millis(500)));
2809 let _b = producer.subscribe(Subscription::default().with_latency_max(Duration::from_secs(10)));
2810
2811 assert_eq!(producer.subscription().unwrap().latency_max, Duration::from_secs(2));
2812 }
2813
2814 #[test]
2815 fn subscriber_control_updates_while_read_future_is_pending() {
2816 let producer = track_producer("test", None);
2817 let mut subscriber = producer.subscribe(None);
2818 let control = subscriber.control();
2819
2820 let mut recv = Box::pin(subscriber.recv_group());
2821 assert!(recv.as_mut().now_or_never().is_none());
2822
2823 control
2824 .update(Subscription::default().with_priority(7).with_ordered(false))
2825 .unwrap();
2826
2827 let aggregate = producer.subscription().expect("expected an active subscription");
2828 assert_eq!(aggregate.priority, 7);
2829 assert!(!aggregate.ordered);
2830 }
2831
2832 #[test]
2833 fn dropped_subscriber_leaves_no_ghost_in_aggregate() {
2834 let mut producer = track_producer("test", None);
2839 let a = producer.subscribe(Subscription::default().with_priority(5));
2840
2841 let waiter = kio::Waiter::noop();
2843 assert!(
2844 matches!(producer.poll_subscription_changed(&waiter), Poll::Ready(Ok(Some(_)))),
2845 "one live subscriber should aggregate to Some",
2846 );
2847
2848 drop(a);
2850
2851 assert!(
2853 matches!(producer.poll_subscription_changed(&waiter), Poll::Ready(Ok(None))),
2854 "a dropped subscriber must not linger in the aggregate",
2855 );
2856
2857 assert!(
2859 producer.subscription().is_none(),
2860 "snapshot must exclude a dropped subscriber",
2861 );
2862 }
2863
2864 #[tokio::test]
2865 async fn out_of_order_max_sequence_at_front() {
2866 tokio::time::pause();
2867
2868 let mut producer = track_producer("test", None);
2869
2870 producer.create_group(group::Info { sequence: 5 }).unwrap();
2872 producer.create_group(group::Info { sequence: 3 }).unwrap();
2873 producer.create_group(group::Info { sequence: 4 }).unwrap();
2874
2875 {
2877 let state = producer.state.read();
2878 assert_eq!(state.max_sequence, Some(5));
2879 }
2880
2881 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
2883
2884 producer.append_group().unwrap(); {
2890 let state = producer.state.read();
2891 assert_eq!(live_groups(&state), 1);
2892 assert_eq!(first_live_sequence(&state), 6);
2893 assert!(!state.duplicates.contains(&3));
2894 assert!(!state.duplicates.contains(&4));
2895 assert!(!state.duplicates.contains(&5));
2896 assert!(state.duplicates.contains(&6));
2897 }
2898 }
2899
2900 #[tokio::test]
2901 async fn max_sequence_at_front_blocks_trim() {
2902 tokio::time::pause();
2903
2904 let mut producer = track_producer("test", None);
2905
2906 producer.create_group(group::Info { sequence: 5 }).unwrap();
2908
2909 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
2910
2911 producer.create_group(group::Info { sequence: 3 }).unwrap();
2913
2914 {
2917 let state = producer.state.read();
2918 assert_eq!(live_groups(&state), 2);
2919 assert_eq!(state.offset, 0);
2920 }
2921
2922 tokio::time::advance(DEFAULT_LATENCY_MAX + Duration::from_secs(1)).await;
2924
2925 producer.create_group(group::Info { sequence: 2 }).unwrap();
2927
2928 {
2933 let state = producer.state.read();
2934 assert_eq!(live_groups(&state), 2);
2935 assert_eq!(state.offset, 0);
2936 assert!(state.duplicates.contains(&5));
2937 assert!(!state.duplicates.contains(&3));
2938 assert!(state.duplicates.contains(&2));
2939 }
2940
2941 let mut consumer = producer.subscribe(None);
2943 let group = consumer.assert_group();
2944 assert_eq!(group.sequence, 5);
2946 }
2947
2948 #[tokio::test]
2949 async fn abort_clears_cached_groups() {
2950 let mut producer = track_producer("test", None);
2951 producer.append_group().unwrap();
2952 producer.append_group().unwrap();
2953
2954 let mut consumer = producer.subscribe(None);
2956 assert_eq!(live_groups(&producer.state.read()), 2);
2957
2958 producer.clone().abort(Error::Cancel).unwrap();
2959
2960 {
2961 let state = producer.state.read();
2962 assert!(state.groups.is_empty(), "cached groups should be dropped on abort");
2963 assert!(state.duplicates.is_empty());
2964 }
2965
2966 let result = consumer.recv_group().now_or_never().expect("should not block");
2968 assert!(matches!(result, Err(Error::Cancel)));
2969 }
2970
2971 #[tokio::test]
2972 async fn drop_unfinished_clears_cached_groups() {
2973 let producer = track_producer("test", None);
2974 let mut writer = producer.clone();
2975 writer.append_group().unwrap();
2976
2977 let mut consumer = producer.subscribe(None);
2979 assert_eq!(live_groups(&producer.state.read()), 1);
2980
2981 drop(writer);
2983 drop(producer);
2984
2985 let result = consumer.recv_group().now_or_never().expect("should not block");
2986 assert!(matches!(result, Err(Error::Dropped)));
2987 }
2988
2989 #[tokio::test]
2990 async fn drop_finished_keeps_cached_groups() {
2991 let mut producer = track_producer("test", None);
2992 producer.append_group().unwrap();
2993 producer.finish().unwrap();
2994
2995 let mut consumer = producer.subscribe(None);
2996 drop(producer);
2997
2998 assert_eq!(consumer.assert_group().sequence, 0);
3000 let done = consumer.recv_group().now_or_never().expect("should not block").unwrap();
3001 assert!(done.is_none(), "consumer should drain then see clean finish");
3002 }
3003
3004 #[test]
3005 fn append_finish_cannot_be_rewritten() {
3006 let mut producer = track_producer("test", None);
3007
3008 assert!(producer.finish().is_ok());
3010 assert!(producer.finish().is_err());
3011 assert!(producer.append_group().is_err());
3012 }
3013
3014 #[test]
3015 fn finish_after_groups() {
3016 let mut producer = track_producer("test", None);
3017
3018 producer.append_group().unwrap();
3019 assert!(producer.finish().is_ok());
3020 assert!(producer.finish().is_err());
3021 assert!(producer.append_group().is_err());
3022 }
3023
3024 #[test]
3025 fn finish_at_rejects_a_boundary_at_or_below_the_live_edge() {
3026 let mut producer = track_producer("test", None);
3027 producer.create_group(group::Info { sequence: 5 }).unwrap();
3028
3029 assert!(producer.finish_at(4).is_err());
3032 assert!(producer.finish_at(5).is_err());
3033 assert!(producer.finish_at(6).is_ok());
3034
3035 {
3036 let state = producer.state.read();
3037 assert_eq!(state.final_sequence, Some(6));
3038 }
3039
3040 assert!(producer.finish_at(6).is_err());
3042 assert!(producer.create_group(group::Info { sequence: 4 }).is_ok());
3043 assert!(producer.create_group(group::Info { sequence: 6 }).is_err());
3044 }
3045
3046 #[test]
3047 fn final_sequence_reports_the_declared_boundary() {
3048 let mut producer = track_producer("test", None);
3049 assert_eq!(producer.final_sequence(), None);
3050
3051 producer.create_group(group::Info { sequence: 5 }).unwrap();
3052 assert_eq!(producer.final_sequence(), None, "a group does not declare a boundary");
3053
3054 producer.finish_at(9).unwrap();
3055 assert_eq!(producer.final_sequence(), Some(9));
3056
3057 assert!(producer.finish().is_err());
3059 }
3060
3061 #[test]
3062 fn final_sequence_reports_the_live_edge_after_finish() {
3063 let mut producer = track_producer("test", None);
3064 producer.create_group(group::Info { sequence: 5 }).unwrap();
3065 producer.finish().unwrap();
3066 assert_eq!(producer.final_sequence(), Some(6));
3067 }
3068
3069 #[tokio::test]
3070 async fn finish_at_declares_a_future_boundary() {
3071 let mut producer = track_producer("test", None);
3072 producer.create_group(group::Info { sequence: 5 }).unwrap();
3073
3074 producer.finish_at(7).unwrap();
3076
3077 let mut consumer = producer.subscribe(None);
3078 assert_eq!(consumer.assert_group().sequence, 5);
3079
3080 let boundary = consumer
3083 .finished()
3084 .now_or_never()
3085 .expect("boundary is known immediately")
3086 .expect("would have errored");
3087 assert_eq!(boundary, 7);
3088 assert!(
3089 consumer.recv_group().now_or_never().is_none(),
3090 "should wait for the outstanding group"
3091 );
3092
3093 producer.create_group(group::Info { sequence: 6 }).unwrap();
3095 assert_eq!(consumer.assert_group().sequence, 6);
3096 let done = consumer
3097 .recv_group()
3098 .now_or_never()
3099 .expect("should not block")
3100 .expect("would have errored");
3101 assert!(done.is_none(), "track completes once the boundary is reached");
3102 }
3103
3104 #[tokio::test]
3105 async fn recv_group_finishes_without_waiting_for_gaps() {
3106 let mut producer = track_producer("test", None);
3107 producer.create_group(group::Info { sequence: 1 }).unwrap();
3108 producer.finish().unwrap();
3109
3110 let mut consumer = producer.subscribe(None);
3111 assert_eq!(consumer.assert_group().sequence, 1);
3112
3113 let done = consumer
3114 .recv_group()
3115 .now_or_never()
3116 .expect("should not block")
3117 .expect("would have errored");
3118 assert!(done.is_none(), "track should finish without waiting for gaps");
3119 }
3120
3121 #[tokio::test]
3122 async fn next_group_skips_late_arrivals() {
3123 let mut producer = track_producer("test", None);
3124 let mut consumer = producer.subscribe(None);
3125
3126 producer.create_group(group::Info { sequence: 5 }).unwrap();
3128 let group = consumer
3129 .next_group()
3130 .now_or_never()
3131 .expect("should not block")
3132 .expect("would have errored")
3133 .expect("track should not be closed");
3134 assert_eq!(group.sequence, 5);
3135
3136 producer.create_group(group::Info { sequence: 3 }).unwrap();
3138 producer.create_group(group::Info { sequence: 4 }).unwrap();
3140 producer.create_group(group::Info { sequence: 7 }).unwrap();
3142
3143 let group = consumer
3144 .next_group()
3145 .now_or_never()
3146 .expect("should not block")
3147 .expect("would have errored")
3148 .expect("track should not be closed");
3149 assert_eq!(group.sequence, 7);
3150
3151 assert!(
3153 consumer.next_group().now_or_never().is_none(),
3154 "should block waiting for a higher sequence"
3155 );
3156 }
3157
3158 #[tokio::test]
3159 async fn next_group_returns_arrivals_in_order() {
3160 let mut producer = track_producer("test", None);
3161 let mut consumer = producer.subscribe(None);
3162
3163 producer.create_group(group::Info { sequence: 3 }).unwrap();
3165 producer.create_group(group::Info { sequence: 5 }).unwrap();
3166
3167 let group = consumer
3168 .next_group()
3169 .now_or_never()
3170 .expect("should not block")
3171 .expect("would have errored")
3172 .expect("track should not be closed");
3173 assert_eq!(group.sequence, 3);
3174
3175 let group = consumer
3176 .next_group()
3177 .now_or_never()
3178 .expect("should not block")
3179 .expect("would have errored")
3180 .expect("track should not be closed");
3181 assert_eq!(group.sequence, 5);
3182 }
3183
3184 #[tokio::test]
3185 async fn next_group_and_recv_group_use_independent_cursors() {
3186 let mut producer = track_producer("test", None);
3187 let mut consumer = producer.subscribe(None);
3188
3189 producer.create_group(group::Info { sequence: 5 }).unwrap();
3191 producer.create_group(group::Info { sequence: 3 }).unwrap();
3192
3193 let group = consumer
3196 .next_group()
3197 .now_or_never()
3198 .expect("should not block")
3199 .expect("would have errored")
3200 .expect("track should not be closed");
3201 assert_eq!(group.sequence, 3);
3202
3203 assert_eq!(consumer.assert_group().sequence, 5);
3206 }
3207
3208 #[tokio::test]
3209 async fn end_at_caps_next_group() {
3210 let mut producer = track_producer("test", None);
3211 let mut consumer = producer.subscribe(None);
3212
3213 for s in 0..6 {
3214 producer.create_group(group::Info { sequence: s }).unwrap();
3215 }
3216
3217 consumer.end_at(2);
3218
3219 assert_eq!(
3221 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3222 0
3223 );
3224 assert_eq!(
3225 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3226 1
3227 );
3228 assert_eq!(
3229 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3230 2
3231 );
3232
3233 assert!(
3235 consumer.next_group().now_or_never().is_none(),
3236 "capped consumer must block instead of returning out-of-range groups"
3237 );
3238 }
3239
3240 #[tokio::test]
3241 async fn end_at_release_drains_cached_groups() {
3242 let mut producer = track_producer("test", None);
3243 let mut consumer = producer.subscribe(None);
3244
3245 for s in 0..6 {
3246 producer.create_group(group::Info { sequence: s }).unwrap();
3247 }
3248
3249 consumer.end_at(1);
3250 assert_eq!(
3251 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3252 0
3253 );
3254 assert_eq!(
3255 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3256 1
3257 );
3258 assert!(consumer.next_group().now_or_never().is_none(), "capped at 1");
3259
3260 consumer.end_at(4);
3262 assert_eq!(
3263 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3264 2
3265 );
3266 assert_eq!(
3267 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3268 3
3269 );
3270 assert_eq!(
3271 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3272 4
3273 );
3274 assert!(consumer.next_group().now_or_never().is_none(), "capped at 4");
3275
3276 consumer.end_at(None);
3278 assert_eq!(
3279 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3280 5
3281 );
3282 assert!(consumer.next_group().now_or_never().is_none(), "no more groups");
3283 }
3284
3285 #[tokio::test]
3286 async fn end_at_lower_than_cursor_parks_consumer() {
3287 let mut producer = track_producer("test", None);
3288 let mut consumer = producer.subscribe(None);
3289
3290 for s in 0..3 {
3291 producer.create_group(group::Info { sequence: s }).unwrap();
3292 }
3293
3294 assert_eq!(
3296 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3297 0
3298 );
3299 assert_eq!(
3300 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3301 1
3302 );
3303 assert_eq!(
3304 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3305 2
3306 );
3307
3308 consumer.end_at(1);
3310 producer.create_group(group::Info { sequence: 3 }).unwrap();
3311 producer.create_group(group::Info { sequence: 4 }).unwrap();
3312 assert!(
3313 consumer.next_group().now_or_never().is_none(),
3314 "cap is below cursor; nothing returnable until cap rises"
3315 );
3316
3317 consumer.end_at(None);
3319 assert_eq!(
3320 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3321 3
3322 );
3323 assert_eq!(
3324 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3325 4
3326 );
3327 }
3328
3329 #[tokio::test]
3330 async fn end_at_toggling_around_late_arrivals() {
3331 let mut producer = track_producer("test", None);
3332 let mut consumer = producer.subscribe(None);
3333
3334 consumer.end_at(5);
3335
3336 producer.create_group(group::Info { sequence: 2 }).unwrap();
3338 producer.create_group(group::Info { sequence: 5 }).unwrap();
3339 producer.create_group(group::Info { sequence: 3 }).unwrap();
3340 producer.create_group(group::Info { sequence: 8 }).unwrap();
3342 producer.create_group(group::Info { sequence: 4 }).unwrap();
3343
3344 assert_eq!(
3346 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3347 2
3348 );
3349 assert_eq!(
3350 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3351 3
3352 );
3353 assert_eq!(
3354 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3355 4
3356 );
3357 assert_eq!(
3358 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3359 5
3360 );
3361 assert!(consumer.next_group().now_or_never().is_none());
3363
3364 consumer.end_at(10);
3366 assert_eq!(
3367 consumer.next_group().now_or_never().unwrap().unwrap().unwrap().sequence,
3368 8
3369 );
3370 }
3371
3372 #[tokio::test]
3373 async fn read_frame_returns_single_frame_per_group() {
3374 let mut producer = track_producer("test", None);
3375 let mut consumer = producer.subscribe(None);
3376
3377 producer.write_frame(Timestamp::ZERO, b"hello".as_slice()).unwrap();
3378 producer.write_frame(Timestamp::ZERO, b"world".as_slice()).unwrap();
3379
3380 let frame = consumer
3381 .read_frame()
3382 .now_or_never()
3383 .expect("should not block")
3384 .expect("would have errored")
3385 .expect("track should not be closed");
3386 assert_eq!(&frame.payload[..], b"hello");
3387
3388 let frame = consumer
3389 .read_frame()
3390 .now_or_never()
3391 .expect("should not block")
3392 .expect("would have errored")
3393 .expect("track should not be closed");
3394 assert_eq!(&frame.payload[..], b"world");
3395 }
3396
3397 #[tokio::test]
3398 async fn read_frame_preserves_timestamp() {
3399 let mut producer = track_producer("test", None);
3400 let mut consumer = producer.subscribe(None);
3401
3402 producer
3403 .write_frame(Timestamp::from_micros(20_000).unwrap(), b"hello".as_slice())
3404 .unwrap();
3405
3406 let frame = consumer
3407 .read_frame()
3408 .now_or_never()
3409 .expect("should not block")
3410 .expect("would have errored")
3411 .expect("track should not be closed");
3412 assert_eq!(frame.timestamp.as_micros(), 20_000);
3413 assert_eq!(&frame.payload[..], b"hello");
3414 }
3415
3416 #[tokio::test]
3417 async fn read_frame_skips_stalled_group_for_newer_ready_frame() {
3418 let mut producer = track_producer("test", None);
3419 let mut consumer = producer.subscribe(None);
3420
3421 let _stalled = producer.create_group(group::Info { sequence: 3 }).unwrap();
3423 let mut g5 = producer.create_group(group::Info { sequence: 5 }).unwrap();
3425 g5.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"later"))
3426 .unwrap();
3427 g5.finish().unwrap();
3428
3429 let frame = consumer
3431 .read_frame()
3432 .now_or_never()
3433 .expect("should not block on stalled earlier group")
3434 .expect("would have errored")
3435 .expect("track should not be closed");
3436 assert_eq!(&frame.payload[..], b"later");
3437 }
3438
3439 #[tokio::test]
3440 async fn read_frame_discards_rest_of_multi_frame_group() {
3441 let mut producer = track_producer("test", None);
3442 let mut consumer = producer.subscribe(None);
3443
3444 let mut g0 = producer.create_group(group::Info { sequence: 0 }).unwrap();
3446 g0.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"one"))
3447 .unwrap();
3448 g0.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"two"))
3449 .unwrap();
3450 g0.finish().unwrap();
3451
3452 producer.write_frame(Timestamp::ZERO, b"next".as_slice()).unwrap();
3454
3455 let frame = consumer
3456 .read_frame()
3457 .now_or_never()
3458 .expect("should not block")
3459 .expect("would have errored")
3460 .expect("track should not be closed");
3461 assert_eq!(&frame.payload[..], b"one");
3462
3463 let frame = consumer
3465 .read_frame()
3466 .now_or_never()
3467 .expect("should not block")
3468 .expect("would have errored")
3469 .expect("track should not be closed");
3470 assert_eq!(&frame.payload[..], b"next");
3471 }
3472
3473 #[tokio::test]
3474 async fn read_frame_waits_for_pending_group_after_finish() {
3475 let mut producer = track_producer("test", None);
3478 let mut consumer = producer.subscribe(None);
3479
3480 let mut g0 = producer.create_group(group::Info { sequence: 0 }).unwrap();
3481 producer.finish().unwrap();
3482
3483 assert!(
3485 consumer.read_frame().now_or_never().is_none(),
3486 "read_frame must block on a pending group even after finish()"
3487 );
3488
3489 g0.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"late"))
3491 .unwrap();
3492 let frame = consumer
3493 .read_frame()
3494 .now_or_never()
3495 .expect("should not block once a frame is written")
3496 .expect("would have errored")
3497 .expect("track should not be closed");
3498 assert_eq!(&frame.payload[..], b"late");
3499 }
3500
3501 #[tokio::test]
3502 async fn read_frame_respects_start_at() {
3503 let mut producer = track_producer("test", None);
3506 let mut consumer = producer.subscribe(None);
3507 consumer.start_at(5);
3508
3509 let mut g3 = producer.create_group(group::Info { sequence: 3 }).unwrap();
3511 g3.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"skip-me"))
3512 .unwrap();
3513 g3.finish().unwrap();
3514
3515 let mut g5 = producer.create_group(group::Info { sequence: 5 }).unwrap();
3516 g5.write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"keep"))
3517 .unwrap();
3518 g5.finish().unwrap();
3519
3520 let frame = consumer
3521 .read_frame()
3522 .now_or_never()
3523 .expect("should not block")
3524 .expect("would have errored")
3525 .expect("track should not be closed");
3526 assert_eq!(&frame.payload[..], b"keep");
3527 }
3528
3529 #[tokio::test]
3530 async fn read_frame_returns_none_when_finished() {
3531 let mut producer = track_producer("test", None);
3532 let mut consumer = producer.subscribe(None);
3533
3534 producer.write_frame(Timestamp::ZERO, b"only".as_slice()).unwrap();
3535 producer.finish().unwrap();
3536
3537 let frame = consumer
3538 .read_frame()
3539 .now_or_never()
3540 .expect("should not block")
3541 .expect("would have errored")
3542 .expect("track should not be closed");
3543 assert_eq!(&frame.payload[..], b"only");
3544
3545 let done = consumer
3546 .read_frame()
3547 .now_or_never()
3548 .expect("should not block")
3549 .expect("would have errored");
3550 assert!(done.is_none());
3551 }
3552
3553 #[test]
3554 fn append_group_returns_bounds_exceeded_on_sequence_overflow() {
3555 let mut producer = track_producer("test", None);
3556 {
3557 let mut state = producer.state.write().ok().unwrap();
3558 state.max_sequence = Some(u64::MAX);
3559 }
3560
3561 assert!(matches!(producer.append_group(), Err(Error::BoundsExceeded(_))));
3562 }
3563
3564 #[tokio::test]
3565 async fn fetch_cache_hit() {
3566 let mut producer = track_producer("test", None);
3567
3568 let mut group = producer.append_group().unwrap(); group
3571 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hello"))
3572 .unwrap();
3573 group.finish().unwrap();
3574
3575 let dynamic = producer.dynamic();
3578 let consumer = producer.consume();
3579 assert!(consumer.peek_group(0).is_some());
3580 let mut g = consumer.fetch_group(0, None).await.unwrap();
3581 assert_eq!(g.sequence, 0);
3582 assert_eq!(&g.read_frame().await.unwrap().unwrap().payload[..], b"hello");
3583
3584 assert!(dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending());
3586 }
3587
3588 #[tokio::test]
3589 async fn fetch_miss_signals_dynamic() {
3590 let producer = track_producer("test", None);
3591 let dynamic = producer.dynamic();
3592 let consumer = producer.consume();
3593
3594 assert!(consumer.peek_group(5).is_none());
3598 let pending = consumer.fetch_group(5, group::Fetch::default().with_priority(7));
3599 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
3600
3601 let req = dynamic
3602 .requested_group()
3603 .now_or_never()
3604 .expect("should not block")
3605 .unwrap();
3606 assert_eq!(req.sequence(), 5);
3607 assert_eq!(req.priority(), 7);
3608
3609 let mut group = req.accept(None).unwrap();
3611 group
3612 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hi"))
3613 .unwrap();
3614 group.finish().unwrap();
3615
3616 let mut g = pending.await.unwrap();
3617 assert_eq!(g.sequence, 5);
3618 assert_eq!(&g.read_frame().await.unwrap().unwrap().payload[..], b"hi");
3619 }
3620
3621 #[tokio::test]
3622 async fn fetch_miss_rejects() {
3623 let producer = track_producer("test", None);
3624 let dynamic = producer.dynamic();
3625 let consumer = producer.consume();
3626
3627 let pending = consumer.fetch_group(5, None);
3628 let req = dynamic
3629 .requested_group()
3630 .now_or_never()
3631 .expect("should not block")
3632 .unwrap();
3633
3634 req.reject(Error::Cancel);
3635 assert!(matches!(pending.await, Err(Error::Cancel)));
3636 let fetch = producer.state.read().fetch.clone();
3637 assert!(fetch.read().is_empty());
3638 }
3639
3640 #[tokio::test]
3641 async fn fetch_miss_drop_rejects() {
3642 let producer = track_producer("test", None);
3643 let dynamic = producer.dynamic();
3644 let consumer = producer.consume();
3645
3646 let pending = consumer.fetch_group(5, None);
3647 let req = dynamic
3648 .requested_group()
3649 .now_or_never()
3650 .expect("should not block")
3651 .unwrap();
3652
3653 drop(req);
3654 assert!(matches!(pending.await, Err(Error::Dropped)));
3655 }
3656
3657 #[tokio::test]
3658 async fn fetch_reject_does_not_poison_retry() {
3659 let producer = track_producer("test", None);
3660 let dynamic = producer.dynamic();
3661 let consumer = producer.consume();
3662
3663 let pending = consumer.fetch_group(5, None);
3664 let req = dynamic
3665 .requested_group()
3666 .now_or_never()
3667 .expect("should not block")
3668 .unwrap();
3669 req.reject(Error::Cancel);
3670 assert!(matches!(pending.await, Err(Error::Cancel)));
3671
3672 let retry = consumer.fetch_group(5, None);
3673 let req = dynamic
3674 .requested_group()
3675 .now_or_never()
3676 .expect("should not block")
3677 .unwrap();
3678 let mut group = req.accept(None).unwrap();
3679 group
3680 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"retry"))
3681 .unwrap();
3682 group.finish().unwrap();
3683
3684 let mut group = retry.await.unwrap();
3685 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"retry");
3686 }
3687
3688 #[tokio::test]
3689 async fn fetch_coalesces_concurrent() {
3690 let producer = track_producer("test", None);
3691 let dynamic = producer.dynamic();
3692 let consumer = producer.consume();
3693
3694 let first = consumer.fetch_group(5, group::Fetch::default().with_priority(1));
3697 let second = consumer.fetch_group(5, group::Fetch::default().with_priority(7));
3698 assert!(kio::Pollable::poll(&*first, &kio::Waiter::noop()).is_pending());
3699
3700 let req = dynamic
3701 .requested_group()
3702 .now_or_never()
3703 .expect("should not block")
3704 .unwrap();
3705 assert_eq!(req.sequence(), 5);
3706 assert_eq!(req.priority(), 7);
3707 assert!(
3708 dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending(),
3709 "the second fetch queued a duplicate request"
3710 );
3711
3712 let third = consumer.fetch_group(5, None);
3714
3715 let mut group = req.accept(None).unwrap();
3717 group
3718 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"hi"))
3719 .unwrap();
3720 group.finish().unwrap();
3721
3722 assert_eq!(first.await.unwrap().sequence, 5);
3723 assert_eq!(second.await.unwrap().sequence, 5);
3724 assert_eq!(third.await.unwrap().sequence, 5);
3725 }
3726
3727 #[tokio::test]
3728 async fn fetch_coalesced_reject_fails_all() {
3729 let producer = track_producer("test", None);
3730 let dynamic = producer.dynamic();
3731 let consumer = producer.consume();
3732
3733 let first = consumer.fetch_group(5, None);
3734 let second = consumer.fetch_group(5, None);
3735 let req = dynamic
3736 .requested_group()
3737 .now_or_never()
3738 .expect("should not block")
3739 .unwrap();
3740 req.reject(Error::Cancel);
3741
3742 assert!(matches!(first.await, Err(Error::Cancel)));
3743 assert!(matches!(second.await, Err(Error::Cancel)));
3744
3745 let retry = consumer.fetch_group(5, None);
3747 assert!(kio::Pollable::poll(&*retry, &kio::Waiter::noop()).is_pending());
3748 let req = dynamic
3749 .requested_group()
3750 .now_or_never()
3751 .expect("should not block")
3752 .unwrap();
3753 assert_eq!(req.sequence(), 5);
3754 }
3755
3756 #[tokio::test]
3757 async fn fetch_queued_fails_when_handlers_leave() {
3758 let producer = track_producer("test", None);
3759 let dynamic = producer.dynamic();
3760 let consumer = producer.consume();
3761
3762 let pending = consumer.fetch_group(5, None);
3764 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
3765 drop(dynamic);
3766 assert!(matches!(pending.await, Err(Error::NotFound)));
3767
3768 let fetch = producer.state.read().fetch.clone();
3770 assert!(fetch.read().is_empty());
3771 }
3772
3773 #[tokio::test]
3774 async fn fetch_miss_no_dynamic_not_found() {
3775 let mut producer = track_producer("test", None);
3778 producer.append_group().unwrap(); let consumer = producer.consume();
3780 assert!(matches!(consumer.fetch_group(5, None).await, Err(Error::NotFound)));
3781 }
3782
3783 #[tokio::test]
3784 async fn fetch_past_final_not_found() {
3785 let mut producer = track_producer("test", None);
3786 producer.append_group().unwrap(); producer.finish().unwrap(); let dynamic = producer.dynamic();
3792 let consumer = producer.consume();
3793 assert!(matches!(consumer.fetch_group(5, None).await, Err(Error::NotFound)));
3794
3795 assert!(dynamic.poll_requested_group(&kio::Waiter::noop()).is_pending());
3797 }
3798
3799 fn pooled_producer(capacity: u64) -> (Producer, cache::Pool) {
3801 let pool = cache::Pool::new(capacity);
3802 let broadcast = broadcast::Info {
3803 origin: crate::origin::Info::default().with_pool(pool.clone()),
3804 ..Default::default()
3805 };
3806 let producer = Producer::new(Arc::new(broadcast), "test", None);
3807 (producer, pool)
3808 }
3809
3810 fn finished_group(producer: &mut Producer, size: usize) -> u64 {
3811 let mut group = producer.append_group().unwrap();
3812 group
3813 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; size]))
3814 .unwrap();
3815 group.finish().unwrap();
3816 group.sequence
3817 }
3818
3819 #[tokio::test]
3820 async fn pool_evicts_oldest_group() {
3821 tokio::time::pause();
3822
3823 let (mut producer, pool) = pooled_producer(3000);
3825
3826 finished_group(&mut producer, 1000); tokio::time::advance(Duration::from_millis(10)).await;
3828 finished_group(&mut producer, 1000); tokio::time::advance(Duration::from_millis(10)).await;
3830 finished_group(&mut producer, 1000); assert!(pool.used() <= 3000, "pool should be back under budget");
3835
3836 let consumer = producer.consume();
3837 assert!(consumer.peek_group(0).is_none(), "evicted group is a cache miss");
3838 assert!(consumer.peek_group(1).is_some());
3839 assert!(consumer.peek_group(2).is_some());
3840
3841 let mut subscriber = producer.subscribe(None);
3843 assert_eq!(subscriber.assert_group().sequence, 1);
3844 assert_eq!(subscriber.assert_group().sequence, 2);
3845 }
3846
3847 #[tokio::test]
3848 async fn pool_never_evicts_latest() {
3849 tokio::time::pause();
3850
3851 let (mut producer, pool) = pooled_producer(100);
3853 finished_group(&mut producer, 1000);
3854
3855 assert!(pool.used() > 100, "pinned latest may exceed the budget");
3856 let mut subscriber = producer.subscribe(None);
3857 let mut group = subscriber.assert_group();
3858 assert_eq!(group.read_frame().await.unwrap().unwrap().payload.len(), 1000);
3859 }
3860
3861 #[tokio::test]
3862 async fn pool_reads_bump_recency() {
3863 tokio::time::pause();
3864
3865 let (mut producer, pool) = pooled_producer(3000);
3866 let mut subscriber = producer.subscribe(None);
3867
3868 finished_group(&mut producer, 1000); tokio::time::advance(Duration::from_millis(10)).await;
3870 finished_group(&mut producer, 1000); tokio::time::advance(Duration::from_millis(10)).await;
3872
3873 let mut group = subscriber.assert_group();
3875 assert_eq!(group.sequence, 0);
3876 group.read_frame().await.unwrap().unwrap();
3877 tokio::time::advance(Duration::from_millis(10)).await;
3878
3879 finished_group(&mut producer, 1000); let consumer = producer.consume();
3883 assert!(
3884 consumer.peek_group(0).is_some(),
3885 "recently read group survives: {:?}",
3886 pool.debug_entries()
3887 );
3888 assert!(
3889 consumer.peek_group(1).is_none(),
3890 "stale group is evicted: {:?}",
3891 pool.debug_entries()
3892 );
3893 }
3894
3895 #[tokio::test]
3896 async fn pool_eviction_aborts_readers() {
3897 tokio::time::pause();
3898
3899 let (mut producer, pool) = pooled_producer(3000);
3900 let mut subscriber = producer.subscribe(None);
3901
3902 finished_group(&mut producer, 1000); let group0 = subscriber.assert_group();
3904
3905 tokio::time::advance(Duration::from_millis(10)).await;
3906 finished_group(&mut producer, 1000); tokio::time::advance(Duration::from_millis(10)).await;
3908 finished_group(&mut producer, 1000); let mut group0 = group0;
3913 let read = group0.read_frame().await;
3914 assert!(
3915 matches!(read, Err(Error::Evicted)),
3916 "expected Evicted, got {read:?}: {:?}",
3917 pool.debug_entries()
3918 );
3919 }
3920
3921 #[tokio::test]
3922 async fn pool_growth_on_old_group_charges() {
3923 tokio::time::pause();
3924
3925 let (mut producer, pool) = pooled_producer(3000);
3926
3927 let mut group0 = producer.append_group().unwrap();
3929 tokio::time::advance(Duration::from_millis(10)).await;
3930 let _group1 = producer.append_group().unwrap();
3932 tokio::time::advance(Duration::from_millis(10)).await;
3933
3934 group0
3937 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 4000]))
3938 .unwrap();
3939 assert!(pool.used() <= 3000, "growth on an old group triggers eviction");
3940 assert!(matches!(group0.abort(Error::Cancel), Err(Error::Evicted)));
3941 }
3942
3943 #[tokio::test]
3944 async fn refetched_latest_group_is_repinned() {
3945 tokio::time::pause();
3946
3947 let (mut producer, pool) = pooled_producer(3000);
3948 let dynamic = producer.dynamic();
3949
3950 let mut straggler = producer.append_group().unwrap();
3952 straggler
3953 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 1000]))
3954 .unwrap();
3955 tokio::time::advance(Duration::from_millis(10)).await;
3956
3957 let latest = producer.append_group().unwrap(); latest.abort(Error::Cancel).unwrap();
3960 tokio::time::advance(Duration::from_millis(10)).await;
3961
3962 let consumer = producer.consume();
3965 let pending = consumer.fetch_group(1, None);
3966 let req = dynamic
3967 .requested_group()
3968 .now_or_never()
3969 .expect("should not block")
3970 .unwrap();
3971 let mut group = req.accept(None).unwrap();
3972 group
3973 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 1000]))
3974 .unwrap();
3975 group.finish().unwrap();
3976 pending.await.unwrap();
3977 tokio::time::advance(Duration::from_millis(10)).await;
3978
3979 straggler
3982 .write_frame(Timestamp::ZERO, bytes::Bytes::from(vec![0u8; 4000]))
3983 .unwrap();
3984
3985 assert!(pool.used() <= 3000);
3986 let mut group = consumer.peek_group(1).expect("refetched latest must stay pinned");
3987 assert_eq!(group.read_frame().await.unwrap().unwrap().payload.len(), 1000);
3988 }
3989
3990 #[tokio::test]
3991 async fn pool_eviction_allows_refetch() {
3992 tokio::time::pause();
3993
3994 let (mut producer, _pool) = pooled_producer(3000);
3995 let dynamic = producer.dynamic();
3996
3997 finished_group(&mut producer, 1000); tokio::time::advance(Duration::from_millis(10)).await;
3999 finished_group(&mut producer, 1000); tokio::time::advance(Duration::from_millis(10)).await;
4001 finished_group(&mut producer, 1000); let consumer = producer.consume();
4006 assert!(consumer.peek_group(0).is_none());
4007 let pending = consumer.fetch_group(0, None);
4008
4009 let req = dynamic
4010 .requested_group()
4011 .now_or_never()
4012 .expect("should not block")
4013 .unwrap();
4014 assert_eq!(req.sequence(), 0);
4015
4016 let mut group = req.accept(None).unwrap();
4018 group
4019 .write_frame(Timestamp::ZERO, bytes::Bytes::from_static(b"refetched"))
4020 .unwrap();
4021 group.finish().unwrap();
4022
4023 let mut group = pending.await.unwrap();
4024 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"refetched");
4025 }
4026
4027 #[tokio::test]
4028 async fn fetch_aborts_with_track() {
4029 let producer = track_producer("test", None);
4030 let dynamic = producer.dynamic();
4031 let consumer = producer.consume();
4032
4033 let pending = consumer.fetch_group(3, None);
4034 assert!(kio::Pollable::poll(&*pending, &kio::Waiter::noop()).is_pending());
4035
4036 producer.abort(Error::Cancel).unwrap();
4037 assert!(pending.await.is_err());
4038 drop(dynamic);
4039 }
4040}