1use crate::cache;
12use crate::frame::{self, Frame, FrameBuf};
13use crate::{Timescale, stats, track};
14use std::collections::VecDeque;
15use std::mem::MaybeUninit;
16use std::sync::Arc;
17use std::task::{Poll, ready};
18
19use crate::{Error, IntoBytes, Result, Timestamp};
20
21const MAX_GROUP_CACHE: u64 = 32 * 1024 * 1024; #[derive(Clone, Copy, Debug, Hash, Eq, PartialEq, Ord, PartialOrd)]
32pub struct Info {
33 pub sequence: u64,
36}
37
38impl Info {
39 #[cfg(test)]
45 pub(crate) fn produce(self) -> Producer {
46 Producer::new(self, track::Info::default())
47 }
48}
49
50impl From<usize> for Info {
51 fn from(sequence: usize) -> Self {
52 Self {
53 sequence: sequence as u64,
54 }
55 }
56}
57
58impl From<u64> for Info {
59 fn from(sequence: u64) -> Self {
60 Self { sequence }
61 }
62}
63
64impl From<u32> for Info {
65 fn from(sequence: u32) -> Self {
66 Self {
67 sequence: sequence as u64,
68 }
69 }
70}
71
72impl From<u16> for Info {
73 fn from(sequence: u16) -> Self {
74 Self {
75 sequence: sequence as u64,
76 }
77 }
78}
79
80pub(crate) struct Partial {
83 timestamp: Timestamp,
84 buf: FrameBuf,
85}
86
87#[derive(Default)]
90pub(crate) struct GroupState {
91 pub(crate) frames: VecDeque<Frame>,
94
95 pub(crate) partial: Option<Partial>,
97
98 pub(crate) offset: usize,
100
101 pub(crate) cache: u64,
103
104 charge: cache::Charge,
107
108 pub(crate) fin: Option<usize>,
111
112 pub(crate) abort: Option<Error>,
114}
115
116impl GroupState {
117 fn poll_frame_source(&self, index: usize) -> Poll<Result<Option<(frame::Info, frame::Source)>>> {
120 if index < self.offset {
121 return Poll::Ready(Err(Error::Lagged));
122 }
123 let local = index - self.offset;
124 if let Some(f) = self.frames.get(local) {
125 self.charge.touch();
126 let info = frame::Info {
127 size: f.payload.len() as u64,
128 timestamp: f.timestamp,
129 };
130 return Poll::Ready(Ok(Some((info, frame::Source::Complete(f.payload.clone())))));
131 }
132 if local == self.frames.len()
133 && let Some(p) = &self.partial
134 {
135 self.charge.touch();
136 let info = frame::Info {
137 size: p.buf.capacity() as u64,
138 timestamp: p.timestamp,
139 };
140 return Poll::Ready(Ok(Some((info, frame::Source::Partial(p.buf.clone())))));
141 }
142 ready!(self.poll_terminal(index))?;
143 Poll::Ready(Ok(None))
144 }
145
146 fn poll_terminal(&self, index: usize) -> Poll<Result<()>> {
153 match (self.fin, &self.abort) {
154 (Some(total), Some(err)) if index < total => Poll::Ready(Err(err.clone())),
155 (Some(_), _) => Poll::Ready(Ok(())),
156 (None, Some(err)) => Poll::Ready(Err(err.clone())),
157 (None, None) => Poll::Pending,
158 }
159 }
160
161 fn poll_finished(&self) -> Poll<Result<u64>> {
162 if let Some(total) = self.fin {
165 Poll::Ready(Ok(total as u64))
166 } else if let Some(err) = &self.abort {
167 Poll::Ready(Err(err.clone()))
168 } else {
169 Poll::Pending
170 }
171 }
172
173 fn evict(&mut self) {
175 while self.cache > MAX_GROUP_CACHE {
176 let Some(frame) = self.frames.pop_front() else {
177 break;
178 };
179 let size = frame.payload.len() as u64;
180 self.cache -= size;
181 self.charge.sub(size);
182 self.offset += 1;
183 }
184 }
185
186 fn release(&mut self) {
188 self.frames.clear();
189 self.partial = None;
190 self.cache = 0;
191 self.charge.clear();
192 }
193}
194
195fn modify(state: &kio::Producer<GroupState>) -> Result<kio::Mut<'_, GroupState>> {
196 state.write().map_err(|r| r.abort.clone().unwrap_or(Error::Dropped))
197}
198
199fn evict(weak: &kio::ProducerWeak<GroupState>) {
202 let Some(producer) = weak.produce() else { return };
204 let Ok(mut state) = producer.write() else { return };
205 if state.abort.is_some() {
206 return;
207 }
208 state.abort = Some(Error::Evicted);
209 state.release();
210 state.close();
211}
212
213pub struct Producer {
219 state: kio::Producer<GroupState>,
221
222 info: Info,
225
226 track: track::Info,
231
232 stats: stats::Meter,
235}
236
237impl std::ops::Deref for Producer {
238 type Target = Info;
239
240 fn deref(&self) -> &Self::Target {
241 &self.info
242 }
243}
244
245impl Producer {
246 pub(crate) fn new(info: Info, track: track::Info) -> Self {
257 let state = kio::Producer::<GroupState>::default();
258 let weak = state.weak();
259 let charge = track.broadcast.origin.pool.register(Box::new(move || evict(&weak)));
260 state.write().ok().expect("a new group is open").charge = charge;
261 Self {
262 info,
263 state,
264 track,
265 stats: stats::Meter::default(),
266 }
267 }
268
269 pub(crate) fn with_meter(mut self, meter: stats::Meter) -> Self {
272 meter.group();
273 self.stats = meter;
274 self
275 }
276
277 pub(crate) fn info(&self) -> Info {
279 self.info
280 }
281
282 pub fn timescale(&self) -> Timescale {
284 self.track.timescale
285 }
286
287 pub fn write_frame<B: IntoBytes>(&mut self, timestamp: Timestamp, data: B) -> Result<()> {
295 let timestamp = timestamp
296 .convert(self.track.timescale)
297 .map_err(|_| Error::TimestampMismatch)?;
298 let payload = data.into_bytes();
299 if payload.len() as u64 > MAX_GROUP_CACHE {
300 return Err(Error::FrameTooLarge);
301 }
302
303 let mut state = modify(&self.state)?;
304 if state.fin.is_some() {
305 return Err(Error::Closed);
306 }
307 debug_assert!(state.partial.is_none(), "a frame is already open");
308 let size = payload.len() as u64;
309 state.cache += size;
310 state.charge.add(size);
311 state.frames.push_back(Frame { timestamp, payload });
312 state.evict();
313
314 drop(state);
317 self.track.broadcast.origin.pool.evict();
318
319 self.stats.frames(1);
321 self.stats.bytes(size);
322 Ok(())
323 }
324
325 pub fn create_frame(&mut self, frame: frame::Info) -> Result<frame::Producer<'_>> {
334 let timestamp = frame
335 .timestamp
336 .convert(self.track.timescale)
337 .map_err(|_| Error::TimestampMismatch)?;
338 if frame.size > MAX_GROUP_CACHE {
339 return Err(Error::FrameTooLarge);
340 }
341 let buf = FrameBuf::new(frame.size as usize);
342
343 let mut state = modify(&self.state)?;
344 if state.fin.is_some() {
345 return Err(Error::Closed);
346 }
347 debug_assert!(state.partial.is_none(), "a frame is already open");
348 state.cache += frame.size;
349 state.charge.add(frame.size);
350 state.partial = Some(Partial {
351 timestamp,
352 buf: buf.clone(),
353 });
354 state.evict();
355
356 drop(state);
359 self.track.broadcast.origin.pool.evict();
360
361 self.stats.frames(1);
364 let meter = self.stats.clone();
365
366 let info = frame::Info {
367 size: frame.size,
368 timestamp,
369 };
370 Ok(frame::Producer::new(self, buf, info).with_meter(meter))
371 }
372
373 pub(crate) fn frame_notify(&self) {
375 let _ = self.state.write();
377 }
378
379 pub(crate) fn frame_commit(&mut self, frame: Frame) -> Result<()> {
381 let mut state = modify(&self.state)?;
382 state.partial = None;
385 state.frames.push_back(frame);
386 Ok(())
387 }
388
389 pub(crate) fn frame_abort(&mut self, err: Error) {
392 let _ = self.clone().abort(err);
393 }
394
395 pub fn frame_count(&self) -> usize {
397 let state = self.state.read();
398 state.offset + state.frames.len() + state.partial.is_some() as usize
399 }
400
401 pub fn finish(&mut self) -> Result<()> {
406 let mut state = modify(&self.state)?;
407 state.fin = Some(state.offset + state.frames.len());
408 Ok(())
409 }
410
411 pub fn abort(self, err: Error) -> Result<()> {
417 let mut guard = modify(&self.state)?;
418 guard.abort = Some(err);
419 guard.release();
420 guard.close();
421 Ok(())
422 }
423
424 pub(crate) fn is_aborted(&self) -> bool {
427 self.state.read().abort.is_some()
428 }
429
430 pub(crate) fn cache_entry(&self) -> Option<Arc<cache::Entry>> {
433 self.state.read().charge.entry()
434 }
435
436 pub fn consume(&self) -> Consumer {
438 Consumer {
439 info: self.info,
440 state: self.state.consume(),
441 track: self.track.clone(),
442 index: 0,
443 prefetch: Prefetch::default(),
444 stats: stats::Meter::default(),
447 }
448 }
449
450 pub async fn closed(&self) -> Error {
452 kio::wait(|waiter| self.poll_closed(waiter)).await
453 }
454
455 pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<Error> {
457 self.state.poll_closed(waiter).map(|()| self.abort_reason())
458 }
459
460 pub async fn unused(&self) -> Result<()> {
462 self.state.unused().await.map_err(|_| self.abort_reason())
463 }
464
465 fn abort_reason(&self) -> Error {
467 self.state.read().abort.clone().unwrap_or(Error::Dropped)
468 }
469}
470
471impl Clone for Producer {
472 fn clone(&self) -> Self {
473 Self {
474 info: self.info,
475 state: self.state.clone(),
476 track: self.track.clone(),
477 stats: self.stats.clone(),
478 }
479 }
480}
481
482impl Drop for Producer {
483 fn drop(&mut self) {
484 if !self.state.is_last() {
488 return;
489 }
490 if let Ok(mut state) = modify(&self.state)
491 && state.fin.is_none()
492 {
493 tracing::warn!(
496 sequence = self.info.sequence,
497 "group::Producer dropped without finish() or abort()"
498 );
499 state.release();
500 }
501 }
502}
503
504struct Prefetch {
512 frames: [MaybeUninit<Frame>; Self::CAP],
514 pos: usize,
515 len: usize,
516}
517
518impl Prefetch {
519 const CAP: usize = 8;
520
521 fn pop(&mut self) -> Option<Frame> {
523 if self.pos == self.len {
524 return None;
525 }
526 let frame = unsafe { self.frames[self.pos].assume_init_read() };
528 self.pos += 1;
529 Some(frame)
530 }
531
532 fn fill(&mut self, frames: impl Iterator<Item = Frame>) {
534 debug_assert_eq!(self.pos, self.len, "fill on a non-empty batch would leak frames");
535 self.pos = 0;
536 self.len = 0;
537 for frame in frames.take(Self::CAP) {
538 self.frames[self.len].write(frame);
539 self.len += 1;
540 }
541 }
542
543 fn buffered(&self) -> (u64, u64) {
546 let mut bytes = 0u64;
547 for slot in &self.frames[self.pos..self.len] {
548 bytes += unsafe { slot.assume_init_ref() }.payload.len() as u64;
550 }
551 ((self.len - self.pos) as u64, bytes)
552 }
553}
554
555impl Default for Prefetch {
556 fn default() -> Self {
557 Self {
558 frames: [const { MaybeUninit::uninit() }; Self::CAP],
559 pos: 0,
560 len: 0,
561 }
562 }
563}
564
565impl Drop for Prefetch {
566 fn drop(&mut self) {
567 for slot in &mut self.frames[self.pos..self.len] {
568 unsafe { slot.assume_init_drop() };
570 }
571 }
572}
573
574pub struct Consumer {
576 state: kio::Consumer<GroupState>,
578
579 info: Info,
581
582 track: track::Info,
585
586 index: usize,
589
590 prefetch: Prefetch,
592
593 stats: stats::Meter,
596}
597
598impl Clone for Consumer {
599 fn clone(&self) -> Self {
600 Self {
603 state: self.state.clone(),
604 info: self.info,
605 track: self.track.clone(),
606 index: self.index,
607 prefetch: Prefetch::default(),
608 stats: self.stats.clone(),
611 }
612 }
613}
614
615impl std::ops::Deref for Consumer {
616 type Target = Info;
617
618 fn deref(&self) -> &Self::Target {
619 &self.info
620 }
621}
622
623impl Consumer {
624 pub(crate) fn with_meter(mut self, meter: stats::Meter) -> Self {
627 meter.group();
628 self.stats = meter;
629 self
630 }
631
632 pub fn timescale(&self) -> Timescale {
634 self.track.timescale
635 }
636
637 fn poll<F, R>(&self, waiter: &kio::Waiter, f: F) -> Poll<Result<R>>
639 where
640 F: Fn(&kio::Ref<'_, GroupState>) -> Poll<Result<R>>,
641 {
642 Poll::Ready(match ready!(self.state.poll(waiter, f)) {
643 Ok(res) => res,
644 Err(state) => Err(state.abort.clone().unwrap_or(Error::Dropped)),
646 })
647 }
648
649 pub async fn next_frame(&mut self) -> Result<Option<frame::Consumer>> {
651 kio::wait(|waiter| self.poll_next_frame(waiter)).await
652 }
653
654 pub fn poll_next_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Consumer>>> {
658 if let Some(frame) = self.prefetch.pop() {
662 self.index += 1;
663 let info = frame::Info {
664 size: frame.payload.len() as u64,
665 timestamp: frame.timestamp,
666 };
667 let source = frame::Source::Complete(frame.payload);
668 return Poll::Ready(Ok(Some(frame::Consumer::new(self.state.clone(), info, source))));
669 }
670
671 let index = self.index;
672 let Some((info, source)) = ready!(self.poll(waiter, |state| state.poll_frame_source(index))?) else {
673 return Poll::Ready(Ok(None));
674 };
675
676 self.index += 1;
677 self.stats.frames(1);
680 Poll::Ready(Ok(Some(
681 frame::Consumer::new(self.state.clone(), info, source).with_meter(self.stats.clone()),
682 )))
683 }
684
685 pub fn poll_read_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Frame>>> {
687 if let Some(frame) = self.prefetch.pop() {
689 self.index += 1;
690 return Poll::Ready(Ok(Some(frame)));
691 }
692
693 let index = self.index;
696 let prefetch = &mut self.prefetch;
697 let res = self.state.poll(waiter, |state| {
698 if index < state.offset {
699 return Poll::Ready(Err(Error::Lagged));
700 }
701 let local = (index - state.offset).min(state.frames.len());
706 prefetch.fill(state.frames.range(local..).cloned());
707 if prefetch.len > 0 {
708 state.charge.touch();
712 return Poll::Ready(Ok(()));
713 }
714 state.poll_terminal(index)
717 });
718
719 match ready!(res) {
720 Ok(Ok(())) => {}
721 Ok(Err(err)) => return Poll::Ready(Err(err)),
722 Err(state) => return Poll::Ready(Err(state.abort.clone().unwrap_or(Error::Dropped))),
723 }
724
725 let (frames, bytes) = self.prefetch.buffered();
728 self.stats.frames(frames);
729 self.stats.bytes(bytes);
730
731 Poll::Ready(Ok(self.prefetch.pop().inspect(|_| {
732 self.index += 1;
733 })))
734 }
735
736 pub async fn read_frame(&mut self) -> Result<Option<frame::Frame>> {
738 if let Some(frame) = self.prefetch.pop() {
740 self.index += 1;
741 return Ok(Some(frame));
742 }
743 kio::wait(|waiter| self.poll_read_frame(waiter)).await
744 }
745
746 pub fn poll_finished(&mut self, waiter: &kio::Waiter) -> Poll<Result<u64>> {
748 self.poll(waiter, |state| state.poll_finished())
749 }
750
751 pub async fn finished(&mut self) -> Result<u64> {
753 kio::wait(|waiter| self.poll_finished(waiter)).await
754 }
755}
756
757#[derive(Clone, Debug, Default)]
759#[non_exhaustive]
760pub struct Fetch {
761 pub priority: u8,
763}
764
765impl Fetch {
766 pub fn with_priority(mut self, priority: u8) -> Self {
768 self.priority = priority;
769 self
770 }
771}
772
773#[cfg(test)]
774mod test {
775 use super::*;
776 use bytes::Bytes;
777 use futures::FutureExt;
778
779 #[test]
780 fn basic_frame_reading() {
781 let mut producer = Info { sequence: 0 }.produce();
782 producer
783 .write_frame(Timestamp::ZERO, Bytes::from_static(b"frame0"))
784 .unwrap();
785 producer
786 .write_frame(Timestamp::ZERO, Bytes::from_static(b"frame1"))
787 .unwrap();
788 producer.finish().unwrap();
789
790 let mut consumer = producer.consume();
791 let f0 = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
792 assert_eq!(f0.size, 6);
793 let f1 = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
794 assert_eq!(f1.size, 6);
795 let end = consumer.next_frame().now_or_never().unwrap().unwrap();
796 assert!(end.is_none());
797 }
798
799 #[test]
800 fn read_frame_all_at_once() {
801 let mut producer = Info { sequence: 0 }.produce();
802 producer
803 .write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
804 .unwrap();
805 producer.finish().unwrap();
806
807 let mut consumer = producer.consume();
808 let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
809 assert_eq!(frame.payload, Bytes::from_static(b"hello"));
810 }
811
812 #[test]
813 fn read_frame_preserves_timestamp() {
814 let mut producer = Info { sequence: 0 }.produce();
815 let timestamp = Timestamp::from_micros(20_000).unwrap();
816 producer.write_frame(timestamp, Bytes::from_static(b"hello")).unwrap();
817 producer.finish().unwrap();
818
819 let mut consumer = producer.consume();
820 let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
821 assert_eq!(frame.timestamp.as_micros(), 20_000);
822 assert_eq!(frame.payload, Bytes::from_static(b"hello"));
823 }
824
825 #[test]
826 fn chunked_frame_reads_whole() {
827 let mut producer = Info { sequence: 0 }.produce();
828 {
829 let mut frame = producer
830 .create_frame(frame::Info {
831 size: 10,
832 timestamp: Timestamp::ZERO,
833 })
834 .unwrap();
835 frame.write(Bytes::from_static(b"hello")).unwrap();
836 frame.write(Bytes::from_static(b"world")).unwrap();
837 frame.finish().unwrap();
838 }
839 producer.finish().unwrap();
840
841 let mut consumer = producer.consume();
844 let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
845 assert_eq!(frame.payload, Bytes::from_static(b"helloworld"));
846 }
847
848 #[test]
849 fn chunked_frame_streams_partial() {
850 let mut producer = Info { sequence: 0 }.produce();
851 let mut consumer = producer.consume();
852
853 let mut frame = producer
854 .create_frame(frame::Info {
855 size: 6,
856 timestamp: Timestamp::ZERO,
857 })
858 .unwrap();
859 frame.write(Bytes::from_static(b"foo")).unwrap();
860
861 let mut f = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
863 let c1 = f.read_chunk().now_or_never().unwrap().unwrap();
864 assert_eq!(c1, Some(Bytes::from_static(b"foo")));
865 assert!(f.read_chunk().now_or_never().is_none());
866
867 frame.write(Bytes::from_static(b"bar")).unwrap();
868 frame.finish().unwrap();
869
870 let c2 = f.read_chunk().now_or_never().unwrap().unwrap();
871 assert_eq!(c2, Some(Bytes::from_static(b"bar")));
872 let c3 = f.read_chunk().now_or_never().unwrap().unwrap();
873 assert_eq!(c3, None);
874 }
875
876 #[test]
877 fn group_finish_returns_none() {
878 let mut producer = Info { sequence: 0 }.produce();
879 producer.finish().unwrap();
880
881 let mut consumer = producer.consume();
882 let end = consumer.next_frame().now_or_never().unwrap().unwrap();
883 assert!(end.is_none());
884 }
885
886 #[test]
887 fn abort_propagates() {
888 let producer = Info { sequence: 0 }.produce();
889 let mut consumer = producer.consume();
890 producer.abort(crate::Error::Cancel).unwrap();
891
892 let result = consumer.next_frame().now_or_never().unwrap();
893 assert!(matches!(result, Err(crate::Error::Cancel)));
894 }
895
896 #[test]
897 fn abort_clears_cached_frames() {
898 let mut producer = Info { sequence: 0 }.produce();
899 producer
900 .write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
901 .unwrap();
902
903 let _consumer = producer.consume();
905 assert_eq!(producer.state.read().frames.len(), 1);
906
907 producer.clone().abort(crate::Error::Cancel).unwrap();
908
909 let state = producer.state.read();
910 assert!(state.frames.is_empty(), "cached frames should be dropped on abort");
911 assert_eq!(state.cache, 0);
912 }
913
914 #[test]
915 fn drop_unfinished_clears_cached_frames() {
916 let producer = Info { sequence: 0 }.produce();
917 let mut writer = producer.clone();
918 writer
919 .write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
920 .unwrap();
921
922 let mut consumer = producer.consume();
924 assert_eq!(producer.state.read().frames.len(), 1);
925
926 drop(writer);
928 drop(producer);
929
930 let result = consumer.next_frame().now_or_never().unwrap();
931 assert!(matches!(result, Err(crate::Error::Dropped)));
932 }
933
934 #[test]
935 fn drop_finished_keeps_cached_frames() {
936 let mut producer = Info { sequence: 0 }.produce();
937 producer
938 .write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
939 .unwrap();
940 producer.finish().unwrap();
941
942 let mut consumer = producer.consume();
943 drop(producer);
944
945 let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
947 assert_eq!(frame.payload, Bytes::from_static(b"data"));
948 }
949
950 #[tokio::test]
951 async fn pending_then_ready() {
952 let mut producer = Info { sequence: 0 }.produce();
953 let mut consumer = producer.consume();
954
955 assert!(consumer.next_frame().now_or_never().is_none());
957
958 producer
959 .write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
960 .unwrap();
961 producer.finish().unwrap();
962
963 let frame = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
964 assert_eq!(frame.size, 4);
965 }
966
967 #[test]
968 fn eviction_drops_old_frames() {
969 let mut producer = Info { sequence: 0 }.produce();
970
971 let big = Bytes::from(vec![0u8; MAX_GROUP_CACHE as usize]);
973 producer.write_frame(Timestamp::ZERO, big.clone()).unwrap();
974 producer.write_frame(Timestamp::ZERO, big).unwrap();
975
976 let state = producer.state.read();
978 assert_eq!(state.offset, 1);
979 assert_eq!(state.frames.len(), 1);
980 assert_eq!(state.frames[0].payload.len(), MAX_GROUP_CACHE as usize);
981 }
982
983 #[test]
984 fn next_frame_returns_cache_full_on_tombstone() {
985 let mut producer = Info { sequence: 0 }.produce();
986
987 let big = Bytes::from(vec![0u8; MAX_GROUP_CACHE as usize]);
988 producer.write_frame(Timestamp::ZERO, big.clone()).unwrap();
989 producer.write_frame(Timestamp::ZERO, big).unwrap();
990
991 let mut consumer = producer.consume();
992 let result = consumer.next_frame().now_or_never().unwrap();
994 assert!(matches!(result, Err(crate::Error::Lagged)));
995 }
996
997 #[test]
998 fn no_eviction_under_budget() {
999 let mut producer = Info { sequence: 0 }.produce();
1000 for _ in 0..100_000 {
1002 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"x")).unwrap();
1003 }
1004 producer.finish().unwrap();
1005
1006 let state = producer.state.read();
1007 assert_eq!(state.offset, 0);
1008 assert_eq!(state.frames.len(), 100_000);
1009 }
1010
1011 #[test]
1012 fn clone_consumer_independent() {
1013 let mut producer = Info { sequence: 0 }.produce();
1014 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"a")).unwrap();
1015
1016 let mut c1 = producer.consume();
1017 let _ = c1.next_frame().now_or_never().unwrap().unwrap().unwrap();
1019
1020 let mut c2 = c1.clone();
1022
1023 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"b")).unwrap();
1024 producer.finish().unwrap();
1025
1026 let f = c2.next_frame().now_or_never().unwrap().unwrap().unwrap();
1028 assert_eq!(f.size, 1); let end = c2.next_frame().now_or_never().unwrap().unwrap();
1031 assert!(end.is_none());
1032 }
1033
1034 #[test]
1037 fn read_frame_crosses_prefetch_batches() {
1038 let n = Prefetch::CAP * 3 + 5;
1039 let mut producer = Info { sequence: 0 }.produce();
1040 for i in 0..n {
1041 producer
1042 .write_frame(Timestamp::ZERO, Bytes::from(vec![i as u8; 4]))
1043 .unwrap();
1044 }
1045 producer.finish().unwrap();
1046
1047 let mut consumer = producer.consume();
1048 for i in 0..n {
1049 let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
1050 assert_eq!(frame.payload, Bytes::from(vec![i as u8; 4]));
1051 }
1052 assert!(consumer.read_frame().now_or_never().unwrap().unwrap().is_none());
1053 }
1054
1055 #[test]
1059 fn abort_after_finish_keeps_the_clean_end_for_a_drained_reader() {
1060 let mut producer = Info { sequence: 0 }.produce();
1061 producer
1062 .write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
1063 .unwrap();
1064 producer.finish().unwrap();
1065
1066 let mut drained = producer.consume();
1067 let mut behind = producer.consume();
1068 let frame = drained.read_frame().now_or_never().unwrap().unwrap().unwrap();
1069 assert_eq!(frame.payload, Bytes::from_static(b"hello"));
1070
1071 producer.abort(Error::Old).unwrap();
1072
1073 assert!(drained.read_frame().now_or_never().unwrap().unwrap().is_none());
1075 assert!(drained.next_frame().now_or_never().unwrap().unwrap().is_none());
1076
1077 assert!(matches!(behind.read_frame().now_or_never().unwrap(), Err(Error::Old)));
1079 }
1080
1081 #[test]
1084 fn finished_survives_a_later_abort() {
1085 let mut producer = Info { sequence: 0 }.produce();
1086 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"a")).unwrap();
1087 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"b")).unwrap();
1088 producer.finish().unwrap();
1089
1090 let mut consumer = producer.consume();
1091 producer.abort(Error::Old).unwrap();
1092
1093 assert_eq!(consumer.finished().now_or_never().unwrap().unwrap(), 2);
1094 }
1095
1096 #[test]
1098 fn interleave_read_and_next_frame() {
1099 let mut producer = Info { sequence: 0 }.produce();
1100 for i in 0..5u8 {
1101 producer.write_frame(Timestamp::ZERO, Bytes::from(vec![i; 1])).unwrap();
1102 }
1103 producer.finish().unwrap();
1104
1105 let mut consumer = producer.consume();
1106 let f0 = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
1108 assert_eq!(f0.payload, Bytes::from(vec![0u8; 1]));
1109
1110 for i in 1..5u8 {
1112 let mut f = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
1113 let data = f.read_all().now_or_never().unwrap().unwrap();
1114 assert_eq!(data, Bytes::from(vec![i; 1]));
1115 }
1116 assert!(consumer.next_frame().now_or_never().unwrap().unwrap().is_none());
1117 }
1118
1119 #[test]
1122 fn read_frame_past_cleared_frames_does_not_panic() {
1123 let mut producer = Info { sequence: 0 }.produce();
1124 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"a")).unwrap();
1125 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"b")).unwrap();
1126
1127 let mut consumer = producer.consume();
1128 consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
1129 consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
1130
1131 producer.abort(Error::Cancel).unwrap();
1134
1135 let result = consumer.read_frame().now_or_never().unwrap();
1136 assert!(matches!(result, Err(Error::Cancel)), "expected Cancel, got {result:?}");
1137 }
1138
1139 #[test]
1142 fn drop_with_partial_batch() {
1143 let mut producer = Info { sequence: 0 }.produce();
1144 for _ in 0..Prefetch::CAP {
1145 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"x")).unwrap();
1146 }
1147 producer.finish().unwrap();
1148
1149 let mut consumer = producer.consume();
1150 let _ = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
1152 drop(consumer);
1153 }
1154
1155 #[test]
1158 fn create_frame_converts_mismatched_scale() {
1159 use crate::{Timescale, Timestamp};
1160
1161 let mut producer = Producer::new(
1162 Info { sequence: 0 },
1163 track::Info::default().with_timescale(Timescale::MICRO),
1164 );
1165 let frame = frame::Info {
1166 size: 3,
1167 timestamp: Timestamp::from_millis(1).unwrap(), };
1169 let writer = producer.create_frame(frame).unwrap();
1170 assert_eq!(writer.timestamp.scale(), Timescale::MICRO);
1171 assert_eq!(writer.timestamp.value(), 1000);
1172 }
1173
1174 #[tokio::test]
1176 async fn create_frame_converts_current_timestamp() {
1177 use crate::Timescale;
1178
1179 let mut producer = Producer::new(
1180 Info { sequence: 0 },
1181 track::Info::default().with_timescale(Timescale::MICRO),
1182 );
1183 let writer = producer
1184 .create_frame(frame::Info {
1185 size: 3,
1186 timestamp: Timestamp::now(),
1187 })
1188 .unwrap();
1189 assert_eq!(writer.timestamp.scale(), Timescale::MICRO);
1190 assert!(!writer.timestamp.is_zero(), "local clock should be non-zero");
1191 }
1192
1193 #[test]
1195 fn create_frame_rejects_oversized() {
1196 let mut producer = Info { sequence: 0 }.produce();
1197 let result = producer.create_frame(frame::Info {
1198 size: MAX_GROUP_CACHE + 1,
1199 timestamp: Timestamp::ZERO,
1200 });
1201 assert!(matches!(result, Err(Error::FrameTooLarge)));
1202 }
1203}