1use crate::cache;
18use crate::frame::{self, Frame, FrameBuf};
19use crate::{Cap, Timescale, stats, track};
20use std::collections::VecDeque;
21use std::mem::MaybeUninit;
22use std::ops::{Bound, RangeBounds};
23use std::sync::Arc;
24use std::sync::atomic::{AtomicBool, Ordering};
25use std::task::{Poll, ready};
26
27use crate::{Error, IntoBytes, Result, Timestamp};
28
29pub const MAX_CACHE_BYTES: u64 = 32 * 1024 * 1024; pub const MAX_GROUP_FRAMES: usize = 8192;
41
42const FRAME_SLOTS: usize = 4;
47
48pub(crate) const CACHE_OVERHEAD: u64 = (kio::Producer::<GroupState>::HEAP
56 + 2 * size_of::<usize>()
59 + size_of::<Alive>()
60 + FRAME_SLOTS * size_of::<Frame>()) as u64;
61
62#[derive(Clone, Copy, Debug, Hash, Eq, PartialEq, Ord, PartialOrd)]
66pub struct Info {
67 pub sequence: u64,
70}
71
72impl Info {
73 #[cfg(test)]
79 pub(crate) fn produce(self) -> Producer {
80 Producer::new(self, track::Info::default(), Default::default())
81 }
82}
83
84impl From<usize> for Info {
85 fn from(sequence: usize) -> Self {
86 Self {
87 sequence: sequence as u64,
88 }
89 }
90}
91
92impl From<u64> for Info {
93 fn from(sequence: u64) -> Self {
94 Self { sequence }
95 }
96}
97
98impl From<u32> for Info {
99 fn from(sequence: u32) -> Self {
100 Self {
101 sequence: sequence as u64,
102 }
103 }
104}
105
106impl From<u16> for Info {
107 fn from(sequence: u16) -> Self {
108 Self {
109 sequence: sequence as u64,
110 }
111 }
112}
113
114pub(crate) struct Partial {
117 timestamp: Timestamp,
118 buf: FrameBuf,
119}
120
121#[derive(Default)]
124pub(crate) struct GroupState {
125 pub(crate) frames: VecDeque<Frame>,
128
129 pub(crate) partial: Option<Partial>,
131
132 pub(crate) offset: usize,
136
137 next_index: usize,
141
142 committed: usize,
146
147 pub(crate) cache: u64,
149
150 charge: cache::Charge,
153
154 timestamp: Option<Timestamp>,
160
161 latest: Option<Timestamp>,
165
166 pub(crate) fin: Option<usize>,
169
170 pub(crate) abort: Option<Error>,
174}
175
176impl GroupState {
177 fn content(&self) -> stats::Content {
179 stats::Content {
180 bytes: self.cache,
181 frames: self.next_index.saturating_sub(self.offset) as u64,
182 groups: 1,
183 datagrams: 0,
184 }
185 }
186
187 pub(crate) fn content_range(&self, start: usize, end: usize) -> stats::Content {
189 let start = start.max(self.offset);
190 let end = end.min(self.next_index);
191 if start >= end {
192 return stats::Content::default();
193 }
194
195 let local_start = start.saturating_sub(self.offset).min(self.frames.len());
196 let local_end = end.saturating_sub(self.offset).min(self.frames.len());
197 let mut bytes = self
198 .frames
199 .range(local_start..local_end)
200 .map(|frame| frame.payload.len() as u64)
201 .sum();
202 if start <= self.committed
203 && self.committed < end
204 && let Some(partial) = &self.partial
205 {
206 bytes += partial.buf.capacity() as u64;
207 }
208
209 stats::Content {
210 bytes,
211 frames: (end - start) as u64,
212 groups: 0,
213 datagrams: 0,
214 }
215 }
216
217 fn poll_frame_source(&self, index: usize) -> Poll<Result<Option<(frame::Info, frame::Source)>>> {
220 if index < self.offset {
221 return Poll::Ready(Err(Error::Lagged));
222 }
223 let local = index - self.offset;
224 if let Some(f) = self.frames.get(local) {
225 self.charge.refresh();
228 let info = frame::Info {
229 size: f.payload.len() as u64,
230 timestamp: f.timestamp,
231 };
232 return Poll::Ready(Ok(Some((info, frame::Source::Complete(f.payload.clone())))));
233 }
234 if local == self.frames.len()
235 && let Some(p) = &self.partial
236 {
237 self.charge.refresh();
238 let info = frame::Info {
239 size: p.buf.capacity() as u64,
240 timestamp: p.timestamp,
241 };
242 return Poll::Ready(Ok(Some((info, frame::Source::Partial(p.buf.clone())))));
243 }
244 ready!(self.poll_terminal(index))?;
245 Poll::Ready(Ok(None))
246 }
247
248 fn poll_terminal(&self, index: usize) -> Poll<Result<()>> {
255 match (self.fin, &self.abort) {
256 (Some(total), Some(err)) if index < total => Poll::Ready(Err(err.clone())),
257 (Some(_), _) => Poll::Ready(Ok(())),
258 (None, Some(err)) => Poll::Ready(Err(err.clone())),
259 (None, None) => Poll::Pending,
260 }
261 }
262
263 fn poll_end(&self, index: usize) -> Poll<Result<()>> {
266 if index < self.offset {
267 return Poll::Ready(Err(Error::Lagged));
268 }
269 self.poll_terminal(index)
270 }
271
272 fn stamp(&mut self, timestamp: Timestamp) {
275 self.timestamp.get_or_insert(timestamp);
276 self.latest = Some(timestamp);
277 }
278
279 fn would_overflow(&self, extra_frames: usize, extra_bytes: u64) -> bool {
281 self.next_index.saturating_sub(self.offset).saturating_add(extra_frames) > MAX_GROUP_FRAMES
282 || self.cache.saturating_add(extra_bytes) > MAX_CACHE_BYTES
283 }
284
285 fn release(&mut self) {
287 self.frames.clear();
288 self.partial = None;
289 self.cache = 0;
290 self.charge.clear();
291 }
292}
293
294fn modify(state: &kio::Producer<GroupState>) -> Result<kio::Mut<'_, GroupState>> {
295 state.write().map_err(|r| r.abort.clone().unwrap_or(Error::Dropped))
296}
297
298pub struct Producer {
304 state: kio::Producer<GroupState>,
306
307 info: Info,
310
311 track: track::Info,
316
317 cache: Arc<cache::Track>,
321
322 stats: stats::Meter,
325
326 alive: Arc<Alive>,
329}
330
331struct Alive {
339 info: Info,
340 state: kio::Producer<GroupState>,
341 aborted: AtomicBool,
349 access: Arc<cache::Access>,
353}
354
355impl Drop for Alive {
356 fn drop(&mut self) {
357 match self.state.write() {
363 Ok(mut state) => {
364 if state.fin.is_some() || state.abort.is_some() {
365 return;
366 }
367 tracing::warn!(
368 sequence = self.info.sequence,
369 "group::Producer dropped without finish() or abort()"
370 );
371 state.release();
372 }
373 Err(state) => {
374 if state.fin.is_some() || state.abort.is_some() {
375 return;
376 }
377 tracing::warn!(
378 sequence = self.info.sequence,
379 "group::Producer dropped without finish() or abort()"
380 );
381 }
382 }
383 }
384}
385
386impl std::ops::Deref for Producer {
387 type Target = Info;
388
389 fn deref(&self) -> &Self::Target {
390 &self.info
391 }
392}
393
394impl Producer {
395 pub(crate) fn new(info: Info, track: track::Info, cache: Arc<cache::Track>) -> Self {
406 let state = kio::Producer::<GroupState>::default();
407 let charge = cache.charge();
408 let access = charge.access();
409 state.write().ok().expect("a new group is open").charge = charge;
410 let alive = Arc::new(Alive {
411 info,
412 state: state.clone(),
413 aborted: AtomicBool::new(false),
414 access,
415 });
416 Self {
417 info,
418 state,
419 track,
420 cache,
421 stats: stats::Meter::default(),
422 alive,
423 }
424 }
425
426 pub(crate) fn with_meter(mut self, meter: stats::Meter) -> Self {
429 meter.group();
430 self.stats = meter;
431 self
432 }
433
434 pub(crate) fn info(&self) -> Info {
436 self.info
437 }
438
439 pub fn timescale(&self) -> Timescale {
441 self.track.timescale
442 }
443
444 pub fn start_at(&mut self, index: u64) -> Result<()> {
458 let index = usize::try_from(index).map_err(|_| Error::BoundsExceeded(crate::coding::BoundsExceeded))?;
459 if index == usize::MAX {
460 return Err(Error::BoundsExceeded(crate::coding::BoundsExceeded));
461 }
462
463 let mut state = modify(&self.state)?;
464 if state.fin.is_some() || state.next_index != state.offset {
467 return Err(Error::Closed);
468 }
469 state.offset = index;
470 state.next_index = index;
471 state.committed = index;
472 Ok(())
473 }
474
475 pub fn write_frame<B: IntoBytes>(&mut self, timestamp: Timestamp, data: B) -> Result<()> {
483 let timestamp = timestamp
484 .convert(self.track.timescale)
485 .map_err(|_| Error::TimestampMismatch)?;
486 let payload = data.into_bytes();
487 if payload.len() as u64 > MAX_CACHE_BYTES {
488 return Err(Error::FrameTooLarge);
489 }
490
491 let mut state = modify(&self.state)?;
492 if state.fin.is_some() {
493 return Err(Error::Closed);
494 }
495 if state.partial.is_some() {
496 return Err(Error::FrameOpen);
497 }
498 let next_index = state
499 .next_index
500 .checked_add(1)
501 .ok_or(Error::BoundsExceeded(crate::coding::BoundsExceeded))?;
502 debug_assert!(state.partial.is_none(), "a frame is already open");
503 let size = payload.len() as u64;
504 if state.would_overflow(1, size) {
505 return Err(self.abort_too_large(state));
506 }
507 state.cache += size;
508 let now = state.charge.add(size);
509 state.frames.push_back(Frame { timestamp, payload });
510 state.next_index = next_index;
511 state.committed = state.next_index;
512 state.stamp(timestamp);
513 drop(state);
514
515 self.cache.settle(now);
518
519 self.stats.frames(1);
521 self.stats.bytes(size);
522 Ok(())
523 }
524
525 pub fn write_frames<const N: usize>(&mut self, frames: &mut frame::Buffer<N>) -> Result<()> {
536 for frame in frames.filled() {
537 frame
538 .timestamp
539 .convert(self.track.timescale)
540 .map_err(|_| Error::TimestampMismatch)?;
541 if frame.payload.len() as u64 > MAX_CACHE_BYTES {
542 return Err(Error::FrameTooLarge);
543 }
544 }
545
546 let count = frames.len();
547 let bytes: u64 = frames.filled().iter().map(|frame| frame.payload.len() as u64).sum();
548 let mut state = modify(&self.state)?;
549 if state.fin.is_some() {
550 return Err(Error::Closed);
551 }
552 if state.partial.is_some() {
553 return Err(Error::FrameOpen);
554 }
555 let next_index = state
556 .next_index
557 .checked_add(count)
558 .ok_or(Error::BoundsExceeded(crate::coding::BoundsExceeded))?;
559 if state.would_overflow(count, bytes) {
560 return Err(self.abort_too_large(state));
561 }
562
563 let mut now = None;
565 for mut frame in frames.drain() {
566 frame.timestamp = frame
567 .timestamp
568 .convert(self.track.timescale)
569 .expect("timestamp scale checked above");
570 let size = frame.payload.len() as u64;
571 state.cache += size;
572 now = state.charge.add(size);
573 state.stamp(frame.timestamp);
574 state.frames.push_back(frame);
575 }
576 state.next_index = next_index;
577 state.committed = next_index;
578 drop(state);
579
580 self.cache.settle(now);
581 self.stats.frames(count as u64);
582 self.stats.bytes(bytes);
583 Ok(())
584 }
585
586 pub fn create_frame(&mut self, frame: frame::Info) -> Result<frame::Producer<'_>> {
595 let timestamp = frame
596 .timestamp
597 .convert(self.track.timescale)
598 .map_err(|_| Error::TimestampMismatch)?;
599 if frame.size > MAX_CACHE_BYTES {
600 return Err(Error::FrameTooLarge);
601 }
602 let buf = FrameBuf::new(frame.size as usize);
603
604 let mut state = modify(&self.state)?;
605 if state.fin.is_some() {
606 return Err(Error::Closed);
607 }
608 if state.partial.is_some() {
609 return Err(Error::FrameOpen);
610 }
611 let next_index = state
612 .next_index
613 .checked_add(1)
614 .ok_or(Error::BoundsExceeded(crate::coding::BoundsExceeded))?;
615 if state.would_overflow(1, frame.size) {
616 return Err(self.abort_too_large(state));
617 }
618 state.cache += frame.size;
619 let now = state.charge.add(frame.size);
620 state.partial = Some(Partial {
621 timestamp,
622 buf: buf.clone(),
623 });
624 state.next_index = next_index;
625 state.stamp(timestamp);
628 drop(state);
629
630 self.cache.settle(now);
633
634 self.stats.frames(1);
637 let meter = self.stats.clone();
638
639 let info = frame::Info {
640 size: frame.size,
641 timestamp,
642 };
643 Ok(frame::Producer::new(self, buf, info).with_meter(meter))
644 }
645
646 pub(crate) fn create_frame_owned(&mut self, frame: frame::Info) -> Result<frame::ProducerOwned> {
651 let timestamp = frame
652 .timestamp
653 .convert(self.track.timescale)
654 .map_err(|_| Error::TimestampMismatch)?;
655 if frame.size > MAX_CACHE_BYTES {
656 return Err(Error::FrameTooLarge);
657 }
658 let buf = FrameBuf::new(frame.size as usize);
659
660 let mut state = modify(&self.state)?;
661 if state.fin.is_some() {
662 return Err(Error::Closed);
663 }
664 if state.partial.is_some() {
665 return Err(Error::FrameOpen);
666 }
667 let next_index = state
668 .next_index
669 .checked_add(1)
670 .ok_or(Error::BoundsExceeded(crate::coding::BoundsExceeded))?;
671 if state.would_overflow(1, frame.size) {
672 return Err(self.abort_too_large(state));
673 }
674 state.cache += frame.size;
675 let now = state.charge.add(frame.size);
676 state.partial = Some(Partial {
677 timestamp,
678 buf: buf.clone(),
679 });
680 state.next_index = next_index;
681 state.stamp(timestamp);
684 drop(state);
685
686 self.cache.settle(now);
689
690 self.stats.frames(1);
693 let meter = self.stats.clone();
694
695 let info = frame::Info {
696 size: frame.size,
697 timestamp,
698 };
699 Ok(frame::ProducerOwned::new(self.clone(), buf, info).with_meter(meter))
700 }
701
702 pub(crate) fn frame_notify(&self) {
704 let now = self
711 .state
712 .write()
713 .ok()
714 .and_then(|mut state| state.charge.record_write());
715 self.cache.settle(now);
718 }
719
720 pub(crate) fn frame_commit(&mut self, frame: Frame) -> Result<()> {
722 let mut state = modify(&self.state)?;
723 state.partial = None;
726 state.frames.push_back(frame);
727 state.committed = state.next_index;
728 let now = state.charge.record_write();
734 drop(state);
735
736 self.cache.settle(now);
739 Ok(())
740 }
741
742 pub(crate) fn frame_abort(&mut self, err: Error) {
745 let _ = self.clone().abort(err);
746 }
747
748 pub fn frame_count(&self) -> usize {
754 self.state.read().next_index
755 }
756
757 pub fn finish(&self) -> Result<()> {
762 let mut state = modify(&self.state)?;
763 if state.partial.is_some() {
764 return Err(Error::FrameOpen);
765 }
766 state.fin = Some(state.next_index);
767 Ok(())
768 }
769
770 pub fn abort(self, err: Error) -> Result<()> {
776 let mut guard = modify(&self.state)?;
777 guard.abort = Some(err);
778 self.alive.aborted.store(true, Ordering::Release);
779 guard.release();
780 guard.close();
781 Ok(())
782 }
783
784 fn abort_too_large(&self, mut state: kio::Mut<'_, GroupState>) -> Error {
787 let err = Error::GroupTooLarge;
788 state.abort = Some(err.clone());
789 self.alive.aborted.store(true, Ordering::Release);
790 state.release();
791 state.close();
792 err
793 }
794
795 pub(crate) fn is_aborted(&self) -> bool {
803 self.alive.aborted.load(Ordering::Acquire)
804 }
805
806 pub(crate) fn is_finished(&self) -> bool {
808 self.state.read().fin.is_some()
809 }
810
811 pub(crate) fn live_first_frame(&self) -> Option<usize> {
821 let state = self.state.read();
822 state.abort.is_none().then_some(state.offset)
823 }
824
825 pub(crate) fn resume_frame(&self) -> Option<usize> {
837 let state = self.state.read();
838 if state.fin.is_some() {
839 return None;
840 }
841 (state.committed > state.offset).then_some(state.committed)
842 }
843
844 pub(crate) fn timestamp(&self) -> Option<Timestamp> {
854 self.state.read().timestamp
855 }
856
857 pub(crate) fn latest(&self) -> Option<Timestamp> {
866 self.state.read().latest
867 }
868
869 pub(crate) fn cache_size(&self) -> u64 {
872 self.state.read().charge.size()
873 }
874
875 pub(crate) fn cache_accessed(&self) -> u64 {
878 self.alive.access.get()
879 }
880
881 pub(crate) fn cache_accessed_tick(&self, now: Option<u64>) -> Option<u64> {
883 self.alive.access.tick(now)
884 }
885
886 pub(crate) fn cache_demote(&self) {
889 if let Ok(mut state) = self.state.write() {
890 state.charge.demote();
891 }
892 }
893
894 pub(crate) fn cache_refresh(&self) {
900 self.state.read().charge.refresh();
901 }
902
903 pub fn consume(&self) -> Consumer {
905 Consumer {
906 info: self.info,
907 track: self.track.clone(),
908 inner: ConsumerKind::Plain(Plain {
909 state: self.state.consume(),
910 index: 0,
911 end: None,
912 prefetch: Prefetch::default(),
913 cache: self.cache.clone(),
914 access: self.alive.access.clone(),
915 refreshed: self.cache.pool().now(),
916 }),
917 stats: stats::Meter::default(),
920 stale_stats: stats::Meter::default(),
921 expiry: None,
922 expired: false,
923 ended: false,
924 stale_counted: Arc::default(),
925 }
926 }
927
928 pub(crate) fn poll_timestamp(&self, waiter: &kio::Waiter) -> Poll<()> {
930 match self.state.poll(waiter, |state| {
931 if state.timestamp.is_some() || state.fin.is_some() || state.abort.is_some() {
932 Poll::Ready(())
933 } else {
934 Poll::Pending
935 }
936 }) {
937 Poll::Ready(_) => Poll::Ready(()),
938 Poll::Pending => Poll::Pending,
939 }
940 }
941
942 pub async fn closed(&self) -> Error {
944 kio::wait(|waiter| self.poll_closed(waiter)).await
945 }
946
947 pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<Error> {
949 self.state.poll_closed(waiter).map(|()| self.abort_reason())
950 }
951
952 pub async fn used(&self) -> Result<()> {
954 self.state.used().await.map_err(|_| self.abort_reason())
955 }
956
957 pub async fn unused(&self) -> Result<()> {
959 self.state.unused().await.map_err(|_| self.abort_reason())
960 }
961
962 fn abort_reason(&self) -> Error {
964 self.state.read().abort.clone().unwrap_or(Error::Dropped)
965 }
966}
967
968impl Clone for Producer {
969 fn clone(&self) -> Self {
970 Self {
971 info: self.info,
972 state: self.state.clone(),
973 track: self.track.clone(),
974 cache: self.cache.clone(),
975 stats: self.stats.clone(),
976 alive: self.alive.clone(),
977 }
978 }
979}
980
981struct Prefetch {
989 frames: [MaybeUninit<Frame>; Self::CAP],
991 pos: usize,
992 len: usize,
993}
994
995impl Prefetch {
996 const CAP: usize = 8;
997
998 fn pop(&mut self) -> Option<Frame> {
1000 if self.pos == self.len {
1001 return None;
1002 }
1003 let frame = unsafe { self.frames[self.pos].assume_init_read() };
1005 self.pos += 1;
1006 Some(frame)
1007 }
1008
1009 fn fill(&mut self, frames: impl Iterator<Item = Frame>) {
1011 debug_assert_eq!(self.pos, self.len, "fill on a non-empty batch would leak frames");
1012 self.pos = 0;
1013 self.len = 0;
1014 for frame in frames.take(Self::CAP) {
1015 self.frames[self.len].write(frame);
1016 self.len += 1;
1017 }
1018 }
1019
1020 fn buffered(&self) -> (u64, u64) {
1023 let mut bytes = 0u64;
1024 for slot in &self.frames[self.pos..self.len] {
1025 bytes += unsafe { slot.assume_init_ref() }.payload.len() as u64;
1027 }
1028 ((self.len - self.pos) as u64, bytes)
1029 }
1030}
1031
1032impl Default for Prefetch {
1033 fn default() -> Self {
1034 Self {
1035 frames: [const { MaybeUninit::uninit() }; Self::CAP],
1036 pos: 0,
1037 len: 0,
1038 }
1039 }
1040}
1041
1042impl Drop for Prefetch {
1043 fn drop(&mut self) {
1044 for slot in &mut self.frames[self.pos..self.len] {
1045 unsafe { slot.assume_init_drop() };
1047 }
1048 }
1049}
1050
1051pub struct Consumer {
1057 inner: ConsumerKind,
1058
1059 info: Info,
1061
1062 track: track::Info,
1065
1066 stats: stats::Meter,
1069 stale_stats: stats::Meter,
1072
1073 expiry: Option<Arc<dyn Expiry>>,
1076 expired: bool,
1077 ended: bool,
1082 stale_counted: Arc<AtomicBool>,
1085}
1086
1087pub(crate) trait Expiry: Send + Sync {
1089 fn is_expired(&self, waiter: &kio::Waiter) -> bool;
1092}
1093
1094#[expect(clippy::large_enum_variant)]
1097enum ConsumerKind {
1098 Plain(Plain),
1099 Spliced(Box<super::resume::Group>),
1101}
1102
1103struct Plain {
1105 state: kio::Consumer<GroupState>,
1107
1108 index: usize,
1111
1112 end: Option<usize>,
1114
1115 prefetch: Prefetch,
1117
1118 cache: Arc<cache::Track>,
1120 access: Arc<cache::Access>,
1121 refreshed: u64,
1122}
1123
1124impl Clone for Plain {
1125 fn clone(&self) -> Self {
1126 Self {
1129 state: self.state.clone(),
1130 index: self.index,
1131 end: self.end,
1132 prefetch: Prefetch::default(),
1133 cache: self.cache.clone(),
1134 access: self.access.clone(),
1135 refreshed: self.refreshed,
1136 }
1137 }
1138}
1139
1140impl Clone for Consumer {
1141 fn clone(&self) -> Self {
1142 Self {
1143 inner: match &self.inner {
1144 ConsumerKind::Plain(plain) => ConsumerKind::Plain(plain.clone()),
1145 ConsumerKind::Spliced(spliced) => ConsumerKind::Spliced(Box::new((**spliced).clone())),
1146 },
1147 info: self.info,
1148 track: self.track.clone(),
1149 stats: self.stats.clone(),
1152 stale_stats: self.stale_stats.clone(),
1153 expiry: self.expiry.clone(),
1154 expired: self.expired,
1155 ended: self.ended,
1156 stale_counted: self.stale_counted.clone(),
1157 }
1158 }
1159}
1160
1161impl std::ops::Deref for Consumer {
1162 type Target = Info;
1163
1164 fn deref(&self) -> &Self::Target {
1165 &self.info
1166 }
1167}
1168
1169impl Consumer {
1170 pub(crate) fn content(&self) -> stats::Content {
1172 match &self.inner {
1173 ConsumerKind::Plain(plain) => plain.state.read().content(),
1174 ConsumerKind::Spliced(_) => stats::Content {
1177 groups: 1,
1178 ..Default::default()
1179 },
1180 }
1181 }
1182
1183 fn unread_content(&self) -> stats::Content {
1185 match &self.inner {
1186 ConsumerKind::Plain(plain) => plain.unread_content(),
1187 ConsumerKind::Spliced(_) => stats::Content::default(),
1189 }
1190 }
1191
1192 pub(crate) fn into_spliced(self, mut spliced: super::resume::Group) -> Self {
1195 spliced.set_stale_meter(self.stale_stats.clone());
1196 Self {
1197 inner: ConsumerKind::Spliced(Box::new(spliced)),
1198 info: self.info,
1199 track: self.track,
1200 stats: self.stats,
1201 stale_stats: self.stale_stats,
1202 expiry: None,
1206 expired: false,
1207 ended: false,
1208 stale_counted: self.stale_counted,
1209 }
1210 }
1211
1212 pub(crate) fn with_meter(mut self, meter: stats::Meter) -> Self {
1215 meter.group();
1216 self.stats = meter.clone();
1217 self.set_stale_meter(meter);
1218 self
1219 }
1220
1221 pub(crate) fn set_stale_meter(&mut self, meter: stats::Meter) {
1223 if let ConsumerKind::Spliced(spliced) = &mut self.inner {
1224 spliced.set_stale_meter(meter.clone());
1225 }
1226 self.stale_stats = meter;
1227 }
1228
1229 pub(crate) fn with_expiry(mut self, expiry: Arc<dyn Expiry>) -> Self {
1231 self.expiry = Some(expiry);
1232 self
1233 }
1234
1235 pub(crate) fn poll_expired(&mut self, waiter: &kio::Waiter) -> bool {
1237 self.poll_expired_while_pending(waiter, false)
1238 }
1239
1240 fn poll_expired_if_blocked(&mut self, waiter: &kio::Waiter) -> Option<bool> {
1255 if self.ended {
1256 return Some(false);
1257 }
1258 if !self.poll_expired(waiter) {
1259 return None;
1260 }
1261 let truncates = self.expired_truncates();
1262 self.ended = !truncates;
1263 Some(truncates)
1264 }
1265
1266 fn expired_truncates(&self) -> bool {
1275 let unread = self.unread_content();
1276 unread.frames > 0 || unread.bytes > 0
1277 }
1278
1279 pub(crate) fn poll_expired_while_pending(&mut self, waiter: &kio::Waiter, pending: bool) -> bool {
1281 if !self.expired
1282 && (pending || self.expiry_pending())
1283 && self.expiry.as_ref().is_some_and(|expiry| expiry.is_expired(waiter))
1284 {
1285 self.expired = true;
1286 if !self.stale_counted.swap(true, Ordering::Relaxed) {
1287 self.stale_stats.stale(self.unread_content());
1288 }
1289 }
1290 self.expired
1291 }
1292
1293 fn expiry_pending(&self) -> bool {
1295 match &self.inner {
1296 ConsumerKind::Plain(plain) => plain.expiry_pending(),
1297 ConsumerKind::Spliced(_) => false,
1299 }
1300 }
1301
1302 pub(crate) fn latency_expired(&self) -> bool {
1304 self.expired
1305 }
1306
1307 pub(crate) fn is_aborted(&self) -> bool {
1313 match &self.inner {
1314 ConsumerKind::Plain(plain) => plain.state.read().abort.is_some(),
1315 ConsumerKind::Spliced(_) => false,
1316 }
1317 }
1318
1319 pub fn keep_alive(&self) {
1321 if let ConsumerKind::Plain(plain) = &self.inner {
1322 plain.state.read().charge.refresh();
1323 }
1324 }
1325
1326 pub(crate) fn cache_refresh(&self) {
1329 self.keep_alive();
1330 }
1331
1332 pub(crate) fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<()> {
1340 match &self.inner {
1341 ConsumerKind::Plain(plain) => plain.state.poll_closed(waiter),
1342 ConsumerKind::Spliced(_) => Poll::Ready(()),
1343 }
1344 }
1345
1346 pub fn timescale(&self) -> Timescale {
1348 self.track.timescale
1349 }
1350
1351 pub fn index(&self) -> u64 {
1356 match &self.inner {
1357 ConsumerKind::Plain(plain) => plain.index as u64,
1358 ConsumerKind::Spliced(spliced) => spliced.index(),
1359 }
1360 }
1361
1362 pub fn set_frames(&mut self, frames: impl RangeBounds<u64>) {
1368 let (start, end) = super::subscription::sequence_bounds(frames);
1369 self.start_at(start);
1370 self.end_at(end.map_or(Bound::Unbounded, Bound::Excluded));
1371 }
1372
1373 pub(crate) fn start_at(&mut self, index: u64) {
1383 match &mut self.inner {
1384 ConsumerKind::Plain(plain) => plain.start_at(index),
1385 ConsumerKind::Spliced(spliced) => spliced.start_at(index),
1386 }
1387 }
1388
1389 pub fn skip_to(&mut self, index: u64) {
1395 match &mut self.inner {
1396 ConsumerKind::Plain(plain) => plain.skip_to(index),
1397 ConsumerKind::Spliced(spliced) => spliced.start_at(index),
1398 }
1399 }
1400
1401 pub(crate) fn end_at(&mut self, end: impl Into<Cap>) {
1408 let end = end.into().exclusive();
1409 match &mut self.inner {
1410 ConsumerKind::Plain(plain) => {
1411 plain.end = end.map(|end| usize::try_from(end).unwrap_or(usize::MAX));
1412 }
1413 ConsumerKind::Spliced(spliced) => spliced.end_at(end),
1414 }
1415 }
1416
1417 pub fn frame_count(&self) -> usize {
1420 match &self.inner {
1421 ConsumerKind::Plain(plain) => {
1422 let state = plain.state.read();
1423 state.fin.unwrap_or(state.next_index)
1424 }
1425 ConsumerKind::Spliced(spliced) => spliced.frame_count(),
1426 }
1427 }
1428
1429 pub async fn next_frame(&mut self) -> Result<Option<frame::Consumer>> {
1431 kio::wait(|waiter| self.poll_next_frame(waiter)).await
1432 }
1433
1434 pub fn poll_next_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Consumer>>> {
1439 if self.ended {
1440 return Poll::Ready(Ok(None));
1441 }
1442 if self.expired {
1443 return Poll::Ready(Err(Error::Old));
1444 }
1445 let stats = self.stats.clone();
1446 let expiry = self
1447 .expiry
1448 .as_ref()
1449 .map(|policy| frame::Expiry::new(policy.clone(), self.stale_stats.clone(), self.stale_counted.clone()));
1450 let res = match &mut self.inner {
1451 ConsumerKind::Plain(plain) => plain.poll_next_frame(waiter, &stats, expiry),
1452 ConsumerKind::Spliced(spliced) => {
1453 let res = ready!(spliced.poll_next_frame(waiter))?;
1456 if res.is_some() {
1457 stats.frames(1);
1458 }
1459 Poll::Ready(Ok(res.map(|frame| frame.with_meter(stats))))
1460 }
1461 };
1462 match res.is_pending().then(|| self.poll_expired_if_blocked(waiter)).flatten() {
1463 Some(true) => Poll::Ready(Err(Error::Old)),
1464 Some(false) => Poll::Ready(Ok(None)),
1465 None => res,
1466 }
1467 }
1468
1469 pub fn poll_read_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Frame>>> {
1471 if self.ended {
1472 return Poll::Ready(Ok(None));
1473 }
1474 if self.expired {
1475 return Poll::Ready(Err(Error::Old));
1476 }
1477 let stats = self.stats.clone();
1478 let res = match &mut self.inner {
1479 ConsumerKind::Plain(plain) => plain.poll_read_frame(waiter, &stats),
1480 ConsumerKind::Spliced(spliced) => {
1481 let res = ready!(spliced.poll_read_frame(waiter))?;
1482 if let Some(frame) = &res {
1483 stats.frames(1);
1484 stats.bytes(frame.payload.len() as u64);
1485 }
1486 Poll::Ready(Ok(res))
1487 }
1488 };
1489 match res.is_pending().then(|| self.poll_expired_if_blocked(waiter)).flatten() {
1490 Some(true) => Poll::Ready(Err(Error::Old)),
1491 Some(false) => Poll::Ready(Ok(None)),
1492 None => res,
1493 }
1494 }
1495
1496 pub async fn read_frame(&mut self) -> Result<Option<frame::Frame>> {
1498 if !self.expired
1501 && let ConsumerKind::Plain(plain) = &mut self.inner
1502 {
1503 if !plain.capped()
1505 && let Some(frame) = plain.prefetch.pop()
1506 {
1507 plain.refresh_if_stale();
1508 plain.index += 1;
1509 return Ok(Some(frame));
1510 }
1511 }
1512 kio::wait(|waiter| self.poll_read_frame(waiter)).await
1513 }
1514
1515 pub fn poll_read_frames<const N: usize>(
1521 &mut self,
1522 waiter: &kio::Waiter,
1523 out: &mut frame::Buffer<N>,
1524 ) -> Poll<Result<usize>> {
1525 out.clear();
1526 if out.capacity() == 0 {
1527 return Poll::Ready(Ok(0));
1528 }
1529
1530 while !out.is_full() {
1531 match self.poll_read_frame(waiter) {
1532 Poll::Ready(Ok(Some(frame))) => out.push(frame).expect("buffer capacity checked"),
1533 Poll::Ready(Ok(None)) => break,
1534 Poll::Ready(Err(err)) => {
1535 if out.is_empty() {
1536 return Poll::Ready(Err(err));
1537 }
1538 break;
1539 }
1540 Poll::Pending if !out.is_empty() => break,
1541 Poll::Pending => return Poll::Pending,
1542 }
1543 }
1544
1545 Poll::Ready(Ok(out.len()))
1546 }
1547
1548 pub async fn read_frames<'a, const N: usize>(
1551 &mut self,
1552 out: &'a mut frame::Buffer<N>,
1553 ) -> Result<&'a mut [frame::Frame]> {
1554 kio::wait(|waiter| self.poll_read_frames(waiter, out)).await?;
1555 Ok(out.filled_mut())
1556 }
1557
1558 pub fn poll_finished(&mut self, waiter: &kio::Waiter) -> Poll<Result<u64>> {
1560 if self.ended {
1561 return Poll::Ready(Ok(self.index()));
1562 }
1563 if self.expired {
1564 return Poll::Ready(Err(Error::Old));
1565 }
1566 let res = match &mut self.inner {
1567 ConsumerKind::Plain(plain) => {
1568 let index = plain.index;
1569 plain
1570 .poll(waiter, |state| state.poll_end(index))
1571 .map(|res| res.map(|()| index as u64))
1572 }
1573 ConsumerKind::Spliced(spliced) => spliced.poll_finished(waiter),
1574 };
1575 match res.is_pending().then(|| self.poll_expired_if_blocked(waiter)).flatten() {
1576 Some(true) => Poll::Ready(Err(Error::Old)),
1577 Some(false) => Poll::Ready(Ok(self.index())),
1579 None => res,
1580 }
1581 }
1582
1583 pub async fn finished(&mut self) -> Result<u64> {
1590 kio::wait(|waiter| self.poll_finished(waiter)).await
1591 }
1592}
1593
1594impl Plain {
1595 fn expiry_pending(&self) -> bool {
1597 if self.capped() {
1598 return false;
1599 }
1600
1601 let state = self.state.read();
1602 state.abort.is_none() && state.fin.is_none_or(|fin| self.index < fin)
1603 }
1604
1605 fn unread_content(&self) -> stats::Content {
1607 let prefetched = self.prefetch.buffered().0 as usize;
1608 let start = self.index.saturating_add(prefetched);
1609 let end = self.end.unwrap_or(usize::MAX);
1610 self.state.read().content_range(start, end)
1611 }
1612
1613 fn refresh_if_stale(&mut self) {
1615 self.access.touch();
1616 let tick = self.cache.pool().now();
1617 if tick != self.refreshed {
1618 self.state.read().charge.refresh();
1619 self.refreshed = tick;
1620 }
1621 }
1622
1623 fn poll<F, R>(&self, waiter: &kio::Waiter, f: F) -> Poll<Result<R>>
1625 where
1626 F: Fn(&kio::Ref<'_, GroupState>) -> Poll<Result<R>>,
1627 {
1628 Poll::Ready(match ready!(self.state.poll(waiter, f)) {
1629 Ok(res) => res,
1630 Err(state) => Err(state.abort.clone().unwrap_or(Error::Dropped)),
1632 })
1633 }
1634
1635 fn capped(&self) -> bool {
1637 self.end.is_some_and(|end| self.index >= end)
1638 }
1639
1640 fn start_at(&mut self, index: u64) {
1641 let index = usize::try_from(index).unwrap_or(usize::MAX);
1642 let index = index.max(self.state.read().offset);
1643 if index <= self.index {
1644 return;
1645 }
1646 self.index = index;
1647 self.prefetch = Prefetch::default();
1649 }
1650
1651 fn skip_to(&mut self, index: u64) {
1652 let index = usize::try_from(index).unwrap_or(usize::MAX);
1653 if index <= self.index {
1654 return;
1655 }
1656 self.index = index;
1657 self.prefetch = Prefetch::default();
1658 }
1659
1660 fn poll_next_frame(
1661 &mut self,
1662 waiter: &kio::Waiter,
1663 stats: &stats::Meter,
1664 expiry: Option<frame::Expiry>,
1665 ) -> Poll<Result<Option<frame::Consumer>>> {
1666 if self.capped() {
1667 return Poll::Ready(Ok(None));
1668 }
1669 let end = self.end.unwrap_or(usize::MAX);
1670
1671 if let Some(frame) = self.prefetch.pop() {
1675 self.refresh_if_stale();
1676 self.index += 1;
1677 let tail = self.index.saturating_add(self.prefetch.buffered().0 as usize)..end;
1678 let info = frame::Info {
1679 size: frame.payload.len() as u64,
1680 timestamp: frame.timestamp,
1681 };
1682 let source = frame::Source::Complete(frame.payload);
1683 let frame = frame::Consumer::new(self.state.clone(), info, source);
1684 return Poll::Ready(Ok(Some(match expiry {
1685 Some(expiry) => frame.with_expiry(expiry.for_frame(tail, false)),
1686 None => frame,
1687 })));
1688 }
1689
1690 let index = self.index;
1691 let Some((info, source)) = ready!(self.poll(waiter, |state| state.poll_frame_source(index))?) else {
1692 return Poll::Ready(Ok(None));
1693 };
1694
1695 self.index += 1;
1696 stats.frames(1);
1699 let frame = frame::Consumer::new(self.state.clone(), info, source).with_meter(stats.clone());
1700 Poll::Ready(Ok(Some(match expiry {
1701 Some(expiry) => frame.with_expiry(expiry.for_frame(self.index..end, true)),
1702 None => frame,
1703 })))
1704 }
1705
1706 fn poll_read_frame(&mut self, waiter: &kio::Waiter, stats: &stats::Meter) -> Poll<Result<Option<frame::Frame>>> {
1707 if self.capped() {
1708 return Poll::Ready(Ok(None));
1709 }
1710
1711 if let Some(frame) = self.prefetch.pop() {
1713 self.refresh_if_stale();
1714 self.index += 1;
1715 return Poll::Ready(Ok(Some(frame)));
1716 }
1717
1718 let index = self.index;
1721 let budget = self.end.map_or(usize::MAX, |end| end.saturating_sub(index));
1724 let prefetch = &mut self.prefetch;
1725 let res = self.state.poll(waiter, |state| {
1726 if index < state.offset {
1727 return Poll::Ready(Err(Error::Lagged));
1728 }
1729 let local = (index - state.offset).min(state.frames.len());
1734 prefetch.fill(state.frames.range(local..).take(budget).cloned());
1735 if prefetch.len > 0 {
1736 state.charge.refresh();
1739 return Poll::Ready(Ok(()));
1740 }
1741 state.poll_terminal(index)
1744 });
1745
1746 match ready!(res) {
1747 Ok(Ok(())) => {}
1748 Ok(Err(err)) => return Poll::Ready(Err(err)),
1749 Err(state) => return Poll::Ready(Err(state.abort.clone().unwrap_or(Error::Dropped))),
1750 }
1751
1752 self.refreshed = self.cache.pool().now();
1754
1755 let (frames, bytes) = self.prefetch.buffered();
1758 stats.frames(frames);
1759 stats.bytes(bytes);
1760
1761 Poll::Ready(Ok(self.prefetch.pop().inspect(|_| {
1762 self.index += 1;
1763 })))
1764 }
1765}
1766
1767#[derive(Clone, Debug, Default)]
1769#[non_exhaustive]
1770pub struct Fetch {
1771 pub priority: u8,
1773
1774 pub frame_start: u64,
1785}
1786
1787impl Fetch {
1788 pub fn with_priority(mut self, priority: u8) -> Self {
1790 self.priority = priority;
1791 self
1792 }
1793
1794 pub fn with_frame_start(mut self, frame_start: u64) -> Self {
1796 self.frame_start = frame_start;
1797 self
1798 }
1799}
1800
1801pub struct Request {
1810 pub(crate) state: kio::Producer<track::TrackState>,
1811 pub(crate) fetch: kio::Shared<track::FetchState>,
1812 pub(crate) sequence: u64,
1813 pub(crate) priority: u8,
1814 pub(crate) frame_start: u64,
1815 pub(crate) result: kio::Producer<track::FetchOutcome>,
1816 pub(crate) done: bool,
1817}
1818
1819#[cfg(test)]
1820mod test {
1821 use super::*;
1822 use crate::model::test_tracing::count_drop_warnings;
1823 use bytes::Bytes;
1824 use futures::FutureExt;
1825
1826 #[test]
1829 fn one_frame_fits_the_charged_slots() {
1830 let mut frames: VecDeque<Frame> = VecDeque::new();
1831 frames.push_back(Frame {
1832 timestamp: Timestamp::ZERO,
1833 payload: Bytes::new(),
1834 });
1835 let capacity = frames.capacity();
1836 assert!(
1837 capacity <= FRAME_SLOTS,
1838 "a one-frame deque now allocates {capacity} slots"
1839 );
1840 }
1841
1842 #[test]
1843 fn basic_frame_reading() {
1844 let mut producer = Info { sequence: 0 }.produce();
1845 producer
1846 .write_frame(Timestamp::ZERO, Bytes::from_static(b"frame0"))
1847 .unwrap();
1848 producer
1849 .write_frame(Timestamp::ZERO, Bytes::from_static(b"frame1"))
1850 .unwrap();
1851 producer.finish().unwrap();
1852
1853 let mut consumer = producer.consume();
1854 let f0 = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
1855 assert_eq!(f0.size, 6);
1856 let f1 = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
1857 assert_eq!(f1.size, 6);
1858 let end = consumer.next_frame().now_or_never().unwrap().unwrap();
1859 assert!(end.is_none());
1860 }
1861
1862 #[test]
1863 fn read_frame_all_at_once() {
1864 let mut producer = Info { sequence: 0 }.produce();
1865 producer
1866 .write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
1867 .unwrap();
1868 producer.finish().unwrap();
1869
1870 let mut consumer = producer.consume();
1871 let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
1872 assert_eq!(frame.payload, Bytes::from_static(b"hello"));
1873 }
1874
1875 #[test]
1876 fn read_frame_preserves_timestamp() {
1877 let mut producer = Info { sequence: 0 }.produce();
1878 let timestamp = Timestamp::from_micros(20_000).unwrap();
1879 producer.write_frame(timestamp, Bytes::from_static(b"hello")).unwrap();
1880 producer.finish().unwrap();
1881
1882 let mut consumer = producer.consume();
1883 let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
1884 assert_eq!(frame.timestamp.as_micros(), 20_000);
1885 assert_eq!(frame.payload, Bytes::from_static(b"hello"));
1886 }
1887
1888 #[test]
1889 fn chunked_frame_reads_whole() {
1890 let mut producer = Info { sequence: 0 }.produce();
1891 {
1892 let mut frame = producer
1893 .create_frame(frame::Info {
1894 size: 10,
1895 timestamp: Timestamp::ZERO,
1896 })
1897 .unwrap();
1898 frame.write(Bytes::from_static(b"hello")).unwrap();
1899 frame.write(Bytes::from_static(b"world")).unwrap();
1900 frame.finish().unwrap();
1901 }
1902 producer.finish().unwrap();
1903
1904 let mut consumer = producer.consume();
1907 let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
1908 assert_eq!(frame.payload, Bytes::from_static(b"helloworld"));
1909 }
1910
1911 #[test]
1912 fn chunked_frame_streams_partial() {
1913 let mut producer = Info { sequence: 0 }.produce();
1914 let mut consumer = producer.consume();
1915
1916 let mut frame = producer
1917 .create_frame(frame::Info {
1918 size: 6,
1919 timestamp: Timestamp::ZERO,
1920 })
1921 .unwrap();
1922 frame.write(Bytes::from_static(b"foo")).unwrap();
1923
1924 let mut f = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
1926 let c1 = f.read_chunk().now_or_never().unwrap().unwrap();
1927 assert_eq!(c1, Some(Bytes::from_static(b"foo")));
1928 assert!(f.read_chunk().now_or_never().is_none());
1929
1930 frame.write(Bytes::from_static(b"bar")).unwrap();
1931 frame.finish().unwrap();
1932
1933 let c2 = f.read_chunk().now_or_never().unwrap().unwrap();
1934 assert_eq!(c2, Some(Bytes::from_static(b"bar")));
1935 let c3 = f.read_chunk().now_or_never().unwrap().unwrap();
1936 assert_eq!(c3, None);
1937 }
1938
1939 #[test]
1940 fn group_finish_returns_none() {
1941 let producer = Info { sequence: 0 }.produce();
1942 producer.finish().unwrap();
1943
1944 let mut consumer = producer.consume();
1945 let end = consumer.next_frame().now_or_never().unwrap().unwrap();
1946 assert!(end.is_none());
1947 }
1948
1949 #[test]
1950 fn abort_propagates() {
1951 let producer = Info { sequence: 0 }.produce();
1952 let mut consumer = producer.consume();
1953 producer.abort(crate::Error::Cancel).unwrap();
1954
1955 let result = consumer.next_frame().now_or_never().unwrap();
1956 assert!(matches!(result, Err(crate::Error::Cancel)));
1957 }
1958
1959 #[test]
1960 fn abort_clears_cached_frames() {
1961 let mut producer = Info { sequence: 0 }.produce();
1962 producer
1963 .write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
1964 .unwrap();
1965
1966 let _consumer = producer.consume();
1968 assert_eq!(producer.state.read().frames.len(), 1);
1969
1970 producer.clone().abort(crate::Error::Cancel).unwrap();
1971
1972 let state = producer.state.read();
1973 assert!(state.frames.is_empty(), "cached frames should be dropped on abort");
1974 assert_eq!(state.cache, 0);
1975 }
1976
1977 #[test]
1978 fn drop_unfinished_clears_cached_frames() {
1979 let producer = Info { sequence: 0 }.produce();
1980 let mut writer = producer.clone();
1981 writer
1982 .write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
1983 .unwrap();
1984
1985 let mut consumer = producer.consume();
1987 assert_eq!(producer.state.read().frames.len(), 1);
1988
1989 drop(writer);
1991 drop(producer);
1992
1993 let result = consumer.next_frame().now_or_never().unwrap();
1994 assert!(matches!(result, Err(crate::Error::Dropped)));
1995 }
1996
1997 #[test]
1998 fn drop_after_abort_does_not_warn() {
1999 let warns = count_drop_warnings("group::Producer dropped without finish", || {
2000 let producer = Info { sequence: 0 }.produce();
2001 let keep = producer.clone();
2002 let mut writer = producer.clone();
2003 writer
2004 .write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
2005 .unwrap();
2006 let _consumer = producer.consume();
2007 writer.abort(crate::Error::Cancel).unwrap();
2008 drop(keep);
2009 });
2010 assert_eq!(warns, 0, "abort-then-drop must not emit unfinished-producer WARN");
2011 }
2012
2013 #[test]
2014 fn drop_unfinished_warns() {
2015 let warns = count_drop_warnings("group::Producer dropped without finish", || {
2016 let producer = Info { sequence: 0 }.produce();
2017 let mut writer = producer.clone();
2018 writer
2019 .write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
2020 .unwrap();
2021 let _consumer = producer.consume();
2022 drop(writer);
2023 drop(producer);
2024 });
2025 assert!(warns >= 1, "unfinished drop must emit unfinished-producer WARN");
2026 }
2027
2028 #[test]
2029 fn drop_finished_keeps_cached_frames() {
2030 let mut producer = Info { sequence: 0 }.produce();
2031 producer
2032 .write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
2033 .unwrap();
2034 producer.finish().unwrap();
2035
2036 let mut consumer = producer.consume();
2037 drop(producer);
2038
2039 let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
2041 assert_eq!(frame.payload, Bytes::from_static(b"data"));
2042 }
2043
2044 #[tokio::test]
2045 async fn pending_then_ready() {
2046 let mut producer = Info { sequence: 0 }.produce();
2047 let mut consumer = producer.consume();
2048
2049 assert!(consumer.next_frame().now_or_never().is_none());
2051
2052 producer
2053 .write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
2054 .unwrap();
2055 producer.finish().unwrap();
2056
2057 let frame = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
2058 assert_eq!(frame.size, 4);
2059 }
2060
2061 #[test]
2062 fn overflow_aborts_the_group() {
2063 let mut producer = Info { sequence: 0 }.produce();
2064 let mut consumer = producer.consume();
2065
2066 let big = Bytes::from(vec![0u8; MAX_CACHE_BYTES as usize]);
2067 producer.write_frame(Timestamp::ZERO, big.clone()).unwrap();
2068 assert!(matches!(
2069 producer.write_frame(Timestamp::ZERO, big),
2070 Err(Error::GroupTooLarge)
2071 ));
2072
2073 {
2074 let state = producer.state.read();
2075 assert!(matches!(state.abort, Some(Error::GroupTooLarge)));
2076 assert!(state.frames.is_empty());
2077 assert_eq!(state.offset, 0);
2078 }
2079
2080 let result = consumer.next_frame().now_or_never().unwrap();
2081 assert!(matches!(result, Err(Error::GroupTooLarge)));
2082 }
2083
2084 #[test]
2085 fn no_overflow_under_budget() {
2086 let mut producer = Info { sequence: 0 }.produce();
2087 for _ in 0..MAX_GROUP_FRAMES {
2089 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"x")).unwrap();
2090 }
2091 producer.finish().unwrap();
2092
2093 let state = producer.state.read();
2094 assert_eq!(state.offset, 0);
2095 assert_eq!(state.frames.len(), MAX_GROUP_FRAMES);
2096 assert!(state.abort.is_none());
2097 }
2098
2099 #[test]
2100 fn writer_sees_group_too_large_on_the_8193rd_frame() {
2101 let mut producer = Info { sequence: 0 }.produce();
2102 for _ in 0..MAX_GROUP_FRAMES {
2103 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"x")).unwrap();
2104 }
2105 assert!(matches!(
2106 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"x")),
2107 Err(Error::GroupTooLarge)
2108 ));
2109 assert!(matches!(producer.state.read().abort, Some(Error::GroupTooLarge)));
2110 }
2111
2112 #[test]
2113 fn clone_consumer_independent() {
2114 let mut producer = Info { sequence: 0 }.produce();
2115 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"a")).unwrap();
2116
2117 let mut c1 = producer.consume();
2118 let _ = c1.next_frame().now_or_never().unwrap().unwrap().unwrap();
2120
2121 let mut c2 = c1.clone();
2123
2124 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"b")).unwrap();
2125 producer.finish().unwrap();
2126
2127 let f = c2.next_frame().now_or_never().unwrap().unwrap().unwrap();
2129 assert_eq!(f.size, 1); let end = c2.next_frame().now_or_never().unwrap().unwrap();
2132 assert!(end.is_none());
2133 }
2134
2135 fn prefetched_consumer(pool: &cache::Pool, max_age: std::time::Duration) -> (Producer, Consumer) {
2136 let cache = cache::Track::new(pool.clone(), kio::Weak::new());
2137 let track = track::Info::default().with_max_age(max_age);
2138 let mut producer = Producer::new(Info { sequence: 0 }, track, cache);
2139 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"a")).unwrap();
2140 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"b")).unwrap();
2141 producer.finish().unwrap();
2142
2143 let mut consumer = producer.consume();
2144 consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
2145 (producer, consumer)
2146 }
2147
2148 #[test]
2149 fn prefetch_refresh_honors_pool_expiry() {
2150 let config = cache::Config::default().with_expiry(std::time::Duration::from_secs(1));
2151 let pool = cache::Pool::new(config);
2152 let (producer, mut consumer) = prefetched_consumer(&pool, std::time::Duration::MAX);
2153 let before = producer.cache_accessed();
2154
2155 crate::model::clock::advance(std::time::Duration::from_millis(600));
2156 consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
2157
2158 assert!(producer.cache_accessed() > before, "the pool cadence is used");
2159 }
2160
2161 #[test]
2162 fn prefetch_refresh_honors_track_max_age() {
2163 let config = cache::Config::default().with_expiry(std::time::Duration::from_secs(30));
2164 let pool = cache::Pool::new(config);
2165 let (producer, mut consumer) = prefetched_consumer(&pool, std::time::Duration::from_secs(1));
2166 let before = producer.cache_accessed();
2167
2168 crate::model::clock::advance(std::time::Duration::from_millis(600));
2169 consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
2170
2171 assert!(producer.cache_accessed() > before, "the track cadence remains in force");
2172 }
2173
2174 #[test]
2177 fn read_frame_crosses_prefetch_batches() {
2178 let n = Prefetch::CAP * 3 + 5;
2179 let mut producer = Info { sequence: 0 }.produce();
2180 for i in 0..n {
2181 producer
2182 .write_frame(Timestamp::ZERO, Bytes::from(vec![i as u8; 4]))
2183 .unwrap();
2184 }
2185 producer.finish().unwrap();
2186
2187 let mut consumer = producer.consume();
2188 for i in 0..n {
2189 let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
2190 assert_eq!(frame.payload, Bytes::from(vec![i as u8; 4]));
2191 }
2192 assert!(consumer.read_frame().now_or_never().unwrap().unwrap().is_none());
2193 }
2194
2195 #[test]
2199 fn abort_after_finish_keeps_the_clean_end_for_a_drained_reader() {
2200 let mut producer = Info { sequence: 0 }.produce();
2201 producer
2202 .write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
2203 .unwrap();
2204 producer.finish().unwrap();
2205
2206 let mut drained = producer.consume();
2207 let mut behind = producer.consume();
2208 let frame = drained.read_frame().now_or_never().unwrap().unwrap().unwrap();
2209 assert_eq!(frame.payload, Bytes::from_static(b"hello"));
2210
2211 producer.abort(Error::Old).unwrap();
2212
2213 assert!(drained.read_frame().now_or_never().unwrap().unwrap().is_none());
2215 assert!(drained.next_frame().now_or_never().unwrap().unwrap().is_none());
2216
2217 assert!(matches!(behind.read_frame().now_or_never().unwrap(), Err(Error::Old)));
2219 }
2220
2221 #[test]
2225 fn finished_answers_for_the_cursor() {
2226 let mut producer = Info { sequence: 0 }.produce();
2227 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"a")).unwrap();
2228 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"b")).unwrap();
2229 producer.finish().unwrap();
2230
2231 let mut drained = producer.consume();
2232 let mut behind = producer.consume();
2233 while drained.read_frame().now_or_never().unwrap().unwrap().is_some() {}
2234 behind.read_frame().now_or_never().unwrap().unwrap().unwrap();
2235
2236 producer.abort(Error::Old).unwrap();
2237
2238 assert_eq!(drained.finished().now_or_never().unwrap().unwrap(), 2);
2239 assert!(matches!(behind.finished().now_or_never().unwrap(), Err(Error::Old)));
2240 assert_eq!(behind.frame_count(), 2);
2241 }
2242
2243 #[test]
2246 fn finished_reports_a_group_too_large() {
2247 let mut producer = Info { sequence: 0 }.produce();
2248 let mut consumer = producer.consume();
2249
2250 let big = Bytes::from(vec![0u8; MAX_CACHE_BYTES as usize]);
2251 producer.write_frame(Timestamp::ZERO, big.clone()).unwrap();
2252 assert!(matches!(
2253 producer.write_frame(Timestamp::ZERO, big),
2254 Err(Error::GroupTooLarge)
2255 ));
2256
2257 assert!(matches!(
2258 consumer.finished().now_or_never().unwrap(),
2259 Err(Error::GroupTooLarge)
2260 ));
2261 }
2262
2263 #[test]
2265 fn interleave_read_and_next_frame() {
2266 let mut producer = Info { sequence: 0 }.produce();
2267 for i in 0..5u8 {
2268 producer.write_frame(Timestamp::ZERO, Bytes::from(vec![i; 1])).unwrap();
2269 }
2270 producer.finish().unwrap();
2271
2272 let mut consumer = producer.consume();
2273 let f0 = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
2275 assert_eq!(f0.payload, Bytes::from(vec![0u8; 1]));
2276
2277 for i in 1..5u8 {
2279 let mut f = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
2280 let data = f.read_all().now_or_never().unwrap().unwrap();
2281 assert_eq!(data, Bytes::from(vec![i; 1]));
2282 }
2283 assert!(consumer.next_frame().now_or_never().unwrap().unwrap().is_none());
2284 }
2285
2286 #[test]
2289 fn read_frame_past_cleared_frames_does_not_panic() {
2290 let mut producer = Info { sequence: 0 }.produce();
2291 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"a")).unwrap();
2292 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"b")).unwrap();
2293
2294 let mut consumer = producer.consume();
2295 consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
2296 consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
2297
2298 producer.abort(Error::Cancel).unwrap();
2301
2302 let result = consumer.read_frame().now_or_never().unwrap();
2303 assert!(matches!(result, Err(Error::Cancel)), "expected Cancel, got {result:?}");
2304 }
2305
2306 #[test]
2309 fn drop_with_partial_batch() {
2310 let mut producer = Info { sequence: 0 }.produce();
2311 for _ in 0..Prefetch::CAP {
2312 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"x")).unwrap();
2313 }
2314 producer.finish().unwrap();
2315
2316 let mut consumer = producer.consume();
2317 let _ = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
2319 drop(consumer);
2320 }
2321
2322 #[tokio::test]
2327 async fn chunk_write_wakes_parked_reader() {
2328 let mut producer = Info { sequence: 0 }.produce();
2329 let mut consumer = producer.consume();
2330 let mut frame = producer
2331 .create_frame(frame::Info {
2332 size: 6,
2333 timestamp: Timestamp::ZERO,
2334 })
2335 .unwrap();
2336 let mut f = consumer.next_frame().await.unwrap().unwrap();
2337 let handle = tokio::spawn(async move { f.read_chunk().await });
2338 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
2340 frame.write(Bytes::from_static(b"foo")).unwrap();
2341 let chunk = tokio::time::timeout(std::time::Duration::from_secs(2), handle)
2342 .await
2343 .expect("parked chunk reader was never woken by the chunk write")
2344 .unwrap()
2345 .unwrap();
2346 assert_eq!(chunk, Some(Bytes::from_static(b"foo")));
2347 }
2348
2349 #[test]
2352 fn create_frame_converts_mismatched_scale() {
2353 use crate::{Timescale, Timestamp};
2354
2355 let mut producer = Producer::new(
2356 Info { sequence: 0 },
2357 track::Info::default().with_timescale(Timescale::MICRO),
2358 Default::default(),
2359 );
2360 let frame = frame::Info {
2361 size: 3,
2362 timestamp: Timestamp::from_millis(1).unwrap(), };
2364 let writer = producer.create_frame(frame).unwrap();
2365 assert_eq!(writer.timestamp.scale(), Timescale::MICRO);
2366 assert_eq!(writer.timestamp.value(), 1000);
2367 }
2368
2369 #[tokio::test]
2371 async fn create_frame_converts_current_timestamp() {
2372 use crate::Timescale;
2373
2374 let mut producer = Producer::new(
2375 Info { sequence: 0 },
2376 track::Info::default().with_timescale(Timescale::MICRO),
2377 Default::default(),
2378 );
2379 let writer = producer
2380 .create_frame(frame::Info {
2381 size: 3,
2382 timestamp: Timestamp::now(),
2383 })
2384 .unwrap();
2385 assert_eq!(writer.timestamp.scale(), Timescale::MICRO);
2386 assert!(!writer.timestamp.is_zero(), "local clock should be non-zero");
2387 }
2388
2389 #[test]
2392 fn start_at_starts_the_group_later() {
2393 let mut producer = Info { sequence: 0 }.produce();
2394 producer.start_at(3).unwrap();
2395 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"d")).unwrap();
2396 producer.finish().unwrap();
2397
2398 assert_eq!(producer.frame_count(), 4);
2400
2401 let mut consumer = producer.consume();
2402 assert_eq!(consumer.frame_count(), 4);
2403
2404 assert!(matches!(
2407 consumer.finished().now_or_never().unwrap(),
2408 Err(Error::Lagged)
2409 ));
2410 assert!(matches!(
2411 consumer.read_frame().now_or_never().unwrap(),
2412 Err(Error::Lagged)
2413 ));
2414 }
2415
2416 #[test]
2419 fn start_at_clamps_up_to_the_first_frame() {
2420 let mut producer = Info { sequence: 0 }.produce();
2421 producer.start_at(3).unwrap();
2422 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"d")).unwrap();
2423 producer.finish().unwrap();
2424
2425 let mut consumer = producer.consume();
2426 consumer.start_at(1);
2427 assert_eq!(consumer.index(), 3, "clamped up to the first frame that exists");
2428 assert_eq!(
2429 consumer.read_frame().now_or_never().unwrap().unwrap().unwrap().payload,
2430 Bytes::from_static(b"d")
2431 );
2432 }
2433
2434 #[test]
2437 fn end_at_caps_and_reopens() {
2438 let mut producer = Info { sequence: 0 }.produce();
2439 for i in 0..4u8 {
2440 producer.write_frame(Timestamp::ZERO, Bytes::from(vec![i])).unwrap();
2441 }
2442 producer.finish().unwrap();
2443
2444 let mut consumer = producer.consume();
2445 consumer.set_frames(..2);
2446 assert_eq!(
2447 consumer.read_frame().now_or_never().unwrap().unwrap().unwrap().payload[0],
2448 0
2449 );
2450 assert_eq!(
2451 consumer.read_frame().now_or_never().unwrap().unwrap().unwrap().payload[0],
2452 1
2453 );
2454 assert!(
2455 consumer.read_frame().now_or_never().unwrap().unwrap().is_none(),
2456 "capped reads end cleanly"
2457 );
2458
2459 consumer.set_frames(..);
2460 assert_eq!(
2461 consumer.read_frame().now_or_never().unwrap().unwrap().unwrap().payload[0],
2462 2
2463 );
2464 }
2465
2466 #[test]
2467 fn frame_ranges_preserve_progress_and_make_inclusion_explicit() {
2468 let mut producer = Info { sequence: 0 }.produce();
2469 for i in 0..4u8 {
2470 producer.write_frame(Timestamp::ZERO, Bytes::from(vec![i])).unwrap();
2471 }
2472 producer.finish().unwrap();
2473 let mut consumer = producer.consume();
2474 consumer.set_frames(1..=1);
2475 assert_eq!(
2476 consumer.read_frame().now_or_never().unwrap().unwrap().unwrap().payload[0],
2477 1
2478 );
2479 assert!(consumer.read_frame().now_or_never().unwrap().unwrap().is_none());
2480 consumer.set_frames(..3);
2481 assert_eq!(
2482 consumer.read_frame().now_or_never().unwrap().unwrap().unwrap().payload[0],
2483 2
2484 );
2485 assert!(consumer.read_frame().now_or_never().unwrap().unwrap().is_none());
2486 consumer.set_frames(0..=3);
2487 assert_eq!(
2488 consumer.read_frame().now_or_never().unwrap().unwrap().unwrap().payload[0],
2489 3
2490 );
2491 }
2492
2493 #[test]
2496 fn end_at_zero_is_empty() {
2497 let mut producer = Info { sequence: 0 }.produce();
2498 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"x")).unwrap();
2499 producer.finish().unwrap();
2500
2501 let mut consumer = producer.consume();
2502 consumer.set_frames(..0);
2503 assert!(
2504 consumer.read_frame().now_or_never().unwrap().unwrap().is_none(),
2505 "empty cap delivers nothing"
2506 );
2507
2508 consumer.set_frames(..1);
2509 assert_eq!(
2510 consumer.read_frame().now_or_never().unwrap().unwrap().unwrap().payload,
2511 Bytes::from_static(b"x")
2512 );
2513 }
2514
2515 #[test]
2517 fn start_at_rejected_after_a_frame() {
2518 let mut producer = Info { sequence: 0 }.produce();
2519 producer.start_at(2).unwrap();
2521 producer.start_at(3).unwrap();
2522
2523 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"a")).unwrap();
2524 assert!(matches!(producer.start_at(4), Err(Error::Closed)));
2525 assert_eq!(producer.frame_count(), 4, "the frame landed at index 3");
2526
2527 let mut producer = Info { sequence: 1 }.produce();
2529 producer.finish().unwrap();
2530 assert!(matches!(producer.start_at(1), Err(Error::Closed)));
2531 }
2532
2533 #[test]
2535 fn start_at_rejects_the_largest_index() {
2536 let mut producer = Info { sequence: 0 }.produce();
2537 assert!(matches!(
2538 producer.start_at(usize::MAX as u64),
2539 Err(Error::BoundsExceeded(_))
2540 ));
2541 }
2542
2543 #[test]
2545 fn create_frame_rejects_oversized() {
2546 let mut producer = Info { sequence: 0 }.produce();
2547 let result = producer.create_frame(frame::Info {
2548 size: MAX_CACHE_BYTES + 1,
2549 timestamp: Timestamp::ZERO,
2550 });
2551 assert!(matches!(result, Err(Error::FrameTooLarge)));
2552 }
2553}