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(), Default::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 let info = frame::Info {
126 size: f.payload.len() as u64,
127 timestamp: f.timestamp,
128 };
129 return Poll::Ready(Ok(Some((info, frame::Source::Complete(f.payload.clone())))));
130 }
131 if local == self.frames.len()
132 && let Some(p) = &self.partial
133 {
134 let info = frame::Info {
135 size: p.buf.capacity() as u64,
136 timestamp: p.timestamp,
137 };
138 return Poll::Ready(Ok(Some((info, frame::Source::Partial(p.buf.clone())))));
139 }
140 ready!(self.poll_terminal(index))?;
141 Poll::Ready(Ok(None))
142 }
143
144 fn poll_terminal(&self, index: usize) -> Poll<Result<()>> {
151 match (self.fin, &self.abort) {
152 (Some(total), Some(err)) if index < total => Poll::Ready(Err(err.clone())),
153 (Some(_), _) => Poll::Ready(Ok(())),
154 (None, Some(err)) => Poll::Ready(Err(err.clone())),
155 (None, None) => Poll::Pending,
156 }
157 }
158
159 fn poll_finished(&self) -> Poll<Result<u64>> {
160 if let Some(total) = self.fin {
163 Poll::Ready(Ok(total as u64))
164 } else if let Some(err) = &self.abort {
165 Poll::Ready(Err(err.clone()))
166 } else {
167 Poll::Pending
168 }
169 }
170
171 fn evict(&mut self) {
173 while self.cache > MAX_GROUP_CACHE {
174 let Some(frame) = self.frames.pop_front() else {
175 break;
176 };
177 let size = frame.payload.len() as u64;
178 self.cache -= size;
179 self.charge.sub(size);
180 self.offset += 1;
181 }
182 }
183
184 fn release(&mut self) {
186 self.frames.clear();
187 self.partial = None;
188 self.cache = 0;
189 self.charge.clear();
190 }
191}
192
193fn modify(state: &kio::Producer<GroupState>) -> Result<kio::Mut<'_, GroupState>> {
194 state.write().map_err(|r| r.abort.clone().unwrap_or(Error::Dropped))
195}
196
197pub struct Producer {
203 state: kio::Producer<GroupState>,
205
206 info: Info,
209
210 track: track::Info,
215
216 cache: Arc<cache::Track>,
220
221 stats: stats::Meter,
224
225 alive: Arc<Alive>,
228}
229
230struct Alive {
238 info: Info,
239 state: kio::Producer<GroupState>,
240}
241
242impl Drop for Alive {
243 fn drop(&mut self) {
244 match self.state.write() {
250 Ok(mut state) => {
251 if state.fin.is_some() || state.abort.is_some() {
252 return;
253 }
254 tracing::warn!(
255 sequence = self.info.sequence,
256 "group::Producer dropped without finish() or abort()"
257 );
258 state.release();
259 }
260 Err(state) => {
261 if state.fin.is_some() || state.abort.is_some() {
262 return;
263 }
264 tracing::warn!(
265 sequence = self.info.sequence,
266 "group::Producer dropped without finish() or abort()"
267 );
268 }
269 }
270 }
271}
272
273impl std::ops::Deref for Producer {
274 type Target = Info;
275
276 fn deref(&self) -> &Self::Target {
277 &self.info
278 }
279}
280
281impl Producer {
282 pub(crate) fn new(info: Info, track: track::Info, cache: Arc<cache::Track>) -> Self {
293 let state = kio::Producer::<GroupState>::default();
294 state.write().ok().expect("a new group is open").charge = cache.charge();
295 let alive = Arc::new(Alive {
296 info,
297 state: state.clone(),
298 });
299 Self {
300 info,
301 state,
302 track,
303 cache,
304 stats: stats::Meter::default(),
305 alive,
306 }
307 }
308
309 pub(crate) fn with_meter(mut self, meter: stats::Meter) -> Self {
312 meter.group();
313 self.stats = meter;
314 self
315 }
316
317 pub(crate) fn info(&self) -> Info {
319 self.info
320 }
321
322 pub fn timescale(&self) -> Timescale {
324 self.track.timescale
325 }
326
327 pub fn write_frame<B: IntoBytes>(&mut self, timestamp: Timestamp, data: B) -> Result<()> {
335 let timestamp = timestamp
336 .convert(self.track.timescale)
337 .map_err(|_| Error::TimestampMismatch)?;
338 let payload = data.into_bytes();
339 if payload.len() as u64 > MAX_GROUP_CACHE {
340 return Err(Error::FrameTooLarge);
341 }
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 let size = payload.len() as u64;
349 state.cache += size;
350 state.charge.add(size);
351 state.frames.push_back(Frame { timestamp, payload });
352 state.evict();
353 drop(state);
354
355 self.cache.settle();
358
359 self.stats.frames(1);
361 self.stats.bytes(size);
362 Ok(())
363 }
364
365 pub fn create_frame(&mut self, frame: frame::Info) -> Result<frame::Producer<'_>> {
374 let timestamp = frame
375 .timestamp
376 .convert(self.track.timescale)
377 .map_err(|_| Error::TimestampMismatch)?;
378 if frame.size > MAX_GROUP_CACHE {
379 return Err(Error::FrameTooLarge);
380 }
381 let buf = FrameBuf::new(frame.size as usize);
382
383 let mut state = modify(&self.state)?;
384 if state.fin.is_some() {
385 return Err(Error::Closed);
386 }
387 debug_assert!(state.partial.is_none(), "a frame is already open");
388 state.cache += frame.size;
389 state.charge.add(frame.size);
390 state.partial = Some(Partial {
391 timestamp,
392 buf: buf.clone(),
393 });
394 state.evict();
395 drop(state);
396
397 self.cache.settle();
400
401 self.stats.frames(1);
404 let meter = self.stats.clone();
405
406 let info = frame::Info {
407 size: frame.size,
408 timestamp,
409 };
410 Ok(frame::Producer::new(self, buf, info).with_meter(meter))
411 }
412
413 pub(crate) fn frame_notify(&self) {
415 let _ = self.state.write();
417 }
418
419 pub(crate) fn frame_commit(&mut self, frame: Frame) -> Result<()> {
421 let mut state = modify(&self.state)?;
422 state.partial = None;
425 state.frames.push_back(frame);
426 Ok(())
427 }
428
429 pub(crate) fn frame_abort(&mut self, err: Error) {
432 let _ = self.clone().abort(err);
433 }
434
435 pub fn frame_count(&self) -> usize {
437 let state = self.state.read();
438 state.offset + state.frames.len() + state.partial.is_some() as usize
439 }
440
441 pub fn finish(&mut self) -> Result<()> {
446 let mut state = modify(&self.state)?;
447 state.fin = Some(state.offset + state.frames.len());
448 Ok(())
449 }
450
451 pub fn abort(self, err: Error) -> Result<()> {
457 let mut guard = modify(&self.state)?;
458 guard.abort = Some(err);
459 guard.release();
460 guard.close();
461 Ok(())
462 }
463
464 pub(crate) fn is_aborted(&self) -> bool {
467 self.state.read().abort.is_some()
468 }
469
470 pub(crate) fn cache_size(&self) -> u64 {
473 self.state.read().charge.size()
474 }
475
476 pub(crate) fn cache_accessed(&self) -> u64 {
479 self.state.read().charge.accessed()
480 }
481
482 pub(crate) fn cache_demote(&self) {
485 if let Ok(mut state) = self.state.write() {
486 state.charge.demote();
487 }
488 }
489
490 pub(crate) fn cache_refresh(&self) {
494 if let Ok(mut state) = self.state.write() {
495 state.charge.refresh();
496 }
497 }
498
499 pub fn consume(&self) -> Consumer {
501 Consumer {
502 info: self.info,
503 state: self.state.consume(),
504 track: self.track.clone(),
505 index: 0,
506 prefetch: Prefetch::default(),
507 stats: stats::Meter::default(),
510 }
511 }
512
513 pub async fn closed(&self) -> Error {
515 kio::wait(|waiter| self.poll_closed(waiter)).await
516 }
517
518 pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<Error> {
520 self.state.poll_closed(waiter).map(|()| self.abort_reason())
521 }
522
523 pub async fn unused(&self) -> Result<()> {
525 self.state.unused().await.map_err(|_| self.abort_reason())
526 }
527
528 fn abort_reason(&self) -> Error {
530 self.state.read().abort.clone().unwrap_or(Error::Dropped)
531 }
532}
533
534impl Clone for Producer {
535 fn clone(&self) -> Self {
536 Self {
537 info: self.info,
538 state: self.state.clone(),
539 track: self.track.clone(),
540 cache: self.cache.clone(),
541 stats: self.stats.clone(),
542 alive: self.alive.clone(),
543 }
544 }
545}
546
547struct Prefetch {
555 frames: [MaybeUninit<Frame>; Self::CAP],
557 pos: usize,
558 len: usize,
559}
560
561impl Prefetch {
562 const CAP: usize = 8;
563
564 fn pop(&mut self) -> Option<Frame> {
566 if self.pos == self.len {
567 return None;
568 }
569 let frame = unsafe { self.frames[self.pos].assume_init_read() };
571 self.pos += 1;
572 Some(frame)
573 }
574
575 fn fill(&mut self, frames: impl Iterator<Item = Frame>) {
577 debug_assert_eq!(self.pos, self.len, "fill on a non-empty batch would leak frames");
578 self.pos = 0;
579 self.len = 0;
580 for frame in frames.take(Self::CAP) {
581 self.frames[self.len].write(frame);
582 self.len += 1;
583 }
584 }
585
586 fn buffered(&self) -> (u64, u64) {
589 let mut bytes = 0u64;
590 for slot in &self.frames[self.pos..self.len] {
591 bytes += unsafe { slot.assume_init_ref() }.payload.len() as u64;
593 }
594 ((self.len - self.pos) as u64, bytes)
595 }
596}
597
598impl Default for Prefetch {
599 fn default() -> Self {
600 Self {
601 frames: [const { MaybeUninit::uninit() }; Self::CAP],
602 pos: 0,
603 len: 0,
604 }
605 }
606}
607
608impl Drop for Prefetch {
609 fn drop(&mut self) {
610 for slot in &mut self.frames[self.pos..self.len] {
611 unsafe { slot.assume_init_drop() };
613 }
614 }
615}
616
617pub struct Consumer {
619 state: kio::Consumer<GroupState>,
621
622 info: Info,
624
625 track: track::Info,
628
629 index: usize,
632
633 prefetch: Prefetch,
635
636 stats: stats::Meter,
639}
640
641impl Clone for Consumer {
642 fn clone(&self) -> Self {
643 Self {
646 state: self.state.clone(),
647 info: self.info,
648 track: self.track.clone(),
649 index: self.index,
650 prefetch: Prefetch::default(),
651 stats: self.stats.clone(),
654 }
655 }
656}
657
658impl std::ops::Deref for Consumer {
659 type Target = Info;
660
661 fn deref(&self) -> &Self::Target {
662 &self.info
663 }
664}
665
666impl Consumer {
667 pub(crate) fn with_meter(mut self, meter: stats::Meter) -> Self {
670 meter.group();
671 self.stats = meter;
672 self
673 }
674
675 pub(crate) fn is_aborted(&self) -> bool {
678 self.state.read().abort.is_some()
679 }
680
681 pub(crate) fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<()> {
685 self.state.poll_closed(waiter)
686 }
687
688 pub fn timescale(&self) -> Timescale {
690 self.track.timescale
691 }
692
693 fn poll<F, R>(&self, waiter: &kio::Waiter, f: F) -> Poll<Result<R>>
695 where
696 F: Fn(&kio::Ref<'_, GroupState>) -> Poll<Result<R>>,
697 {
698 Poll::Ready(match ready!(self.state.poll(waiter, f)) {
699 Ok(res) => res,
700 Err(state) => Err(state.abort.clone().unwrap_or(Error::Dropped)),
702 })
703 }
704
705 pub async fn next_frame(&mut self) -> Result<Option<frame::Consumer>> {
707 kio::wait(|waiter| self.poll_next_frame(waiter)).await
708 }
709
710 pub fn poll_next_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Consumer>>> {
714 if let Some(frame) = self.prefetch.pop() {
718 self.index += 1;
719 let info = frame::Info {
720 size: frame.payload.len() as u64,
721 timestamp: frame.timestamp,
722 };
723 let source = frame::Source::Complete(frame.payload);
724 return Poll::Ready(Ok(Some(frame::Consumer::new(self.state.clone(), info, source))));
725 }
726
727 let index = self.index;
728 let Some((info, source)) = ready!(self.poll(waiter, |state| state.poll_frame_source(index))?) else {
729 return Poll::Ready(Ok(None));
730 };
731
732 self.index += 1;
733 self.stats.frames(1);
736 Poll::Ready(Ok(Some(
737 frame::Consumer::new(self.state.clone(), info, source).with_meter(self.stats.clone()),
738 )))
739 }
740
741 pub fn poll_read_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Frame>>> {
743 if let Some(frame) = self.prefetch.pop() {
745 self.index += 1;
746 return Poll::Ready(Ok(Some(frame)));
747 }
748
749 let index = self.index;
752 let prefetch = &mut self.prefetch;
753 let res = self.state.poll(waiter, |state| {
754 if index < state.offset {
755 return Poll::Ready(Err(Error::Lagged));
756 }
757 let local = (index - state.offset).min(state.frames.len());
762 prefetch.fill(state.frames.range(local..).cloned());
763 if prefetch.len > 0 {
764 return Poll::Ready(Ok(()));
765 }
766 state.poll_terminal(index)
769 });
770
771 match ready!(res) {
772 Ok(Ok(())) => {}
773 Ok(Err(err)) => return Poll::Ready(Err(err)),
774 Err(state) => return Poll::Ready(Err(state.abort.clone().unwrap_or(Error::Dropped))),
775 }
776
777 let (frames, bytes) = self.prefetch.buffered();
780 self.stats.frames(frames);
781 self.stats.bytes(bytes);
782
783 Poll::Ready(Ok(self.prefetch.pop().inspect(|_| {
784 self.index += 1;
785 })))
786 }
787
788 pub async fn read_frame(&mut self) -> Result<Option<frame::Frame>> {
790 if let Some(frame) = self.prefetch.pop() {
792 self.index += 1;
793 return Ok(Some(frame));
794 }
795 kio::wait(|waiter| self.poll_read_frame(waiter)).await
796 }
797
798 pub fn poll_finished(&mut self, waiter: &kio::Waiter) -> Poll<Result<u64>> {
800 self.poll(waiter, |state| state.poll_finished())
801 }
802
803 pub async fn finished(&mut self) -> Result<u64> {
805 kio::wait(|waiter| self.poll_finished(waiter)).await
806 }
807}
808
809#[derive(Clone, Debug, Default)]
811#[non_exhaustive]
812pub struct Fetch {
813 pub priority: u8,
815}
816
817impl Fetch {
818 pub fn with_priority(mut self, priority: u8) -> Self {
820 self.priority = priority;
821 self
822 }
823}
824
825#[cfg(test)]
826mod test {
827 use super::*;
828 use crate::model::test_tracing::count_drop_warnings;
829 use bytes::Bytes;
830 use futures::FutureExt;
831
832 #[test]
833 fn basic_frame_reading() {
834 let mut producer = Info { sequence: 0 }.produce();
835 producer
836 .write_frame(Timestamp::ZERO, Bytes::from_static(b"frame0"))
837 .unwrap();
838 producer
839 .write_frame(Timestamp::ZERO, Bytes::from_static(b"frame1"))
840 .unwrap();
841 producer.finish().unwrap();
842
843 let mut consumer = producer.consume();
844 let f0 = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
845 assert_eq!(f0.size, 6);
846 let f1 = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
847 assert_eq!(f1.size, 6);
848 let end = consumer.next_frame().now_or_never().unwrap().unwrap();
849 assert!(end.is_none());
850 }
851
852 #[test]
853 fn read_frame_all_at_once() {
854 let mut producer = Info { sequence: 0 }.produce();
855 producer
856 .write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
857 .unwrap();
858 producer.finish().unwrap();
859
860 let mut consumer = producer.consume();
861 let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
862 assert_eq!(frame.payload, Bytes::from_static(b"hello"));
863 }
864
865 #[test]
866 fn read_frame_preserves_timestamp() {
867 let mut producer = Info { sequence: 0 }.produce();
868 let timestamp = Timestamp::from_micros(20_000).unwrap();
869 producer.write_frame(timestamp, Bytes::from_static(b"hello")).unwrap();
870 producer.finish().unwrap();
871
872 let mut consumer = producer.consume();
873 let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
874 assert_eq!(frame.timestamp.as_micros(), 20_000);
875 assert_eq!(frame.payload, Bytes::from_static(b"hello"));
876 }
877
878 #[test]
879 fn chunked_frame_reads_whole() {
880 let mut producer = Info { sequence: 0 }.produce();
881 {
882 let mut frame = producer
883 .create_frame(frame::Info {
884 size: 10,
885 timestamp: Timestamp::ZERO,
886 })
887 .unwrap();
888 frame.write(Bytes::from_static(b"hello")).unwrap();
889 frame.write(Bytes::from_static(b"world")).unwrap();
890 frame.finish().unwrap();
891 }
892 producer.finish().unwrap();
893
894 let mut consumer = producer.consume();
897 let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
898 assert_eq!(frame.payload, Bytes::from_static(b"helloworld"));
899 }
900
901 #[test]
902 fn chunked_frame_streams_partial() {
903 let mut producer = Info { sequence: 0 }.produce();
904 let mut consumer = producer.consume();
905
906 let mut frame = producer
907 .create_frame(frame::Info {
908 size: 6,
909 timestamp: Timestamp::ZERO,
910 })
911 .unwrap();
912 frame.write(Bytes::from_static(b"foo")).unwrap();
913
914 let mut f = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
916 let c1 = f.read_chunk().now_or_never().unwrap().unwrap();
917 assert_eq!(c1, Some(Bytes::from_static(b"foo")));
918 assert!(f.read_chunk().now_or_never().is_none());
919
920 frame.write(Bytes::from_static(b"bar")).unwrap();
921 frame.finish().unwrap();
922
923 let c2 = f.read_chunk().now_or_never().unwrap().unwrap();
924 assert_eq!(c2, Some(Bytes::from_static(b"bar")));
925 let c3 = f.read_chunk().now_or_never().unwrap().unwrap();
926 assert_eq!(c3, None);
927 }
928
929 #[test]
930 fn group_finish_returns_none() {
931 let mut producer = Info { sequence: 0 }.produce();
932 producer.finish().unwrap();
933
934 let mut consumer = producer.consume();
935 let end = consumer.next_frame().now_or_never().unwrap().unwrap();
936 assert!(end.is_none());
937 }
938
939 #[test]
940 fn abort_propagates() {
941 let producer = Info { sequence: 0 }.produce();
942 let mut consumer = producer.consume();
943 producer.abort(crate::Error::Cancel).unwrap();
944
945 let result = consumer.next_frame().now_or_never().unwrap();
946 assert!(matches!(result, Err(crate::Error::Cancel)));
947 }
948
949 #[test]
950 fn abort_clears_cached_frames() {
951 let mut producer = Info { sequence: 0 }.produce();
952 producer
953 .write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
954 .unwrap();
955
956 let _consumer = producer.consume();
958 assert_eq!(producer.state.read().frames.len(), 1);
959
960 producer.clone().abort(crate::Error::Cancel).unwrap();
961
962 let state = producer.state.read();
963 assert!(state.frames.is_empty(), "cached frames should be dropped on abort");
964 assert_eq!(state.cache, 0);
965 }
966
967 #[test]
968 fn drop_unfinished_clears_cached_frames() {
969 let producer = Info { sequence: 0 }.produce();
970 let mut writer = producer.clone();
971 writer
972 .write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
973 .unwrap();
974
975 let mut consumer = producer.consume();
977 assert_eq!(producer.state.read().frames.len(), 1);
978
979 drop(writer);
981 drop(producer);
982
983 let result = consumer.next_frame().now_or_never().unwrap();
984 assert!(matches!(result, Err(crate::Error::Dropped)));
985 }
986
987 #[test]
988 fn drop_after_abort_does_not_warn() {
989 let warns = count_drop_warnings("group::Producer dropped without finish", || {
990 let producer = Info { sequence: 0 }.produce();
991 let keep = producer.clone();
992 let mut writer = producer.clone();
993 writer
994 .write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
995 .unwrap();
996 let _consumer = producer.consume();
997 writer.abort(crate::Error::Cancel).unwrap();
998 drop(keep);
999 });
1000 assert_eq!(warns, 0, "abort-then-drop must not emit unfinished-producer WARN");
1001 }
1002
1003 #[test]
1004 fn drop_unfinished_warns() {
1005 let warns = count_drop_warnings("group::Producer dropped without finish", || {
1006 let producer = Info { sequence: 0 }.produce();
1007 let mut writer = producer.clone();
1008 writer
1009 .write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
1010 .unwrap();
1011 let _consumer = producer.consume();
1012 drop(writer);
1013 drop(producer);
1014 });
1015 assert!(warns >= 1, "unfinished drop must emit unfinished-producer WARN");
1016 }
1017
1018 #[test]
1019 fn drop_finished_keeps_cached_frames() {
1020 let mut producer = Info { sequence: 0 }.produce();
1021 producer
1022 .write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
1023 .unwrap();
1024 producer.finish().unwrap();
1025
1026 let mut consumer = producer.consume();
1027 drop(producer);
1028
1029 let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
1031 assert_eq!(frame.payload, Bytes::from_static(b"data"));
1032 }
1033
1034 #[tokio::test]
1035 async fn pending_then_ready() {
1036 let mut producer = Info { sequence: 0 }.produce();
1037 let mut consumer = producer.consume();
1038
1039 assert!(consumer.next_frame().now_or_never().is_none());
1041
1042 producer
1043 .write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
1044 .unwrap();
1045 producer.finish().unwrap();
1046
1047 let frame = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
1048 assert_eq!(frame.size, 4);
1049 }
1050
1051 #[test]
1052 fn eviction_drops_old_frames() {
1053 let mut producer = Info { sequence: 0 }.produce();
1054
1055 let big = Bytes::from(vec![0u8; MAX_GROUP_CACHE as usize]);
1057 producer.write_frame(Timestamp::ZERO, big.clone()).unwrap();
1058 producer.write_frame(Timestamp::ZERO, big).unwrap();
1059
1060 let state = producer.state.read();
1062 assert_eq!(state.offset, 1);
1063 assert_eq!(state.frames.len(), 1);
1064 assert_eq!(state.frames[0].payload.len(), MAX_GROUP_CACHE as usize);
1065 }
1066
1067 #[test]
1068 fn next_frame_returns_cache_full_on_tombstone() {
1069 let mut producer = Info { sequence: 0 }.produce();
1070
1071 let big = Bytes::from(vec![0u8; MAX_GROUP_CACHE as usize]);
1072 producer.write_frame(Timestamp::ZERO, big.clone()).unwrap();
1073 producer.write_frame(Timestamp::ZERO, big).unwrap();
1074
1075 let mut consumer = producer.consume();
1076 let result = consumer.next_frame().now_or_never().unwrap();
1078 assert!(matches!(result, Err(crate::Error::Lagged)));
1079 }
1080
1081 #[test]
1082 fn no_eviction_under_budget() {
1083 let mut producer = Info { sequence: 0 }.produce();
1084 for _ in 0..100_000 {
1086 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"x")).unwrap();
1087 }
1088 producer.finish().unwrap();
1089
1090 let state = producer.state.read();
1091 assert_eq!(state.offset, 0);
1092 assert_eq!(state.frames.len(), 100_000);
1093 }
1094
1095 #[test]
1096 fn clone_consumer_independent() {
1097 let mut producer = Info { sequence: 0 }.produce();
1098 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"a")).unwrap();
1099
1100 let mut c1 = producer.consume();
1101 let _ = c1.next_frame().now_or_never().unwrap().unwrap().unwrap();
1103
1104 let mut c2 = c1.clone();
1106
1107 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"b")).unwrap();
1108 producer.finish().unwrap();
1109
1110 let f = c2.next_frame().now_or_never().unwrap().unwrap().unwrap();
1112 assert_eq!(f.size, 1); let end = c2.next_frame().now_or_never().unwrap().unwrap();
1115 assert!(end.is_none());
1116 }
1117
1118 #[test]
1121 fn read_frame_crosses_prefetch_batches() {
1122 let n = Prefetch::CAP * 3 + 5;
1123 let mut producer = Info { sequence: 0 }.produce();
1124 for i in 0..n {
1125 producer
1126 .write_frame(Timestamp::ZERO, Bytes::from(vec![i as u8; 4]))
1127 .unwrap();
1128 }
1129 producer.finish().unwrap();
1130
1131 let mut consumer = producer.consume();
1132 for i in 0..n {
1133 let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
1134 assert_eq!(frame.payload, Bytes::from(vec![i as u8; 4]));
1135 }
1136 assert!(consumer.read_frame().now_or_never().unwrap().unwrap().is_none());
1137 }
1138
1139 #[test]
1143 fn abort_after_finish_keeps_the_clean_end_for_a_drained_reader() {
1144 let mut producer = Info { sequence: 0 }.produce();
1145 producer
1146 .write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
1147 .unwrap();
1148 producer.finish().unwrap();
1149
1150 let mut drained = producer.consume();
1151 let mut behind = producer.consume();
1152 let frame = drained.read_frame().now_or_never().unwrap().unwrap().unwrap();
1153 assert_eq!(frame.payload, Bytes::from_static(b"hello"));
1154
1155 producer.abort(Error::Old).unwrap();
1156
1157 assert!(drained.read_frame().now_or_never().unwrap().unwrap().is_none());
1159 assert!(drained.next_frame().now_or_never().unwrap().unwrap().is_none());
1160
1161 assert!(matches!(behind.read_frame().now_or_never().unwrap(), Err(Error::Old)));
1163 }
1164
1165 #[test]
1168 fn finished_survives_a_later_abort() {
1169 let mut producer = Info { sequence: 0 }.produce();
1170 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"a")).unwrap();
1171 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"b")).unwrap();
1172 producer.finish().unwrap();
1173
1174 let mut consumer = producer.consume();
1175 producer.abort(Error::Old).unwrap();
1176
1177 assert_eq!(consumer.finished().now_or_never().unwrap().unwrap(), 2);
1178 }
1179
1180 #[test]
1182 fn interleave_read_and_next_frame() {
1183 let mut producer = Info { sequence: 0 }.produce();
1184 for i in 0..5u8 {
1185 producer.write_frame(Timestamp::ZERO, Bytes::from(vec![i; 1])).unwrap();
1186 }
1187 producer.finish().unwrap();
1188
1189 let mut consumer = producer.consume();
1190 let f0 = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
1192 assert_eq!(f0.payload, Bytes::from(vec![0u8; 1]));
1193
1194 for i in 1..5u8 {
1196 let mut f = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
1197 let data = f.read_all().now_or_never().unwrap().unwrap();
1198 assert_eq!(data, Bytes::from(vec![i; 1]));
1199 }
1200 assert!(consumer.next_frame().now_or_never().unwrap().unwrap().is_none());
1201 }
1202
1203 #[test]
1206 fn read_frame_past_cleared_frames_does_not_panic() {
1207 let mut producer = Info { sequence: 0 }.produce();
1208 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"a")).unwrap();
1209 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"b")).unwrap();
1210
1211 let mut consumer = producer.consume();
1212 consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
1213 consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
1214
1215 producer.abort(Error::Cancel).unwrap();
1218
1219 let result = consumer.read_frame().now_or_never().unwrap();
1220 assert!(matches!(result, Err(Error::Cancel)), "expected Cancel, got {result:?}");
1221 }
1222
1223 #[test]
1226 fn drop_with_partial_batch() {
1227 let mut producer = Info { sequence: 0 }.produce();
1228 for _ in 0..Prefetch::CAP {
1229 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"x")).unwrap();
1230 }
1231 producer.finish().unwrap();
1232
1233 let mut consumer = producer.consume();
1234 let _ = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
1236 drop(consumer);
1237 }
1238
1239 #[test]
1242 fn create_frame_converts_mismatched_scale() {
1243 use crate::{Timescale, Timestamp};
1244
1245 let mut producer = Producer::new(
1246 Info { sequence: 0 },
1247 track::Info::default().with_timescale(Timescale::MICRO),
1248 Default::default(),
1249 );
1250 let frame = frame::Info {
1251 size: 3,
1252 timestamp: Timestamp::from_millis(1).unwrap(), };
1254 let writer = producer.create_frame(frame).unwrap();
1255 assert_eq!(writer.timestamp.scale(), Timescale::MICRO);
1256 assert_eq!(writer.timestamp.value(), 1000);
1257 }
1258
1259 #[tokio::test]
1261 async fn create_frame_converts_current_timestamp() {
1262 use crate::Timescale;
1263
1264 let mut producer = Producer::new(
1265 Info { sequence: 0 },
1266 track::Info::default().with_timescale(Timescale::MICRO),
1267 Default::default(),
1268 );
1269 let writer = producer
1270 .create_frame(frame::Info {
1271 size: 3,
1272 timestamp: Timestamp::now(),
1273 })
1274 .unwrap();
1275 assert_eq!(writer.timestamp.scale(), Timescale::MICRO);
1276 assert!(!writer.timestamp.is_zero(), "local clock should be non-zero");
1277 }
1278
1279 #[test]
1281 fn create_frame_rejects_oversized() {
1282 let mut producer = Info { sequence: 0 }.produce();
1283 let result = producer.create_frame(frame::Info {
1284 size: MAX_GROUP_CACHE + 1,
1285 timestamp: Timestamp::ZERO,
1286 });
1287 assert!(matches!(result, Err(Error::FrameTooLarge)));
1288 }
1289}