1use std::collections::VecDeque;
2use std::task::{Poll, ready};
3
4use moq_net::Timestamp;
5
6use super::{Container, Frame};
7
8pub struct Consumer<F: Container> {
43 track: moq_net::track::Subscriber,
44
45 format: F,
46
47 current: u64,
49
50 pending: VecDeque<GroupBuffer>,
52
53 startup: bool,
56
57 latency: std::time::Duration,
59
60 rewind: Rewind,
62}
63
64#[derive(Default)]
71struct Rewind {
72 live_edge: Option<(u64, Timestamp)>,
75
76 boundary: Option<Reset>,
79
80 discontinuity: u64,
83}
84
85#[derive(Clone, Copy)]
92struct Reset {
93 prev_max: u64,
96
97 group: u64,
100
101 timestamp: Timestamp,
105}
106
107impl Reset {
108 fn by_sequence(&self, sequence: u64) -> Option<bool> {
111 if sequence <= self.prev_max {
112 Some(true)
113 } else if sequence >= self.group {
114 Some(false)
115 } else {
116 None
117 }
118 }
119
120 fn is_stale(&self, sequence: u64, timestamp: Timestamp) -> bool {
124 self.by_sequence(sequence).unwrap_or(timestamp >= self.timestamp)
125 }
126}
127
128impl<F: Container> Consumer<F> {
129 pub fn new(track: moq_net::track::Subscriber, format: F) -> Self {
133 Self {
134 track,
135 format,
136 current: 0,
137 pending: VecDeque::new(),
138 startup: true,
139 latency: std::time::Duration::ZERO,
140 rewind: Rewind::default(),
141 }
142 }
143
144 pub fn with_latency(mut self, latency: std::time::Duration) -> Self {
149 self.latency = latency;
150 self
151 }
152
153 pub fn discontinuity(&self) -> u64 {
162 self.rewind.discontinuity
163 }
164
165 pub async fn read(&mut self) -> Result<Option<Frame>, F::Error> {
173 kio::wait(|waiter| self.poll_read(waiter)).await
174 }
175
176 pub fn poll_read(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Frame>, F::Error>> {
181 let finished = self.poll_read_finish(waiter)?.is_ready();
183
184 if self.startup {
186 for (i, group) in self.pending.iter_mut().enumerate() {
188 if !matches!(group.poll_min_timestamp(waiter, &self.format), Poll::Ready(Ok(_))) {
191 continue;
192 }
193
194 self.current = group.sequence;
196 self.startup = false;
197 self.pending.drain(0..i);
198 break;
199 }
200 }
201
202 loop {
203 if self.poll_reset(waiter)? {
206 continue;
207 }
208
209 self.poll_classify(waiter)?;
211
212 while let Some(group) = self.pending.front_mut()
215 && group.sequence <= self.current
216 {
217 match group.poll_read(waiter, &self.format) {
218 Poll::Ready(Ok(Some(frame))) => {
219 let seq = group.group.sequence;
222 let ts = frame.timestamp;
223 if self.rewind.live_edge.is_none_or(|(_, high)| ts > high) {
224 self.rewind.live_edge = Some((seq, ts));
225 }
226 return Poll::Ready(Ok(Some(frame)));
227 }
228 Poll::Pending => break,
230 Poll::Ready(Err(e)) => {
231 if !group.poll_aborted(waiter) {
238 return Poll::Ready(Err(e));
239 }
240 tracing::warn!(error = ?e, "current group evicted; skipping to next buffered group");
246 self.pending.pop_front();
247 self.current = self.pending.front().map_or(self.current + 1, |g| g.sequence);
248 }
249 Poll::Ready(Ok(None)) => {
251 self.pending.pop_front();
252 self.current += 1;
253 }
254 }
255 }
256
257 let (oldest_timestamp, current_end) = if let Some(current) = self.pending.front_mut()
260 && current.sequence <= self.current
261 {
262 match current.poll_min_timestamp(waiter, &self.format) {
263 Poll::Ready(Ok(ts)) => (Some(std::time::Duration::from(ts)), current.max_end),
264 _ => (None, None),
265 }
266 } else {
267 (None, None)
268 };
269
270 let mut next_group = None;
272 for (i, group) in self.pending.iter_mut().enumerate() {
273 if group.sequence <= self.current {
274 continue;
275 }
276
277 if let Poll::Ready(Ok(ts)) = group.poll_min_timestamp(waiter, &self.format) {
278 next_group = Some((i, std::time::Duration::from(ts)));
279 break;
280 }
281 }
282
283 let mut max_timestamp = std::time::Duration::ZERO;
285 for group in self.pending.iter_mut().rev() {
286 if group.sequence <= self.current {
287 break;
288 }
289
290 if let Poll::Ready(Ok(ts)) = group.poll_max_timestamp(waiter, &self.format) {
291 max_timestamp = max_timestamp.max(ts.into());
292 break; }
294 }
295
296 let should_skip = if let Some((_, next_start)) = next_group {
297 if let Some(oldest) = oldest_timestamp {
298 let over_latency = max_timestamp.saturating_sub(oldest) >= self.latency;
303 let covered = current_end.is_some_and(|end| end >= next_start);
304 over_latency || covered
305 } else {
306 finished || self.pending.front().is_some_and(|g| g.sequence > self.current)
313 }
314 } else {
315 false
316 };
317
318 if let Some((new_idx, _)) = next_group
319 && should_skip
320 {
321 self.pending.drain(0..new_idx);
322 let new_current = self.pending.front().map(|g| g.sequence).unwrap();
323
324 tracing::debug!(old = self.current, new = new_current, "skipping slow groups");
325
326 self.current = new_current;
327 continue;
328 }
329
330 if finished && self.pending.is_empty() {
331 return Poll::Ready(Ok(None));
332 }
333
334 return Poll::Pending;
335 }
336 }
337
338 fn poll_read_finish(&mut self, waiter: &kio::Waiter) -> Poll<Result<(), F::Error>> {
342 loop {
343 let Some(group) = ready!(self.track.poll_recv_group(waiter)?) else {
344 return Poll::Ready(Ok(()));
346 };
347
348 let reader = GroupBuffer::new(group);
349 let sequence = reader.group.sequence;
350
351 let drop = match &self.rewind.boundary {
356 Some(reset) => match reset.by_sequence(sequence) {
357 Some(true) => true, Some(false) => sequence < self.current, None => false, },
361 None => sequence < self.current,
362 };
363 if drop {
364 tracing::debug!(old = ?sequence, current = ?self.current, "skipping old group");
365 continue;
366 }
367
368 let idx = self
369 .pending
370 .partition_point(|g| g.group.sequence < reader.group.sequence);
371 self.pending.insert(idx, reader);
372 }
373 }
374
375 fn poll_reset(&mut self, waiter: &kio::Waiter) -> Result<bool, F::Error> {
387 let Some((prev_max, live_edge)) = self.rewind.live_edge else {
388 return Ok(false);
389 };
390
391 let reset = {
397 let mut found = None;
398 for group in self.pending.iter_mut().rev() {
399 if group.group.sequence <= self.current {
401 break;
402 }
403
404 let Poll::Ready(Ok(min)) = group.poll_min_timestamp(waiter, &self.format) else {
406 continue;
407 };
408
409 if min < live_edge {
410 found = Some(Reset {
411 prev_max,
412 group: group.group.sequence,
413 timestamp: min,
414 });
415 break;
416 }
417 }
418
419 let Some(reset) = found else {
420 return Ok(false);
421 };
422 reset
423 };
424
425 self.pending.retain(|g| match reset.by_sequence(g.group.sequence) {
428 Some(stale) => !stale,
429 None => g.min_timestamp.is_none_or(|ts| !reset.is_stale(g.group.sequence, ts)),
430 });
431
432 self.rewind.discontinuity += 1;
433 tracing::debug!(
434 prev_max = reset.prev_max,
435 group = reset.group,
436 discontinuity = self.rewind.discontinuity,
437 "buffer reset: group timestamps rewound"
438 );
439 self.rewind.boundary = Some(reset);
440 self.current = self.pending.front().map_or(reset.group, |g| g.group.sequence);
442 self.rewind.live_edge = Some((reset.group, reset.timestamp));
443
444 Ok(true)
445 }
446
447 fn poll_classify(&mut self, waiter: &kio::Waiter) -> Result<(), F::Error> {
454 let Some(reset) = self.rewind.boundary else {
455 return Ok(());
456 };
457
458 let mut i = 0;
459 while i < self.pending.len() {
460 let group = &mut self.pending[i];
461 if reset.by_sequence(group.group.sequence).is_some() {
463 i += 1;
464 continue;
465 }
466
467 match group.poll_min_timestamp(waiter, &self.format) {
468 Poll::Ready(Ok(min)) if reset.is_stale(group.group.sequence, min) => {
469 self.pending.remove(i);
470 }
471 _ => i += 1,
472 }
473 }
474
475 Ok(())
476 }
477
478 pub fn set_latency(&mut self, latency: std::time::Duration) {
480 self.latency = latency;
481 }
482}
483
484struct GroupBuffer {
489 group: moq_net::group::Consumer,
490
491 index: usize,
493
494 buffered: VecDeque<Frame>,
496
497 min_timestamp: Option<Timestamp>,
499
500 max_timestamp: Option<Timestamp>,
502
503 max_end: Option<std::time::Duration>,
507}
508
509impl GroupBuffer {
510 fn new(group: moq_net::group::Consumer) -> Self {
511 Self {
512 group,
513 index: 0,
514 buffered: VecDeque::new(),
515 max_timestamp: None,
516 min_timestamp: None,
517 max_end: None,
518 }
519 }
520
521 fn poll_read<F: Container>(&mut self, waiter: &kio::Waiter, format: &F) -> Poll<Result<Option<Frame>, F::Error>> {
523 if let Some(frame) = self.buffered.pop_front() {
524 return Poll::Ready(Ok(Some(frame)));
525 }
526
527 match ready!(self.buffer_one(waiter, format)?) {
528 true => Poll::Ready(Ok(Some(self.buffered.pop_front().unwrap()))),
529 false => Poll::Ready(Ok(None)),
530 }
531 }
532
533 fn buffer_once<F: Container>(&mut self, waiter: &kio::Waiter, format: &F) -> Poll<Result<bool, F::Error>> {
537 let Some(frames) = ready!(format.poll_read(&mut self.group, waiter)?) else {
538 return Poll::Ready(Ok(false));
539 };
540
541 for mut frame in frames {
542 self.min_timestamp = Some(match self.min_timestamp {
543 Some(existing) => existing.min(frame.timestamp),
544 None => frame.timestamp,
545 });
546
547 self.max_timestamp = Some(match self.max_timestamp {
548 Some(existing) => existing.max(frame.timestamp),
549 None => frame.timestamp,
550 });
551
552 let duration = frame.duration.map(std::time::Duration::from).unwrap_or_default();
556 let end = std::time::Duration::from(frame.timestamp) + duration;
557 self.max_end = Some(match self.max_end {
558 Some(existing) => existing.max(end),
559 None => end,
560 });
561
562 frame.keyframe = frame.keyframe || self.index == 0;
565 self.index += 1;
566
567 self.buffered.push_back(frame);
568 }
569
570 Poll::Ready(Ok(true))
571 }
572
573 fn buffer_one<F: Container>(&mut self, waiter: &kio::Waiter, format: &F) -> Poll<Result<bool, F::Error>> {
574 loop {
575 if !self.buffered.is_empty() {
576 return Poll::Ready(Ok(true));
577 }
578 if !ready!(self.buffer_once(waiter, format)?) {
579 return Poll::Ready(Ok(false));
580 }
581 }
584 }
585
586 fn buffer_all<F: Container>(&mut self, waiter: &kio::Waiter, format: &F) -> Poll<Result<(), F::Error>> {
587 while ready!(self.buffer_once(waiter, format)?) {}
588 Poll::Ready(Ok(()))
589 }
590
591 fn poll_max_timestamp<F: Container>(
593 &mut self,
594 waiter: &kio::Waiter,
595 format: &F,
596 ) -> Poll<Result<Timestamp, F::Error>> {
597 let _ = self.buffer_all(waiter, format)?;
599
600 if let Some(max) = self.max_timestamp {
601 return Poll::Ready(Ok(max));
602 }
603
604 if let Poll::Ready(_frames) = self.group.poll_finished(waiter)? {
605 return Poll::Ready(Err(moq_net::Error::Decode(moq_net::DecodeError::Short).into()));
606 }
607
608 Poll::Pending
609 }
610
611 fn poll_min_timestamp<F: Container>(
612 &mut self,
613 waiter: &kio::Waiter,
614 format: &F,
615 ) -> Poll<Result<Timestamp, F::Error>> {
616 let _ = self.buffer_one(waiter, format)?;
617
618 if let Some(min) = self.min_timestamp {
619 return Poll::Ready(Ok(min));
620 }
621
622 if let Poll::Ready(_frames) = self.group.poll_finished(waiter)? {
623 return Poll::Ready(Err(moq_net::Error::Decode(moq_net::DecodeError::Short).into()));
624 }
625
626 Poll::Pending
627 }
628
629 fn poll_aborted(&mut self, waiter: &kio::Waiter) -> bool {
635 matches!(self.group.poll_finished(waiter), Poll::Ready(Err(_)))
636 }
637}
638
639impl std::ops::Deref for GroupBuffer {
640 type Target = moq_net::group::Consumer;
641
642 fn deref(&self) -> &Self::Target {
643 &self.group
644 }
645}
646
647#[cfg(test)]
648mod tests {
649 use super::Container as ContainerTrait;
650 use super::*;
651 use crate::catalog::hang::Container;
652 use std::time::Duration;
653
654 use bytes::Bytes;
655
656 fn track_producer(
659 name: impl Into<std::sync::Arc<str>>,
660 info: impl Into<Option<moq_net::track::Info>>,
661 ) -> moq_net::track::Producer {
662 moq_net::broadcast::Info::new()
663 .produce()
664 .create_track(name, info)
665 .unwrap()
666 }
667
668 fn ts(micros: u64) -> Timestamp {
669 Timestamp::from_micros(micros).unwrap()
670 }
671
672 struct DurationWire;
676
677 fn encode_duration_frame(timestamp: Timestamp, duration: Timestamp) -> Vec<u8> {
679 let mut buf = Vec::with_capacity(18);
680 buf.extend_from_slice(&(timestamp.as_micros() as u64).to_le_bytes());
681 buf.extend_from_slice(&(duration.as_micros() as u64).to_le_bytes());
682 buf.extend_from_slice(&[0xDE, 0xAD]);
683 buf
684 }
685
686 impl ContainerTrait for DurationWire {
687 type Error = crate::Error;
688
689 fn write(&self, group: &mut moq_net::group::Producer, frames: &[Frame]) -> Result<(), Self::Error> {
690 for frame in frames {
693 group.write_frame(frame.timestamp, encode_duration_frame(frame.timestamp, ts(0)))?;
694 }
695 Ok(())
696 }
697
698 fn poll_read(
699 &self,
700 group: &mut moq_net::group::Consumer,
701 waiter: &kio::Waiter,
702 ) -> Poll<Result<Option<Vec<Frame>>, Self::Error>> {
703 use bytes::Buf;
704
705 let Some(mut data) = ready!(group.poll_read_frame(waiter)?).map(|f| f.payload) else {
706 return Poll::Ready(Ok(None));
707 };
708
709 let timestamp = ts(data.get_u64_le());
710 let duration = ts(data.get_u64_le());
711 let payload = data.copy_to_bytes(data.remaining());
712
713 Poll::Ready(Ok(Some(vec![Frame {
714 timestamp,
715 payload,
716 keyframe: false,
717 duration: Some(duration),
718 }])))
719 }
720 }
721
722 fn write_duration_frame(group: &mut moq_net::group::Producer, timestamp: Timestamp, duration: Timestamp) {
724 group
725 .write_frame(timestamp, encode_duration_frame(timestamp, duration))
726 .unwrap();
727 }
728
729 fn write_group(track: &mut moq_net::track::Producer, sequence: u64, timestamps: &[Timestamp]) {
731 let mut group = track.create_group(moq_net::group::Info { sequence }).unwrap();
732 for ×tamp in timestamps {
733 let frame = Frame {
734 timestamp,
735 payload: Bytes::from_static(&[0xDE, 0xAD]),
736 keyframe: false,
737 duration: None,
738 };
739 Container::Legacy.write(&mut group, &[frame]).unwrap();
740 }
741 group.finish().unwrap();
742 }
743
744 async fn read_all(consumer: &mut Consumer<Container>) -> Result<Vec<Frame>, crate::Error> {
746 let mut frames = Vec::new();
747 loop {
748 match tokio::time::timeout(Duration::from_millis(200), consumer.read()).await {
749 Ok(Ok(Some(frame))) => frames.push(frame),
750 Ok(Ok(None)) => break,
751 Ok(Err(e)) => return Err(e),
752 Err(_) => panic!(
753 "read_all: Consumer::read timed out after 200ms ({} frames collected so far)",
754 frames.len()
755 ),
756 }
757 }
758 Ok(frames)
759 }
760
761 #[tokio::test]
764 async fn read_single_group() {
765 let mut track = track_producer("test", hang::container::track_info());
766 let consumer_track = track.subscribe(None);
767 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(500));
768
769 write_group(&mut track, 0, &[ts(0)]);
770 track.finish().unwrap();
771
772 let frames = read_all(&mut consumer).await.unwrap();
773 assert_eq!(frames.len(), 1);
774 assert_eq!(frames[0].timestamp, ts(0));
775 assert!(frames[0].keyframe);
776
777 assert!(consumer.read().await.unwrap().is_none());
779 }
780
781 #[tokio::test]
782 async fn read_multiple_frames_single_group() {
783 let mut track = track_producer("test", hang::container::track_info());
784 let consumer_track = track.subscribe(None);
785 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(500));
786
787 write_group(&mut track, 0, &[ts(0), ts(33_000), ts(66_000)]);
788 track.finish().unwrap();
789
790 let frames = read_all(&mut consumer).await.unwrap();
791 assert_eq!(frames.len(), 3);
792 assert_eq!(frames[0].timestamp, ts(0));
793 assert_eq!(frames[1].timestamp, ts(33_000));
794 assert_eq!(frames[2].timestamp, ts(66_000));
795
796 assert!(frames[0].keyframe);
797 }
798
799 #[tokio::test]
800 async fn read_multiple_groups_within_latency() {
801 let mut track = track_producer("test", hang::container::track_info());
802 let consumer_track = track.subscribe(None);
803 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(500));
804
805 for i in 0..5u64 {
807 write_group(&mut track, i, &[ts(i * 20_000)]);
808 }
809 track.finish().unwrap();
810
811 let frames = read_all(&mut consumer).await.unwrap();
812 assert_eq!(frames.len(), 5);
813 }
814
815 #[tokio::test]
818 async fn latency_skip_delivers_recent_groups() {
819 tokio::time::pause();
820 let mut track = track_producer("test", hang::container::track_info());
821 let consumer_track = track.subscribe(None);
822 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(100));
823
824 let mut group0 = track.create_group(moq_net::group::Info { sequence: 0 }).unwrap();
826 for f in 0..5u64 {
827 Container::Legacy
828 .write(
829 &mut group0,
830 &[Frame {
831 timestamp: ts(f * 2_000),
832 payload: Bytes::from_static(&[0xDE, 0xAD]),
833 keyframe: false,
834 duration: None,
835 }],
836 )
837 .unwrap();
838 }
839
840 for g in 1..20u64 {
842 let timestamps: Vec<_> = (0..5).map(|f| ts(g * 15_000 + f * 2_000)).collect();
843 write_group(&mut track, g, ×tamps);
844 }
845 track.finish().unwrap();
846
847 let finisher = tokio::spawn(async move {
849 tokio::time::sleep(Duration::from_millis(50)).await;
850 group0.finish().unwrap();
851 });
852
853 let frames = read_all(&mut consumer).await.unwrap();
854 assert!(frames.len() >= 25, "Expected >= 25 frames, got {}", frames.len());
856 finisher.await.expect("finisher task panicked");
857 }
858
859 #[tokio::test]
860 async fn zero_latency_skips_aggressively() {
861 tokio::time::pause();
862 let mut track = track_producer("test", hang::container::track_info());
863 let consumer_track = track.subscribe(None);
864 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::ZERO);
865
866 let mut group0 = track.create_group(moq_net::group::Info { sequence: 0 }).unwrap();
869 Container::Legacy
870 .write(
871 &mut group0,
872 &[Frame {
873 timestamp: ts(0),
874 payload: Bytes::from_static(&[0xDE, 0xAD]),
875 keyframe: false,
876 duration: None,
877 }],
878 )
879 .unwrap();
880
881 for g in 1..10u64 {
882 let timestamps: Vec<_> = (0..3).map(|f| ts(g * 50_000 + f * 5_000)).collect();
883 write_group(&mut track, g, ×tamps);
884 }
885 track.finish().unwrap();
886
887 let finisher = tokio::spawn(async move {
888 tokio::time::sleep(Duration::from_millis(50)).await;
889 group0.finish().unwrap();
890 });
891
892 let frames = read_all(&mut consumer).await.unwrap();
893 assert_eq!(frames.len(), 28, "Expected group 0 frame + groups 1-9");
894 assert!(!frames.is_empty(), "Expected at least some frames");
895 finisher.await.expect("finisher task panicked");
896 }
897
898 #[tokio::test]
899 async fn latency_skip_correctness() {
900 tokio::time::pause();
901 let mut track = track_producer("test", hang::container::track_info());
902 let consumer_track = track.subscribe(None);
903 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(100));
904
905 let mut group0 = track.create_group(moq_net::group::Info { sequence: 0 }).unwrap();
906 Container::Legacy
907 .write(
908 &mut group0,
909 &[Frame {
910 timestamp: ts(0),
911 payload: Bytes::from_static(&[0xDE, 0xAD]),
912 keyframe: false,
913 duration: None,
914 }],
915 )
916 .unwrap();
917
918 for g in 1..10u64 {
919 write_group(&mut track, g, &[ts(g * 30_000)]);
920 }
921 track.finish().unwrap();
922
923 let finisher = tokio::spawn(async move {
924 tokio::time::sleep(Duration::from_millis(50)).await;
925 group0.finish().unwrap();
926 });
927
928 let frames = read_all(&mut consumer).await.unwrap();
929 assert!(!frames.is_empty(), "Expected at least some frames");
930 assert_eq!(frames.len(), 10, "Expected group 0 frame + groups 1-9");
931 assert_eq!(frames[0].timestamp, ts(0));
932
933 for i in 1..10u64 {
934 assert_eq!(frames[i as usize].timestamp, ts(i * 30_000));
935 }
936 finisher.await.expect("finisher task panicked");
937 }
938
939 #[test]
944 fn reset_classifies_out_of_order_groups() {
945 let reset = Reset {
946 prev_max: 55,
947 group: 58,
948 timestamp: ts(90),
949 };
950
951 assert!(!reset.is_stale(57, ts(88)));
953 assert!(reset.is_stale(52, ts(86)));
956 assert!(reset.is_stale(56, ts(105)));
958 assert!(!reset.is_stale(58, ts(90)));
960 assert!(!reset.is_stale(59, ts(92)));
961 assert!(reset.is_stale(55, ts(100)));
963 }
964
965 #[tokio::test]
968 async fn reset_keeps_out_of_order_new_group() {
969 tokio::time::pause();
970 let mut track = track_producer("test", hang::container::track_info());
971 let consumer_track = track.subscribe(None);
972 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_secs(10));
973
974 write_group(&mut track, 0, &[ts(0)]);
976 write_group(&mut track, 1, &[ts(100_000)]);
977 write_group(&mut track, 2, &[ts(200_000)]);
978 write_group(&mut track, 5, &[ts(3_000)]);
980
981 let finisher = tokio::spawn(async move {
983 tokio::time::sleep(Duration::from_millis(50)).await;
984 write_group(&mut track, 3, &[ts(1_000)]);
985 write_group(&mut track, 4, &[ts(2_000)]);
986 track.finish().unwrap();
987 });
988
989 let frames = read_all(&mut consumer).await.unwrap();
990 let micros: Vec<u128> = frames.iter().map(|f| f.timestamp.as_micros()).collect();
991
992 assert!(micros.contains(&100_000), "old epoch played before the reset");
995 assert!(
996 micros.contains(&1_000) && micros.contains(&2_000) && micros.contains(&3_000),
997 "out-of-order new-epoch groups kept, got {micros:?}"
998 );
999 assert_eq!(consumer.discontinuity(), 1, "one rewind detected");
1000 finisher.await.expect("finisher task panicked");
1001 }
1002
1003 #[tokio::test]
1007 async fn reset_detected_behind_forward_newest_group() {
1008 tokio::time::pause();
1009 let mut track = track_producer("test", hang::container::track_info());
1010 let consumer_track = track.subscribe(None);
1011 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_secs(10));
1012
1013 write_group(&mut track, 0, &[ts(0)]);
1015 write_group(&mut track, 1, &[ts(100_000)]);
1016 write_group(&mut track, 2, &[ts(200_000)]);
1017 write_group(&mut track, 6, &[ts(250_000)]);
1019 write_group(&mut track, 5, &[ts(50_000)]);
1021 track.finish().unwrap();
1022
1023 let frames = read_all(&mut consumer).await.unwrap();
1024 let micros: Vec<u128> = frames.iter().map(|f| f.timestamp.as_micros()).collect();
1025
1026 assert_eq!(
1027 consumer.discontinuity(),
1028 1,
1029 "rewind detected behind a forward newest group"
1030 );
1031 assert!(micros.contains(&50_000), "resumed at the rewound group, got {micros:?}");
1032 assert!(
1033 !micros.contains(&200_000),
1034 "the reneged tail was dropped, got {micros:?}"
1035 );
1036 }
1037
1038 #[tokio::test]
1042 async fn backwards_timestamp_resets_buffer() {
1043 tokio::time::pause();
1044 let mut track = track_producer("test", hang::container::track_info());
1045 let consumer_track = track.subscribe(None);
1046 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_secs(10));
1048
1049 for i in 0..5u64 {
1051 write_group(&mut track, i, &[ts(i * 100_000)]);
1052 }
1053 write_group(&mut track, 5, &[ts(0), ts(20_000)]);
1055 track.finish().unwrap();
1056
1057 let frames = read_all(&mut consumer).await.unwrap();
1058 let timestamps: Vec<_> = frames.iter().map(|f| f.timestamp).collect();
1059
1060 assert_eq!(timestamps, vec![ts(0), ts(100_000), ts(0), ts(20_000)]);
1063 assert_eq!(consumer.discontinuity(), 1);
1064 }
1065
1066 #[tokio::test]
1069 async fn backwards_timestamp_always_resets() {
1070 let mut track = track_producer("test", hang::container::track_info());
1071 let consumer_track = track.subscribe(None);
1072 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_secs(10));
1073
1074 write_group(&mut track, 0, &[ts(0)]);
1075 write_group(&mut track, 1, &[ts(500_000)]);
1076 write_group(&mut track, 2, &[ts(0)]); track.finish().unwrap();
1078
1079 let frames = read_all(&mut consumer).await.unwrap();
1080 let timestamps: Vec<_> = frames.iter().map(|f| f.timestamp).collect();
1081
1082 assert_eq!(timestamps, vec![ts(0), ts(500_000), ts(0)]);
1083 assert_eq!(consumer.discontinuity(), 1, "the backwards group triggered a reset");
1084 }
1085
1086 fn write_marker(group: &mut moq_net::group::Producer, timestamp: Timestamp) {
1091 let frame = Frame {
1092 timestamp,
1093 payload: Bytes::new(),
1094 keyframe: false,
1095 duration: None,
1096 };
1097 Container::Legacy.write(group, &[frame]).unwrap();
1098 }
1099
1100 #[tokio::test]
1104 async fn empty_payload_is_skipped() {
1105 let mut track = track_producer("test", hang::container::track_info());
1106 let consumer_track = track.subscribe(None);
1107 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(500));
1108
1109 let mut group = track.create_group(moq_net::group::Info { sequence: 0 }).unwrap();
1110 let media = |timestamp| Frame {
1111 timestamp,
1112 payload: Bytes::from_static(&[0xDE, 0xAD]),
1113 keyframe: false,
1114 duration: None,
1115 };
1116 Container::Legacy.write(&mut group, &[media(ts(0))]).unwrap();
1117 write_marker(&mut group, ts(16_000)); Container::Legacy.write(&mut group, &[media(ts(33_000))]).unwrap();
1119 write_marker(&mut group, ts(50_000)); group.finish().unwrap();
1121 track.finish().unwrap();
1122
1123 let frames = read_all(&mut consumer).await.unwrap();
1124 assert_eq!(frames.len(), 2, "markers are not surfaced as media");
1125 assert_eq!(frames[0].timestamp, ts(0));
1126 assert_eq!(frames[1].timestamp, ts(33_000));
1127 }
1128
1129 #[tokio::test]
1133 async fn consecutive_markers_do_not_stall() {
1134 let mut track = track_producer("test", hang::container::track_info());
1135 let consumer_track = track.subscribe(None);
1136 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(500));
1137
1138 let mut group = track.create_group(moq_net::group::Info { sequence: 0 }).unwrap();
1139 Container::Legacy
1140 .write(
1141 &mut group,
1142 &[Frame {
1143 timestamp: ts(0),
1144 payload: Bytes::from_static(&[0xDE, 0xAD]),
1145 keyframe: false,
1146 duration: None,
1147 }],
1148 )
1149 .unwrap();
1150 for i in 1..5u64 {
1151 write_marker(&mut group, ts(i * 1_000));
1152 }
1153 group.finish().unwrap();
1154 write_group(&mut track, 1, &[ts(100_000)]);
1155 track.finish().unwrap();
1156
1157 let frames = read_all(&mut consumer).await.unwrap();
1158 let micros: Vec<u128> = frames.iter().map(|f| f.timestamp.as_micros()).collect();
1159 assert_eq!(micros, vec![0, 100_000], "markers skipped, next group reached");
1160 }
1161
1162 #[tokio::test]
1165 async fn groups_delivered_in_sequence_order() {
1166 tokio::time::pause();
1167 let mut track = track_producer("test", hang::container::track_info());
1168 let consumer_track = track.subscribe(None);
1169 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(500));
1170
1171 let mut group0 = track.create_group(moq_net::group::Info { sequence: 0 }).unwrap();
1172 Container::Legacy
1173 .write(
1174 &mut group0,
1175 &[Frame {
1176 timestamp: ts(0),
1177 payload: Bytes::from_static(&[0xDE, 0xAD]),
1178 keyframe: false,
1179 duration: None,
1180 }],
1181 )
1182 .unwrap();
1183
1184 write_group(&mut track, 2, &[ts(60_000)]);
1185 write_group(&mut track, 1, &[ts(30_000)]);
1186 track.finish().unwrap();
1187
1188 let finisher = tokio::spawn(async move {
1189 tokio::time::sleep(Duration::from_millis(10)).await;
1190 group0.finish().unwrap();
1191 });
1192
1193 let frames = read_all(&mut consumer).await.unwrap();
1194 assert_eq!(frames.len(), 3);
1195 assert_eq!(frames[0].timestamp, ts(0));
1196 assert_eq!(frames[1].timestamp, ts(30_000));
1197 assert_eq!(frames[2].timestamp, ts(60_000));
1198 finisher.await.expect("finisher task panicked");
1199 }
1200
1201 #[tokio::test]
1202 async fn adjacent_group_flushed_immediately() {
1203 let mut track = track_producer("test", hang::container::track_info());
1204 let consumer_track = track.subscribe(None);
1205 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(500));
1206
1207 write_group(&mut track, 0, &[ts(0)]);
1208 write_group(&mut track, 1, &[ts(30_000)]);
1209 track.finish().unwrap();
1210
1211 let frames = read_all(&mut consumer).await.unwrap();
1212 assert_eq!(frames.len(), 2);
1213 assert_eq!(frames[0].timestamp, ts(0));
1214 assert_eq!(frames[1].timestamp, ts(30_000));
1215 }
1216
1217 #[tokio::test]
1220 async fn bframes_within_group() {
1221 let mut track = track_producer("test", hang::container::track_info());
1222 let consumer_track = track.subscribe(None);
1223 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(500));
1224
1225 write_group(&mut track, 0, &[ts(0), ts(66_000), ts(33_000)]);
1226 track.finish().unwrap();
1227
1228 let frames = read_all(&mut consumer).await.unwrap();
1229 assert_eq!(frames.len(), 3);
1230 assert_eq!(frames[0].timestamp, ts(0));
1231 assert_eq!(frames[1].timestamp, ts(66_000));
1232 assert_eq!(frames[2].timestamp, ts(33_000));
1233 }
1234
1235 #[tokio::test]
1238 async fn empty_track_returns_none() {
1239 tokio::time::pause();
1240 let mut track = track_producer("test", hang::container::track_info());
1241 let consumer_track = track.subscribe(None);
1242 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(500));
1243
1244 track.finish().unwrap();
1245
1246 let result = tokio::time::timeout(Duration::from_millis(200), consumer.read()).await;
1247 match result {
1248 Ok(Ok(None)) => {} Ok(Ok(Some(_))) => panic!("expected None for empty track, got Some"),
1250 Ok(Err(e)) => panic!("expected None for empty track, got error: {e}"),
1251 Err(_) => panic!("should not hang on empty track"),
1252 }
1253 }
1254
1255 #[tokio::test]
1256 async fn track_closed_with_error() {
1257 tokio::time::pause();
1258 let mut track = track_producer("test", hang::container::track_info());
1259 let consumer_track = track.subscribe(None);
1260 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(500));
1261
1262 write_group(&mut track, 0, &[ts(0)]);
1263 track.abort(moq_net::Error::Cancel).unwrap();
1264
1265 let result = tokio::time::timeout(Duration::from_millis(500), async {
1266 let mut frames = Vec::new();
1267 while let Ok(Some(frame)) = consumer.read().await {
1268 frames.push(frame);
1269 }
1270 frames
1271 })
1272 .await;
1273
1274 assert!(result.is_ok(), "Consumer should not hang after track error");
1275 }
1276
1277 #[tokio::test]
1280 async fn gap_in_group_sequence_recovery() {
1281 let mut track = track_producer("test", hang::container::track_info());
1282 let consumer_track = track.subscribe(None);
1283 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(100));
1284
1285 write_group(&mut track, 0, &[ts(0), ts(20_000)]);
1286 write_group(&mut track, 1, &[ts(40_000), ts(60_000)]);
1287 write_group(&mut track, 3, &[ts(120_000), ts(140_000)]);
1288 write_group(&mut track, 4, &[ts(160_000), ts(180_000)]);
1289 write_group(&mut track, 5, &[ts(200_000), ts(220_000)]);
1290 write_group(&mut track, 6, &[ts(240_000), ts(260_000)]);
1291 track.finish().unwrap();
1292
1293 let frames = read_all(&mut consumer).await.unwrap();
1294 assert!(frames.len() >= 4, "Expected >= 4 frames, got {}", frames.len());
1295 }
1296
1297 #[tokio::test]
1298 async fn gap_at_start_of_sequence() {
1299 let mut track = track_producer("test", hang::container::track_info());
1300 let consumer_track = track.subscribe(None);
1301 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(80));
1302
1303 write_group(&mut track, 5, &[ts(0), ts(20_000)]);
1304 write_group(&mut track, 7, &[ts(80_000), ts(100_000)]);
1305 write_group(&mut track, 8, &[ts(120_000), ts(140_000)]);
1306 write_group(&mut track, 9, &[ts(160_000), ts(180_000)]);
1307 track.finish().unwrap();
1308
1309 let frames = read_all(&mut consumer).await.unwrap();
1310 assert!(frames.len() >= 4, "Expected >= 4 frames, got {}", frames.len());
1311 }
1312
1313 #[tokio::test]
1321 async fn evicted_group_with_gap_skips_to_live() {
1322 let mut track = track_producer("test", hang::container::track_info());
1323 let consumer_track = track.subscribe(None);
1324 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(100));
1325
1326 let mut group0 = track.create_group(moq_net::group::Info { sequence: 0 }).unwrap();
1328 Container::Legacy
1329 .write(
1330 &mut group0,
1331 &[Frame {
1332 timestamp: ts(0),
1333 payload: Bytes::from_static(&[0xDE, 0xAD]),
1334 keyframe: false,
1335 duration: None,
1336 }],
1337 )
1338 .unwrap();
1339 let first = consumer.read().await.unwrap().unwrap();
1340 assert_eq!(first.timestamp, ts(0));
1341
1342 write_group(&mut track, 5, &[ts(150_000)]);
1345
1346 group0.abort(moq_net::Error::Old).unwrap();
1348
1349 let next = tokio::time::timeout(Duration::from_secs(1), consumer.read())
1352 .await
1353 .expect("consumer hung on an evicted group / gap")
1354 .unwrap()
1355 .unwrap();
1356 assert_eq!(next.timestamp, ts(150_000), "skipped the evicted gap to the live group");
1357 }
1358
1359 #[tokio::test]
1364 async fn missing_sequence_skips_on_live_track() {
1365 let mut track = track_producer("test", hang::container::track_info());
1366 let consumer_track = track.subscribe(None);
1367 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(100));
1368
1369 write_group(&mut track, 0, &[ts(0), ts(20_000)]);
1372 write_group(&mut track, 2, &[ts(200_000)]);
1373
1374 let reached = tokio::time::timeout(Duration::from_secs(1), async {
1376 loop {
1377 let frame = consumer.read().await.unwrap().unwrap();
1378 if frame.timestamp == ts(200_000) {
1379 return;
1380 }
1381 }
1382 })
1383 .await;
1384 assert!(reached.is_ok(), "consumer hung on a missing sequence on a live track");
1385 }
1386
1387 struct FailingDecode;
1393
1394 impl ContainerTrait for FailingDecode {
1395 type Error = crate::Error;
1396
1397 fn write(&self, group: &mut moq_net::group::Producer, frames: &[Frame]) -> Result<(), Self::Error> {
1398 for frame in frames {
1399 group.write_frame(moq_net::Timestamp::ZERO, frame.payload.clone())?;
1400 }
1401 Ok(())
1402 }
1403
1404 fn poll_read(
1405 &self,
1406 group: &mut moq_net::group::Consumer,
1407 waiter: &kio::Waiter,
1408 ) -> Poll<Result<Option<Vec<Frame>>, Self::Error>> {
1409 use bytes::Buf;
1410
1411 let Some(mut data) = ready!(group.poll_read_frame(waiter)?).map(|f| f.payload) else {
1412 return Poll::Ready(Ok(None));
1413 };
1414 if data.as_ref() == b"FAIL" {
1415 return Poll::Ready(Err(crate::Error::UnknownFormat("malformed payload".into())));
1416 }
1417 Poll::Ready(Ok(Some(vec![Frame {
1418 timestamp: ts(data.get_u64_le()),
1419 payload: Bytes::new(),
1420 keyframe: false,
1421 duration: None,
1422 }])))
1423 }
1424 }
1425
1426 #[tokio::test]
1430 async fn decode_error_propagates() {
1431 tokio::time::pause();
1432 let mut track = track_producer("test", None);
1433 let consumer_track = track.subscribe(None);
1434 let mut consumer = Consumer::new(consumer_track, FailingDecode);
1435
1436 let mut group = track.create_group(moq_net::group::Info { sequence: 0 }).unwrap();
1438 group
1439 .write_frame(moq_net::Timestamp::ZERO, Bytes::from(0u64.to_le_bytes().to_vec()))
1440 .unwrap();
1441 group
1442 .write_frame(moq_net::Timestamp::ZERO, Bytes::from_static(b"FAIL"))
1443 .unwrap();
1444 group.finish().unwrap();
1445 track.finish().unwrap();
1446
1447 let first = consumer.read().await;
1449 assert!(matches!(first, Ok(Some(_))), "first frame should decode, got {first:?}");
1450
1451 let second = tokio::time::timeout(Duration::from_millis(200), consumer.read())
1452 .await
1453 .expect("consumer hung on a decode error");
1454 assert!(
1455 matches!(second, Err(crate::Error::UnknownFormat(_))),
1456 "decode error must propagate, got {second:?}"
1457 );
1458 }
1459
1460 #[tokio::test]
1463 async fn frame_timestamp_and_index_decoding() {
1464 let mut track = track_producer("test", hang::container::track_info());
1465 let consumer_track = track.subscribe(None);
1466 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(500));
1467
1468 write_group(&mut track, 0, &[ts(0), ts(33_333), ts(66_666)]);
1469 track.finish().unwrap();
1470
1471 let frames = read_all(&mut consumer).await.unwrap();
1472 assert_eq!(frames.len(), 3);
1473
1474 assert_eq!(frames[0].timestamp, ts(0));
1475 assert!(frames[0].keyframe);
1476
1477 assert_eq!(frames[1].timestamp, ts(33_333));
1478
1479 assert_eq!(frames[2].timestamp, ts(66_666));
1480 }
1481
1482 #[tokio::test]
1483 async fn frame_payload_preserved() {
1484 let mut track = track_producer("test", hang::container::track_info());
1485 let consumer_track = track.subscribe(None);
1486 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(500));
1487
1488 let payload_bytes = vec![0x01, 0x02, 0x03, 0x04, 0x05];
1489 let mut group = track.create_group(moq_net::group::Info { sequence: 0 }).unwrap();
1490 Container::Legacy
1491 .write(
1492 &mut group,
1493 &[Frame {
1494 timestamp: ts(0),
1495 payload: Bytes::from(payload_bytes.clone()),
1496
1497 keyframe: false,
1498 duration: None,
1499 }],
1500 )
1501 .unwrap();
1502 group.finish().unwrap();
1503 track.finish().unwrap();
1504
1505 let frames = read_all(&mut consumer).await.unwrap();
1506 assert_eq!(frames.len(), 1);
1507
1508 use bytes::Buf;
1509 let mut received = Vec::new();
1510 let mut payload = frames[0].payload.clone();
1511 while payload.has_remaining() {
1512 received.push(payload.get_u8());
1513 }
1514 assert_eq!(received, payload_bytes);
1515 }
1516
1517 #[tokio::test]
1520 async fn no_infinite_loop_with_buffered_frames() {
1521 tokio::time::pause();
1522 let mut track = track_producer("test", hang::container::track_info());
1523 let consumer_track = track.subscribe(None);
1524 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_secs(10));
1525
1526 let mut group0 = track.create_group(moq_net::group::Info { sequence: 0 }).unwrap();
1527 Container::Legacy
1528 .write(
1529 &mut group0,
1530 &[Frame {
1531 timestamp: ts(0),
1532 payload: Bytes::from_static(&[0xDE, 0xAD]),
1533 keyframe: false,
1534 duration: None,
1535 }],
1536 )
1537 .unwrap();
1538
1539 write_group(&mut track, 1, &[ts(100_000)]);
1540
1541 let finisher = tokio::spawn(async move {
1542 tokio::time::sleep(Duration::from_millis(20)).await;
1543 write_group(&mut track, 2, &[ts(200_000)]);
1545 tokio::time::sleep(Duration::from_millis(20)).await;
1546 group0.finish().unwrap();
1547 track.finish().unwrap();
1548 });
1549
1550 let frames = tokio::time::timeout(Duration::from_secs(2), async {
1551 let mut frames = Vec::new();
1552 while let Some(frame) = consumer.read().await.unwrap() {
1553 frames.push(frame);
1554 }
1555 frames
1556 })
1557 .await
1558 .expect("consumer hung — possible infinite loop regression");
1559
1560 assert_eq!(frames.len(), 3);
1561 finisher.await.expect("finisher task panicked");
1562 }
1563
1564 #[tokio::test]
1567 async fn large_timestamps() {
1568 let mut track = track_producer("test", hang::container::track_info());
1569 let consumer_track = track.subscribe(None);
1570 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_secs(3700));
1571
1572 let one_hour = 3_600_000_000u64;
1573 write_group(&mut track, 0, &[ts(one_hour)]);
1574 track.finish().unwrap();
1575
1576 let frames = read_all(&mut consumer).await.unwrap();
1577 assert_eq!(frames.len(), 1);
1578 assert_eq!(frames[0].timestamp, ts(one_hour));
1579 assert_eq!(frames[0].timestamp.as_micros(), one_hour as u128);
1580 }
1581
1582 #[tokio::test]
1583 async fn set_latency_changes_behavior() {
1584 let mut track = track_producer("test", hang::container::track_info());
1585 let consumer_track = track.subscribe(None);
1586 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_secs(10));
1587
1588 write_group(&mut track, 0, &[ts(0)]);
1589 track.finish().unwrap();
1590
1591 let frame = consumer.read().await.unwrap().unwrap();
1592 assert_eq!(frame.timestamp, ts(0));
1593
1594 consumer.set_latency(Duration::from_millis(100));
1595
1596 assert!(consumer.read().await.unwrap().is_none());
1597 }
1598
1599 #[tokio::test]
1600 async fn max_timestamp_tracks_through_bframes() {
1601 tokio::time::pause();
1602 let mut track = track_producer("test", hang::container::track_info());
1603 let consumer_track = track.subscribe(None);
1604 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(110));
1607
1608 let mut group0 = track.create_group(moq_net::group::Info { sequence: 0 }).unwrap();
1609 for ×tamp in &[ts(0), ts(66_000), ts(33_000)] {
1610 Container::Legacy
1611 .write(
1612 &mut group0,
1613 &[Frame {
1614 timestamp,
1615 payload: Bytes::from_static(&[0xDE, 0xAD]),
1616 keyframe: false,
1617 duration: None,
1618 }],
1619 )
1620 .unwrap();
1621 }
1622
1623 write_group(&mut track, 1, &[ts(100_000)]);
1624 track.finish().unwrap();
1625
1626 let finisher = tokio::spawn(async move {
1627 tokio::time::sleep(Duration::from_millis(50)).await;
1628 group0.finish().unwrap();
1629 });
1630
1631 let frames = tokio::time::timeout(Duration::from_secs(2), async {
1632 let mut frames = Vec::new();
1633 while let Some(frame) = consumer.read().await.unwrap() {
1634 frames.push(frame);
1635 }
1636 frames
1637 })
1638 .await
1639 .expect("consumer hung — max_timestamp regression");
1640
1641 assert_eq!(frames.len(), 4, "Expected all 4 frames, got {}", frames.len());
1642 assert_eq!(frames[0].timestamp, ts(0));
1643 assert_eq!(frames[1].timestamp, ts(66_000));
1644 assert_eq!(frames[2].timestamp, ts(33_000));
1645 assert_eq!(frames[3].timestamp, ts(100_000));
1646 finisher.await.expect("finisher task panicked");
1647 }
1648
1649 #[tokio::test]
1652 async fn startup_selects_earliest_group() {
1653 tokio::time::pause();
1654 let mut track = track_producer("test", hang::container::track_info());
1655 let consumer_track = track.subscribe(None);
1656 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(100));
1657
1658 write_group(&mut track, 3, &[ts(0)]);
1659 write_group(&mut track, 5, &[ts(150_000)]);
1660
1661 let mut group7 = track.create_group(moq_net::group::Info { sequence: 7 }).unwrap();
1662 Container::Legacy
1663 .write(
1664 &mut group7,
1665 &[Frame {
1666 timestamp: ts(300_000),
1667 payload: Bytes::from_static(&[0xDE, 0xAD]),
1668 keyframe: false,
1669 duration: None,
1670 }],
1671 )
1672 .unwrap();
1673
1674 let finisher = tokio::spawn(async move {
1675 tokio::time::sleep(Duration::from_millis(50)).await;
1676 Container::Legacy
1677 .write(
1678 &mut group7,
1679 &[Frame {
1680 timestamp: ts(400_000),
1681 payload: Bytes::from_static(&[0xBE, 0xEF]),
1682 keyframe: false,
1683 duration: None,
1684 }],
1685 )
1686 .unwrap();
1687 group7.finish().unwrap();
1688 track.finish().unwrap();
1689 });
1690
1691 let _frames = tokio::time::timeout(Duration::from_secs(2), async {
1692 let mut frames = Vec::new();
1693 while let Some(frame) = consumer.read().await.unwrap() {
1694 frames.push(frame);
1695 }
1696 frames
1697 })
1698 .await
1699 .expect("should not hang");
1700
1701 finisher.await.unwrap();
1702 }
1703
1704 #[tokio::test]
1705 async fn startup_skips_groups_without_data() {
1706 tokio::time::pause();
1707 let mut track = track_producer("test", hang::container::track_info());
1708 let consumer_track = track.subscribe(None);
1709 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(500));
1710
1711 let _group5 = track.create_group(moq_net::group::Info { sequence: 5 }).unwrap();
1712 write_group(&mut track, 7, &[ts(210_000)]);
1713 track.finish().unwrap();
1714
1715 let frames = tokio::time::timeout(Duration::from_millis(500), async {
1716 let mut frames = Vec::new();
1717 while let Some(frame) = consumer.read().await.unwrap() {
1718 frames.push(frame);
1719 }
1720 frames
1721 })
1722 .await
1723 .expect("should not hang");
1724
1725 assert!(!frames.is_empty());
1726 }
1727
1728 #[tokio::test]
1729 async fn startup_single_group_mid_stream() {
1730 let mut track = track_producer("test", hang::container::track_info());
1731 let consumer_track = track.subscribe(None);
1732 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(500));
1733
1734 write_group(&mut track, 100, &[ts(3_000_000)]);
1735 track.finish().unwrap();
1736
1737 let frames = read_all(&mut consumer).await.unwrap();
1738 assert_eq!(frames.len(), 1);
1739 }
1740
1741 #[tokio::test]
1742 async fn multiple_sequential_latency_skips() {
1743 tokio::time::pause();
1744 let mut track = track_producer("test", hang::container::track_info());
1745 let consumer_track = track.subscribe(None);
1746 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(50));
1747
1748 let mut group0 = track.create_group(moq_net::group::Info { sequence: 0 }).unwrap();
1749 Container::Legacy
1750 .write(
1751 &mut group0,
1752 &[Frame {
1753 timestamp: ts(0),
1754 payload: Bytes::from_static(&[0xAA]),
1755
1756 keyframe: false,
1757 duration: None,
1758 }],
1759 )
1760 .unwrap();
1761
1762 write_group(&mut track, 1, &[ts(100_000)]);
1763 write_group(&mut track, 2, &[ts(200_000)]);
1764 write_group(&mut track, 3, &[ts(300_000)]);
1765 track.finish().unwrap();
1766
1767 let finisher = tokio::spawn(async move {
1768 tokio::time::sleep(Duration::from_millis(20)).await;
1769 group0.finish().unwrap();
1770 });
1771
1772 let frames = read_all(&mut consumer).await.unwrap();
1773 assert!(!frames.is_empty());
1774 finisher.await.unwrap();
1775 }
1776
1777 #[tokio::test]
1778 async fn latency_skip_boundary_exact() {
1779 tokio::time::pause();
1780 let mut track = track_producer("test", hang::container::track_info());
1781 let consumer_track = track.subscribe(None);
1782 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(100));
1783
1784 let mut group0 = track.create_group(moq_net::group::Info { sequence: 0 }).unwrap();
1785 Container::Legacy
1786 .write(
1787 &mut group0,
1788 &[Frame {
1789 timestamp: ts(0),
1790 payload: Bytes::from_static(&[0xAA]),
1791
1792 keyframe: false,
1793 duration: None,
1794 }],
1795 )
1796 .unwrap();
1797
1798 write_group(&mut track, 1, &[ts(100_000)]);
1799 track.finish().unwrap();
1800
1801 let finisher = tokio::spawn(async move {
1802 tokio::time::sleep(Duration::from_millis(20)).await;
1803 group0.finish().unwrap();
1804 });
1805
1806 let frames = read_all(&mut consumer).await.unwrap();
1807 assert!(!frames.is_empty());
1808 finisher.await.unwrap();
1809 }
1810
1811 #[tokio::test]
1816 async fn single_newer_group_triggers_skip() {
1817 tokio::time::pause();
1818 let mut track = track_producer("test", hang::container::track_info());
1819 let consumer_track = track.subscribe(None);
1820 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(100));
1821
1822 let mut group0 = track.create_group(moq_net::group::Info { sequence: 0 }).unwrap();
1824 Container::Legacy
1825 .write(
1826 &mut group0,
1827 &[Frame {
1828 timestamp: ts(0),
1829 payload: Bytes::from_static(&[0xDE, 0xAD]),
1830 keyframe: false,
1831 duration: None,
1832 }],
1833 )
1834 .unwrap();
1835
1836 write_group(&mut track, 1, &[ts(200_000)]);
1838 track.finish().unwrap();
1839
1840 let finisher = tokio::spawn(async move {
1841 tokio::time::sleep(Duration::from_millis(50)).await;
1842 group0.finish().unwrap();
1843 });
1844
1845 let frames = read_all(&mut consumer).await.unwrap();
1846 assert_eq!(frames.len(), 2, "Expected group 0 frame + group 1 frame");
1847 finisher.await.unwrap();
1848 }
1849
1850 #[tokio::test]
1854 async fn single_missing_sequence_near_eof_skips() {
1855 tokio::time::pause();
1856 let mut track = track_producer("test", hang::container::track_info());
1857 let consumer_track = track.subscribe(None);
1858 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(100));
1859
1860 write_group(&mut track, 0, &[ts(0), ts(20_000)]);
1862 write_group(&mut track, 2, &[ts(200_000)]);
1864 track.finish().unwrap();
1865
1866 let frames = read_all(&mut consumer).await.unwrap();
1867 assert_eq!(frames.len(), 3, "Expected group 0 (2 frames) + group 2 (1 frame)");
1868 }
1869
1870 #[tokio::test]
1871 async fn group_error_skips_to_next() {
1872 let mut track = track_producer("test", hang::container::track_info());
1873 let consumer_track = track.subscribe(None);
1874 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(500));
1875
1876 let group0 = track.create_group(moq_net::group::Info { sequence: 0 }).unwrap();
1877 group0.abort(moq_net::Error::Cancel).unwrap();
1878
1879 write_group(&mut track, 1, &[ts(30_000)]);
1880 track.finish().unwrap();
1881
1882 let frames = read_all(&mut consumer).await.unwrap();
1883 assert_eq!(frames.len(), 1);
1884 }
1885
1886 #[tokio::test]
1887 async fn track_finishes_while_reading() {
1888 tokio::time::pause();
1889 let mut track = track_producer("test", hang::container::track_info());
1890 let consumer_track = track.subscribe(None);
1891 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(500));
1892
1893 write_group(&mut track, 0, &[ts(0)]);
1894
1895 let finisher = tokio::spawn(async move {
1896 tokio::time::sleep(Duration::from_millis(20)).await;
1897 write_group(&mut track, 1, &[ts(30_000)]);
1898 tokio::time::sleep(Duration::from_millis(20)).await;
1899 track.finish().unwrap();
1900 });
1901
1902 let frames = tokio::time::timeout(Duration::from_secs(2), async {
1903 let mut frames = Vec::new();
1904 while let Some(frame) = consumer.read().await.unwrap() {
1905 frames.push(frame);
1906 }
1907 frames
1908 })
1909 .await
1910 .expect("consumer should not hang");
1911
1912 assert_eq!(frames.len(), 2);
1913 finisher.await.unwrap();
1914 }
1915
1916 #[tokio::test]
1917 async fn empty_group_advances() {
1918 let mut track = track_producer("test", hang::container::track_info());
1919 let consumer_track = track.subscribe(None);
1920 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(500));
1921
1922 let mut group0 = track.create_group(moq_net::group::Info { sequence: 0 }).unwrap();
1923 group0.finish().unwrap();
1924
1925 write_group(&mut track, 1, &[ts(30_000)]);
1926 track.finish().unwrap();
1927
1928 let frames = read_all(&mut consumer).await.unwrap();
1929 assert_eq!(frames.len(), 1);
1930 }
1931
1932 #[tokio::test]
1935 async fn video_container_legacy() {
1936 tokio::time::pause();
1937
1938 let mut track = track_producer("video", hang::container::track_info());
1939 let consumer_track = track.subscribe(None);
1940 let mut consumer = Consumer::new(consumer_track, Container::Legacy).with_latency(Duration::from_millis(500));
1941
1942 let mut group = track.create_group(moq_net::group::Info { sequence: 0 }).unwrap();
1944 for i in 0..3u64 {
1945 let frame = Frame {
1946 timestamp: ts(i * 33_333),
1947 payload: Bytes::from_static(&[0xDE, 0xAD]),
1948 keyframe: false,
1949 duration: None,
1950 };
1951 Container::Legacy.write(&mut group, &[frame]).unwrap();
1952 }
1953 group.finish().unwrap();
1954 track.finish().unwrap();
1955
1956 let mut frames = Vec::new();
1957 while let Some(frame) = consumer.read().await.unwrap() {
1958 frames.push(frame);
1959 }
1960
1961 assert_eq!(frames.len(), 3);
1962 assert_eq!(frames[0].timestamp, ts(0));
1963 assert!(frames[0].keyframe);
1964 assert_eq!(frames[1].timestamp, ts(33_333));
1965 assert!(!frames[1].keyframe);
1966 assert_eq!(frames[2].timestamp, ts(66_666));
1967 assert!(!frames[2].keyframe);
1968 }
1969
1970 #[tokio::test]
1976 async fn duration_skip_advances_to_next_group() {
1977 tokio::time::pause();
1978 let mut track = track_producer("test", None);
1981 let consumer_track = track.subscribe(None);
1982 let mut consumer = Consumer::new(consumer_track, DurationWire).with_latency(Duration::from_secs(10));
1984
1985 let mut group0 = track.create_group(moq_net::group::Info { sequence: 0 }).unwrap();
1987 write_duration_frame(&mut group0, ts(0), ts(33_000));
1988
1989 let mut group1 = track.create_group(moq_net::group::Info { sequence: 1 }).unwrap();
1991 write_duration_frame(&mut group1, ts(33_000), ts(33_000));
1992 group1.finish().unwrap();
1993
1994 track.finish().unwrap();
1995
1996 let frames = tokio::time::timeout(Duration::from_secs(2), async {
1997 let mut frames = Vec::new();
1998 while let Some(frame) = consumer.read().await.unwrap() {
1999 frames.push(frame);
2000 }
2001 frames
2002 })
2003 .await
2004 .expect("consumer hung — duration skip regression");
2005
2006 assert_eq!(frames.len(), 2);
2007 assert_eq!(frames[0].timestamp, ts(0));
2008 assert_eq!(frames[1].timestamp, ts(33_000));
2009
2010 drop(group0);
2012 }
2013
2014 #[tokio::test]
2018 async fn duration_below_gap_does_not_skip() {
2019 tokio::time::pause();
2020 let mut track = track_producer("test", None);
2022 let consumer_track = track.subscribe(None);
2023 let mut consumer = Consumer::new(consumer_track, DurationWire).with_latency(Duration::from_secs(10));
2024
2025 let mut group0 = track.create_group(moq_net::group::Info { sequence: 0 }).unwrap();
2027 write_duration_frame(&mut group0, ts(0), ts(10_000));
2028
2029 let mut group1 = track.create_group(moq_net::group::Info { sequence: 1 }).unwrap();
2031 write_duration_frame(&mut group1, ts(33_000), ts(33_000));
2032 group1.finish().unwrap();
2033 track.finish().unwrap();
2034
2035 let finisher = tokio::spawn(async move {
2038 tokio::time::sleep(Duration::from_millis(20)).await;
2039 write_duration_frame(&mut group0, ts(20_000), ts(10_000));
2040 group0.finish().unwrap();
2041 });
2042
2043 let frames = tokio::time::timeout(Duration::from_secs(2), async {
2044 let mut frames = Vec::new();
2045 while let Some(frame) = consumer.read().await.unwrap() {
2046 frames.push(frame);
2047 }
2048 frames
2049 })
2050 .await
2051 .expect("consumer hung");
2052
2053 assert_eq!(frames.len(), 3);
2055 assert_eq!(frames[0].timestamp, ts(0));
2056 assert_eq!(frames[1].timestamp, ts(20_000));
2057 assert_eq!(frames[2].timestamp, ts(33_000));
2058 finisher.await.unwrap();
2059 }
2060}