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
21pub(super) const 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 self.charge.refresh();
128 let info = frame::Info {
129 size: f.payload.len() as u64,
130 timestamp: f.timestamp,
131 };
132 return Poll::Ready(Ok(Some((info, frame::Source::Complete(f.payload.clone())))));
133 }
134 if local == self.frames.len()
135 && let Some(p) = &self.partial
136 {
137 self.charge.refresh();
138 let info = frame::Info {
139 size: p.buf.capacity() as u64,
140 timestamp: p.timestamp,
141 };
142 return Poll::Ready(Ok(Some((info, frame::Source::Partial(p.buf.clone())))));
143 }
144 ready!(self.poll_terminal(index))?;
145 Poll::Ready(Ok(None))
146 }
147
148 fn poll_terminal(&self, index: usize) -> Poll<Result<()>> {
155 match (self.fin, &self.abort) {
156 (Some(total), Some(err)) if index < total => Poll::Ready(Err(err.clone())),
157 (Some(_), _) => Poll::Ready(Ok(())),
158 (None, Some(err)) => Poll::Ready(Err(err.clone())),
159 (None, None) => Poll::Pending,
160 }
161 }
162
163 fn poll_finished(&self) -> Poll<Result<u64>> {
164 if let Some(total) = self.fin {
167 Poll::Ready(Ok(total as u64))
168 } else if let Some(err) = &self.abort {
169 Poll::Ready(Err(err.clone()))
170 } else {
171 Poll::Pending
172 }
173 }
174
175 fn evict(&mut self) {
177 while self.cache > MAX_GROUP_CACHE {
178 let Some(frame) = self.frames.pop_front() else {
179 break;
180 };
181 let size = frame.payload.len() as u64;
182 self.cache -= size;
183 self.charge.sub(size);
184 self.offset += 1;
185 }
186 }
187
188 fn release(&mut self) {
190 self.frames.clear();
191 self.partial = None;
192 self.cache = 0;
193 self.charge.clear();
194 }
195}
196
197fn modify(state: &kio::Producer<GroupState>) -> Result<kio::Mut<'_, GroupState>> {
198 state.write().map_err(|r| r.abort.clone().unwrap_or(Error::Dropped))
199}
200
201pub struct Producer {
207 state: kio::Producer<GroupState>,
209
210 info: Info,
213
214 track: track::Info,
219
220 cache: Arc<cache::Track>,
224
225 stats: stats::Meter,
228
229 alive: Arc<Alive>,
232}
233
234struct Alive {
242 info: Info,
243 state: kio::Producer<GroupState>,
244}
245
246impl Drop for Alive {
247 fn drop(&mut self) {
248 match self.state.write() {
254 Ok(mut state) => {
255 if state.fin.is_some() || state.abort.is_some() {
256 return;
257 }
258 tracing::warn!(
259 sequence = self.info.sequence,
260 "group::Producer dropped without finish() or abort()"
261 );
262 state.release();
263 }
264 Err(state) => {
265 if state.fin.is_some() || state.abort.is_some() {
266 return;
267 }
268 tracing::warn!(
269 sequence = self.info.sequence,
270 "group::Producer dropped without finish() or abort()"
271 );
272 }
273 }
274 }
275}
276
277impl std::ops::Deref for Producer {
278 type Target = Info;
279
280 fn deref(&self) -> &Self::Target {
281 &self.info
282 }
283}
284
285impl Producer {
286 pub(crate) fn new(info: Info, track: track::Info, cache: Arc<cache::Track>) -> Self {
297 let state = kio::Producer::<GroupState>::default();
298 state.write().ok().expect("a new group is open").charge = cache.charge();
299 let alive = Arc::new(Alive {
300 info,
301 state: state.clone(),
302 });
303 Self {
304 info,
305 state,
306 track,
307 cache,
308 stats: stats::Meter::default(),
309 alive,
310 }
311 }
312
313 pub(crate) fn with_meter(mut self, meter: stats::Meter) -> Self {
316 meter.group();
317 self.stats = meter;
318 self
319 }
320
321 pub(crate) fn info(&self) -> Info {
323 self.info
324 }
325
326 pub fn timescale(&self) -> Timescale {
328 self.track.timescale
329 }
330
331 pub fn write_frame<B: IntoBytes>(&mut self, timestamp: Timestamp, data: B) -> Result<()> {
339 let timestamp = timestamp
340 .convert(self.track.timescale)
341 .map_err(|_| Error::TimestampMismatch)?;
342 let payload = data.into_bytes();
343 if payload.len() as u64 > MAX_GROUP_CACHE {
344 return Err(Error::FrameTooLarge);
345 }
346
347 let mut state = modify(&self.state)?;
348 if state.fin.is_some() {
349 return Err(Error::Closed);
350 }
351 debug_assert!(state.partial.is_none(), "a frame is already open");
352 let size = payload.len() as u64;
353 state.cache += size;
354 state.charge.add(size);
355 state.frames.push_back(Frame { timestamp, payload });
356 state.evict();
357 drop(state);
358
359 self.cache.settle();
362
363 self.stats.frames(1);
365 self.stats.bytes(size);
366 Ok(())
367 }
368
369 pub fn create_frame(&mut self, frame: frame::Info) -> Result<frame::Producer<'_>> {
378 let timestamp = frame
379 .timestamp
380 .convert(self.track.timescale)
381 .map_err(|_| Error::TimestampMismatch)?;
382 if frame.size > MAX_GROUP_CACHE {
383 return Err(Error::FrameTooLarge);
384 }
385 let buf = FrameBuf::new(frame.size as usize);
386
387 let mut state = modify(&self.state)?;
388 if state.fin.is_some() {
389 return Err(Error::Closed);
390 }
391 debug_assert!(state.partial.is_none(), "a frame is already open");
392 state.cache += frame.size;
393 state.charge.add(frame.size);
394 state.partial = Some(Partial {
395 timestamp,
396 buf: buf.clone(),
397 });
398 state.evict();
399 drop(state);
400
401 self.cache.settle();
404
405 self.stats.frames(1);
408 let meter = self.stats.clone();
409
410 let info = frame::Info {
411 size: frame.size,
412 timestamp,
413 };
414 Ok(frame::Producer::new(self, buf, info).with_meter(meter))
415 }
416
417 pub(crate) fn frame_notify(&self) {
419 if let Ok(mut state) = self.state.write() {
426 state.charge.record_write();
427 }
428 }
429
430 pub(crate) fn frame_commit(&mut self, frame: Frame) -> Result<()> {
432 let mut state = modify(&self.state)?;
433 state.partial = None;
436 state.frames.push_back(frame);
437 Ok(())
438 }
439
440 pub(crate) fn frame_abort(&mut self, err: Error) {
443 let _ = self.clone().abort(err);
444 }
445
446 pub fn frame_count(&self) -> usize {
448 let state = self.state.read();
449 state.offset + state.frames.len() + state.partial.is_some() as usize
450 }
451
452 pub fn finish(&mut self) -> Result<()> {
457 let mut state = modify(&self.state)?;
458 state.fin = Some(state.offset + state.frames.len());
459 Ok(())
460 }
461
462 pub fn abort(self, err: Error) -> Result<()> {
468 let mut guard = modify(&self.state)?;
469 guard.abort = Some(err);
470 guard.release();
471 guard.close();
472 Ok(())
473 }
474
475 pub(crate) fn is_aborted(&self) -> bool {
478 self.state.read().abort.is_some()
479 }
480
481 pub(crate) fn cache_size(&self) -> u64 {
484 self.state.read().charge.size()
485 }
486
487 pub(crate) fn cache_accessed(&self) -> u64 {
490 self.state.read().charge.accessed()
491 }
492
493 pub(crate) fn cache_demote(&self) {
496 if let Ok(mut state) = self.state.write() {
497 state.charge.demote();
498 }
499 }
500
501 pub(crate) fn cache_refresh(&self) {
507 self.state.read().charge.refresh();
508 }
509
510 pub fn consume(&self) -> Consumer {
512 Consumer {
513 info: self.info,
514 state: self.state.consume(),
515 track: self.track.clone(),
516 index: 0,
517 prefetch: Prefetch::default(),
518 last_refresh: web_async::time::Instant::now(),
519 stats: stats::Meter::default(),
522 }
523 }
524
525 pub async fn closed(&self) -> Error {
527 kio::wait(|waiter| self.poll_closed(waiter)).await
528 }
529
530 pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<Error> {
532 self.state.poll_closed(waiter).map(|()| self.abort_reason())
533 }
534
535 pub async fn unused(&self) -> Result<()> {
537 self.state.unused().await.map_err(|_| self.abort_reason())
538 }
539
540 fn abort_reason(&self) -> Error {
542 self.state.read().abort.clone().unwrap_or(Error::Dropped)
543 }
544}
545
546impl Clone for Producer {
547 fn clone(&self) -> Self {
548 Self {
549 info: self.info,
550 state: self.state.clone(),
551 track: self.track.clone(),
552 cache: self.cache.clone(),
553 stats: self.stats.clone(),
554 alive: self.alive.clone(),
555 }
556 }
557}
558
559struct Prefetch {
567 frames: [MaybeUninit<Frame>; Self::CAP],
569 pos: usize,
570 len: usize,
571}
572
573impl Prefetch {
574 const CAP: usize = 8;
575
576 fn pop(&mut self) -> Option<Frame> {
578 if self.pos == self.len {
579 return None;
580 }
581 let frame = unsafe { self.frames[self.pos].assume_init_read() };
583 self.pos += 1;
584 Some(frame)
585 }
586
587 fn fill(&mut self, frames: impl Iterator<Item = Frame>) {
589 debug_assert_eq!(self.pos, self.len, "fill on a non-empty batch would leak frames");
590 self.pos = 0;
591 self.len = 0;
592 for frame in frames.take(Self::CAP) {
593 self.frames[self.len].write(frame);
594 self.len += 1;
595 }
596 }
597
598 fn buffered(&self) -> (u64, u64) {
601 let mut bytes = 0u64;
602 for slot in &self.frames[self.pos..self.len] {
603 bytes += unsafe { slot.assume_init_ref() }.payload.len() as u64;
605 }
606 ((self.len - self.pos) as u64, bytes)
607 }
608}
609
610impl Default for Prefetch {
611 fn default() -> Self {
612 Self {
613 frames: [const { MaybeUninit::uninit() }; Self::CAP],
614 pos: 0,
615 len: 0,
616 }
617 }
618}
619
620impl Drop for Prefetch {
621 fn drop(&mut self) {
622 for slot in &mut self.frames[self.pos..self.len] {
623 unsafe { slot.assume_init_drop() };
625 }
626 }
627}
628
629pub struct Consumer {
631 state: kio::Consumer<GroupState>,
633
634 info: Info,
636
637 track: track::Info,
640
641 index: usize,
644
645 prefetch: Prefetch,
647
648 last_refresh: web_async::time::Instant,
653
654 stats: stats::Meter,
657}
658
659impl Clone for Consumer {
660 fn clone(&self) -> Self {
661 Self {
664 state: self.state.clone(),
665 info: self.info,
666 track: self.track.clone(),
667 index: self.index,
668 prefetch: Prefetch::default(),
669 last_refresh: self.last_refresh,
670 stats: self.stats.clone(),
673 }
674 }
675}
676
677impl std::ops::Deref for Consumer {
678 type Target = Info;
679
680 fn deref(&self) -> &Self::Target {
681 &self.info
682 }
683}
684
685impl Consumer {
686 pub(crate) fn with_meter(mut self, meter: stats::Meter) -> Self {
689 meter.group();
690 self.stats = meter;
691 self
692 }
693
694 pub(crate) fn is_aborted(&self) -> bool {
697 self.state.read().abort.is_some()
698 }
699
700 pub(crate) fn cache_refresh(&self) {
703 self.state.read().charge.refresh();
704 }
705
706 pub(crate) fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<()> {
710 self.state.poll_closed(waiter)
711 }
712
713 pub fn timescale(&self) -> Timescale {
715 self.track.timescale
716 }
717
718 fn refresh_if_stale(&mut self) {
724 if self.last_refresh.elapsed() * 2 < self.track.latency_max {
725 return;
726 }
727 self.state.read().charge.refresh();
728 self.last_refresh = web_async::time::Instant::now();
729 }
730
731 fn poll<F, R>(&self, waiter: &kio::Waiter, f: F) -> Poll<Result<R>>
733 where
734 F: Fn(&kio::Ref<'_, GroupState>) -> Poll<Result<R>>,
735 {
736 Poll::Ready(match ready!(self.state.poll(waiter, f)) {
737 Ok(res) => res,
738 Err(state) => Err(state.abort.clone().unwrap_or(Error::Dropped)),
740 })
741 }
742
743 pub async fn next_frame(&mut self) -> Result<Option<frame::Consumer>> {
745 kio::wait(|waiter| self.poll_next_frame(waiter)).await
746 }
747
748 pub fn poll_next_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Consumer>>> {
752 if let Some(frame) = self.prefetch.pop() {
756 self.refresh_if_stale();
757 self.index += 1;
758 let info = frame::Info {
759 size: frame.payload.len() as u64,
760 timestamp: frame.timestamp,
761 };
762 let source = frame::Source::Complete(frame.payload);
763 return Poll::Ready(Ok(Some(frame::Consumer::new(self.state.clone(), info, source))));
764 }
765
766 let index = self.index;
767 let Some((info, source)) = ready!(self.poll(waiter, |state| state.poll_frame_source(index))?) else {
768 return Poll::Ready(Ok(None));
769 };
770
771 self.index += 1;
772 self.stats.frames(1);
775 Poll::Ready(Ok(Some(
776 frame::Consumer::new(self.state.clone(), info, source).with_meter(self.stats.clone()),
777 )))
778 }
779
780 pub fn poll_read_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Frame>>> {
782 if let Some(frame) = self.prefetch.pop() {
784 self.refresh_if_stale();
785 self.index += 1;
786 return Poll::Ready(Ok(Some(frame)));
787 }
788
789 let index = self.index;
792 let prefetch = &mut self.prefetch;
793 let res = self.state.poll(waiter, |state| {
794 if index < state.offset {
795 return Poll::Ready(Err(Error::Lagged));
796 }
797 let local = (index - state.offset).min(state.frames.len());
802 prefetch.fill(state.frames.range(local..).cloned());
803 if prefetch.len > 0 {
804 state.charge.refresh();
807 return Poll::Ready(Ok(()));
808 }
809 state.poll_terminal(index)
812 });
813
814 match ready!(res) {
815 Ok(Ok(())) => {}
816 Ok(Err(err)) => return Poll::Ready(Err(err)),
817 Err(state) => return Poll::Ready(Err(state.abort.clone().unwrap_or(Error::Dropped))),
818 }
819
820 self.last_refresh = web_async::time::Instant::now();
823
824 let (frames, bytes) = self.prefetch.buffered();
827 self.stats.frames(frames);
828 self.stats.bytes(bytes);
829
830 Poll::Ready(Ok(self.prefetch.pop().inspect(|_| {
831 self.index += 1;
832 })))
833 }
834
835 pub async fn read_frame(&mut self) -> Result<Option<frame::Frame>> {
837 if let Some(frame) = self.prefetch.pop() {
839 self.refresh_if_stale();
840 self.index += 1;
841 return Ok(Some(frame));
842 }
843 kio::wait(|waiter| self.poll_read_frame(waiter)).await
844 }
845
846 pub fn poll_finished(&mut self, waiter: &kio::Waiter) -> Poll<Result<u64>> {
848 self.poll(waiter, |state| state.poll_finished())
849 }
850
851 pub async fn finished(&mut self) -> Result<u64> {
853 kio::wait(|waiter| self.poll_finished(waiter)).await
854 }
855}
856
857#[derive(Clone, Debug, Default)]
859#[non_exhaustive]
860pub struct Fetch {
861 pub priority: u8,
863}
864
865impl Fetch {
866 pub fn with_priority(mut self, priority: u8) -> Self {
868 self.priority = priority;
869 self
870 }
871}
872
873#[cfg(test)]
874mod test {
875 use super::*;
876 use crate::model::test_tracing::count_drop_warnings;
877 use bytes::Bytes;
878 use futures::FutureExt;
879
880 #[test]
881 fn basic_frame_reading() {
882 let mut producer = Info { sequence: 0 }.produce();
883 producer
884 .write_frame(Timestamp::ZERO, Bytes::from_static(b"frame0"))
885 .unwrap();
886 producer
887 .write_frame(Timestamp::ZERO, Bytes::from_static(b"frame1"))
888 .unwrap();
889 producer.finish().unwrap();
890
891 let mut consumer = producer.consume();
892 let f0 = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
893 assert_eq!(f0.size, 6);
894 let f1 = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
895 assert_eq!(f1.size, 6);
896 let end = consumer.next_frame().now_or_never().unwrap().unwrap();
897 assert!(end.is_none());
898 }
899
900 #[test]
901 fn read_frame_all_at_once() {
902 let mut producer = Info { sequence: 0 }.produce();
903 producer
904 .write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
905 .unwrap();
906 producer.finish().unwrap();
907
908 let mut consumer = producer.consume();
909 let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
910 assert_eq!(frame.payload, Bytes::from_static(b"hello"));
911 }
912
913 #[test]
914 fn read_frame_preserves_timestamp() {
915 let mut producer = Info { sequence: 0 }.produce();
916 let timestamp = Timestamp::from_micros(20_000).unwrap();
917 producer.write_frame(timestamp, Bytes::from_static(b"hello")).unwrap();
918 producer.finish().unwrap();
919
920 let mut consumer = producer.consume();
921 let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
922 assert_eq!(frame.timestamp.as_micros(), 20_000);
923 assert_eq!(frame.payload, Bytes::from_static(b"hello"));
924 }
925
926 #[test]
927 fn chunked_frame_reads_whole() {
928 let mut producer = Info { sequence: 0 }.produce();
929 {
930 let mut frame = producer
931 .create_frame(frame::Info {
932 size: 10,
933 timestamp: Timestamp::ZERO,
934 })
935 .unwrap();
936 frame.write(Bytes::from_static(b"hello")).unwrap();
937 frame.write(Bytes::from_static(b"world")).unwrap();
938 frame.finish().unwrap();
939 }
940 producer.finish().unwrap();
941
942 let mut consumer = producer.consume();
945 let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
946 assert_eq!(frame.payload, Bytes::from_static(b"helloworld"));
947 }
948
949 #[test]
950 fn chunked_frame_streams_partial() {
951 let mut producer = Info { sequence: 0 }.produce();
952 let mut consumer = producer.consume();
953
954 let mut frame = producer
955 .create_frame(frame::Info {
956 size: 6,
957 timestamp: Timestamp::ZERO,
958 })
959 .unwrap();
960 frame.write(Bytes::from_static(b"foo")).unwrap();
961
962 let mut f = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
964 let c1 = f.read_chunk().now_or_never().unwrap().unwrap();
965 assert_eq!(c1, Some(Bytes::from_static(b"foo")));
966 assert!(f.read_chunk().now_or_never().is_none());
967
968 frame.write(Bytes::from_static(b"bar")).unwrap();
969 frame.finish().unwrap();
970
971 let c2 = f.read_chunk().now_or_never().unwrap().unwrap();
972 assert_eq!(c2, Some(Bytes::from_static(b"bar")));
973 let c3 = f.read_chunk().now_or_never().unwrap().unwrap();
974 assert_eq!(c3, None);
975 }
976
977 #[test]
978 fn group_finish_returns_none() {
979 let mut producer = Info { sequence: 0 }.produce();
980 producer.finish().unwrap();
981
982 let mut consumer = producer.consume();
983 let end = consumer.next_frame().now_or_never().unwrap().unwrap();
984 assert!(end.is_none());
985 }
986
987 #[test]
988 fn abort_propagates() {
989 let producer = Info { sequence: 0 }.produce();
990 let mut consumer = producer.consume();
991 producer.abort(crate::Error::Cancel).unwrap();
992
993 let result = consumer.next_frame().now_or_never().unwrap();
994 assert!(matches!(result, Err(crate::Error::Cancel)));
995 }
996
997 #[test]
998 fn abort_clears_cached_frames() {
999 let mut producer = Info { sequence: 0 }.produce();
1000 producer
1001 .write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
1002 .unwrap();
1003
1004 let _consumer = producer.consume();
1006 assert_eq!(producer.state.read().frames.len(), 1);
1007
1008 producer.clone().abort(crate::Error::Cancel).unwrap();
1009
1010 let state = producer.state.read();
1011 assert!(state.frames.is_empty(), "cached frames should be dropped on abort");
1012 assert_eq!(state.cache, 0);
1013 }
1014
1015 #[test]
1016 fn drop_unfinished_clears_cached_frames() {
1017 let producer = Info { sequence: 0 }.produce();
1018 let mut writer = producer.clone();
1019 writer
1020 .write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
1021 .unwrap();
1022
1023 let mut consumer = producer.consume();
1025 assert_eq!(producer.state.read().frames.len(), 1);
1026
1027 drop(writer);
1029 drop(producer);
1030
1031 let result = consumer.next_frame().now_or_never().unwrap();
1032 assert!(matches!(result, Err(crate::Error::Dropped)));
1033 }
1034
1035 #[test]
1036 fn drop_after_abort_does_not_warn() {
1037 let warns = count_drop_warnings("group::Producer dropped without finish", || {
1038 let producer = Info { sequence: 0 }.produce();
1039 let keep = producer.clone();
1040 let mut writer = producer.clone();
1041 writer
1042 .write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
1043 .unwrap();
1044 let _consumer = producer.consume();
1045 writer.abort(crate::Error::Cancel).unwrap();
1046 drop(keep);
1047 });
1048 assert_eq!(warns, 0, "abort-then-drop must not emit unfinished-producer WARN");
1049 }
1050
1051 #[test]
1052 fn drop_unfinished_warns() {
1053 let warns = count_drop_warnings("group::Producer dropped without finish", || {
1054 let producer = Info { sequence: 0 }.produce();
1055 let mut writer = producer.clone();
1056 writer
1057 .write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
1058 .unwrap();
1059 let _consumer = producer.consume();
1060 drop(writer);
1061 drop(producer);
1062 });
1063 assert!(warns >= 1, "unfinished drop must emit unfinished-producer WARN");
1064 }
1065
1066 #[test]
1067 fn drop_finished_keeps_cached_frames() {
1068 let mut producer = Info { sequence: 0 }.produce();
1069 producer
1070 .write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
1071 .unwrap();
1072 producer.finish().unwrap();
1073
1074 let mut consumer = producer.consume();
1075 drop(producer);
1076
1077 let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
1079 assert_eq!(frame.payload, Bytes::from_static(b"data"));
1080 }
1081
1082 #[tokio::test]
1083 async fn pending_then_ready() {
1084 let mut producer = Info { sequence: 0 }.produce();
1085 let mut consumer = producer.consume();
1086
1087 assert!(consumer.next_frame().now_or_never().is_none());
1089
1090 producer
1091 .write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
1092 .unwrap();
1093 producer.finish().unwrap();
1094
1095 let frame = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
1096 assert_eq!(frame.size, 4);
1097 }
1098
1099 #[test]
1100 fn eviction_drops_old_frames() {
1101 let mut producer = Info { sequence: 0 }.produce();
1102
1103 let big = Bytes::from(vec![0u8; MAX_GROUP_CACHE as usize]);
1105 producer.write_frame(Timestamp::ZERO, big.clone()).unwrap();
1106 producer.write_frame(Timestamp::ZERO, big).unwrap();
1107
1108 let state = producer.state.read();
1110 assert_eq!(state.offset, 1);
1111 assert_eq!(state.frames.len(), 1);
1112 assert_eq!(state.frames[0].payload.len(), MAX_GROUP_CACHE as usize);
1113 }
1114
1115 #[test]
1116 fn next_frame_returns_cache_full_on_tombstone() {
1117 let mut producer = Info { sequence: 0 }.produce();
1118
1119 let big = Bytes::from(vec![0u8; MAX_GROUP_CACHE as usize]);
1120 producer.write_frame(Timestamp::ZERO, big.clone()).unwrap();
1121 producer.write_frame(Timestamp::ZERO, big).unwrap();
1122
1123 let mut consumer = producer.consume();
1124 let result = consumer.next_frame().now_or_never().unwrap();
1126 assert!(matches!(result, Err(crate::Error::Lagged)));
1127 }
1128
1129 #[test]
1130 fn no_eviction_under_budget() {
1131 let mut producer = Info { sequence: 0 }.produce();
1132 for _ in 0..100_000 {
1134 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"x")).unwrap();
1135 }
1136 producer.finish().unwrap();
1137
1138 let state = producer.state.read();
1139 assert_eq!(state.offset, 0);
1140 assert_eq!(state.frames.len(), 100_000);
1141 }
1142
1143 #[test]
1144 fn clone_consumer_independent() {
1145 let mut producer = Info { sequence: 0 }.produce();
1146 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"a")).unwrap();
1147
1148 let mut c1 = producer.consume();
1149 let _ = c1.next_frame().now_or_never().unwrap().unwrap().unwrap();
1151
1152 let mut c2 = c1.clone();
1154
1155 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"b")).unwrap();
1156 producer.finish().unwrap();
1157
1158 let f = c2.next_frame().now_or_never().unwrap().unwrap().unwrap();
1160 assert_eq!(f.size, 1); let end = c2.next_frame().now_or_never().unwrap().unwrap();
1163 assert!(end.is_none());
1164 }
1165
1166 #[test]
1169 fn read_frame_crosses_prefetch_batches() {
1170 let n = Prefetch::CAP * 3 + 5;
1171 let mut producer = Info { sequence: 0 }.produce();
1172 for i in 0..n {
1173 producer
1174 .write_frame(Timestamp::ZERO, Bytes::from(vec![i as u8; 4]))
1175 .unwrap();
1176 }
1177 producer.finish().unwrap();
1178
1179 let mut consumer = producer.consume();
1180 for i in 0..n {
1181 let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
1182 assert_eq!(frame.payload, Bytes::from(vec![i as u8; 4]));
1183 }
1184 assert!(consumer.read_frame().now_or_never().unwrap().unwrap().is_none());
1185 }
1186
1187 #[test]
1191 fn abort_after_finish_keeps_the_clean_end_for_a_drained_reader() {
1192 let mut producer = Info { sequence: 0 }.produce();
1193 producer
1194 .write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
1195 .unwrap();
1196 producer.finish().unwrap();
1197
1198 let mut drained = producer.consume();
1199 let mut behind = producer.consume();
1200 let frame = drained.read_frame().now_or_never().unwrap().unwrap().unwrap();
1201 assert_eq!(frame.payload, Bytes::from_static(b"hello"));
1202
1203 producer.abort(Error::Old).unwrap();
1204
1205 assert!(drained.read_frame().now_or_never().unwrap().unwrap().is_none());
1207 assert!(drained.next_frame().now_or_never().unwrap().unwrap().is_none());
1208
1209 assert!(matches!(behind.read_frame().now_or_never().unwrap(), Err(Error::Old)));
1211 }
1212
1213 #[test]
1216 fn finished_survives_a_later_abort() {
1217 let mut producer = Info { sequence: 0 }.produce();
1218 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"a")).unwrap();
1219 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"b")).unwrap();
1220 producer.finish().unwrap();
1221
1222 let mut consumer = producer.consume();
1223 producer.abort(Error::Old).unwrap();
1224
1225 assert_eq!(consumer.finished().now_or_never().unwrap().unwrap(), 2);
1226 }
1227
1228 #[test]
1230 fn interleave_read_and_next_frame() {
1231 let mut producer = Info { sequence: 0 }.produce();
1232 for i in 0..5u8 {
1233 producer.write_frame(Timestamp::ZERO, Bytes::from(vec![i; 1])).unwrap();
1234 }
1235 producer.finish().unwrap();
1236
1237 let mut consumer = producer.consume();
1238 let f0 = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
1240 assert_eq!(f0.payload, Bytes::from(vec![0u8; 1]));
1241
1242 for i in 1..5u8 {
1244 let mut f = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
1245 let data = f.read_all().now_or_never().unwrap().unwrap();
1246 assert_eq!(data, Bytes::from(vec![i; 1]));
1247 }
1248 assert!(consumer.next_frame().now_or_never().unwrap().unwrap().is_none());
1249 }
1250
1251 #[test]
1254 fn read_frame_past_cleared_frames_does_not_panic() {
1255 let mut producer = Info { sequence: 0 }.produce();
1256 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"a")).unwrap();
1257 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"b")).unwrap();
1258
1259 let mut consumer = producer.consume();
1260 consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
1261 consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
1262
1263 producer.abort(Error::Cancel).unwrap();
1266
1267 let result = consumer.read_frame().now_or_never().unwrap();
1268 assert!(matches!(result, Err(Error::Cancel)), "expected Cancel, got {result:?}");
1269 }
1270
1271 #[test]
1274 fn drop_with_partial_batch() {
1275 let mut producer = Info { sequence: 0 }.produce();
1276 for _ in 0..Prefetch::CAP {
1277 producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"x")).unwrap();
1278 }
1279 producer.finish().unwrap();
1280
1281 let mut consumer = producer.consume();
1282 let _ = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
1284 drop(consumer);
1285 }
1286
1287 #[tokio::test]
1292 async fn chunk_write_wakes_parked_reader() {
1293 let mut producer = Info { sequence: 0 }.produce();
1294 let mut consumer = producer.consume();
1295 let mut frame = producer
1296 .create_frame(frame::Info {
1297 size: 6,
1298 timestamp: Timestamp::ZERO,
1299 })
1300 .unwrap();
1301 let mut f = consumer.next_frame().await.unwrap().unwrap();
1302 let handle = tokio::spawn(async move { f.read_chunk().await });
1303 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
1305 frame.write(Bytes::from_static(b"foo")).unwrap();
1306 let chunk = tokio::time::timeout(std::time::Duration::from_secs(2), handle)
1307 .await
1308 .expect("parked chunk reader was never woken by the chunk write")
1309 .unwrap()
1310 .unwrap();
1311 assert_eq!(chunk, Some(Bytes::from_static(b"foo")));
1312 }
1313
1314 #[test]
1317 fn create_frame_converts_mismatched_scale() {
1318 use crate::{Timescale, Timestamp};
1319
1320 let mut producer = Producer::new(
1321 Info { sequence: 0 },
1322 track::Info::default().with_timescale(Timescale::MICRO),
1323 Default::default(),
1324 );
1325 let frame = frame::Info {
1326 size: 3,
1327 timestamp: Timestamp::from_millis(1).unwrap(), };
1329 let writer = producer.create_frame(frame).unwrap();
1330 assert_eq!(writer.timestamp.scale(), Timescale::MICRO);
1331 assert_eq!(writer.timestamp.value(), 1000);
1332 }
1333
1334 #[tokio::test]
1336 async fn create_frame_converts_current_timestamp() {
1337 use crate::Timescale;
1338
1339 let mut producer = Producer::new(
1340 Info { sequence: 0 },
1341 track::Info::default().with_timescale(Timescale::MICRO),
1342 Default::default(),
1343 );
1344 let writer = producer
1345 .create_frame(frame::Info {
1346 size: 3,
1347 timestamp: Timestamp::now(),
1348 })
1349 .unwrap();
1350 assert_eq!(writer.timestamp.scale(), Timescale::MICRO);
1351 assert!(!writer.timestamp.is_zero(), "local clock should be non-zero");
1352 }
1353
1354 #[test]
1356 fn create_frame_rejects_oversized() {
1357 let mut producer = Info { sequence: 0 }.produce();
1358 let result = producer.create_frame(frame::Info {
1359 size: MAX_GROUP_CACHE + 1,
1360 timestamp: Timestamp::ZERO,
1361 });
1362 assert!(matches!(result, Err(Error::FrameTooLarge)));
1363 }
1364}