moq_binary/snapshot/
mod.rs1use crate::Compression;
15
16pub mod consumer;
17pub mod producer;
18
19pub use consumer::Consumer;
20pub use producer::Producer;
21
22#[derive(Debug, Clone, Default)]
27#[non_exhaustive]
28pub struct Config {
29 pub compression: Compression,
35}
36
37#[cfg(test)]
38mod test {
39 use std::task::Poll;
40
41 use bytes::Bytes;
42
43 use super::*;
44
45 fn cfg(compression: bool) -> Config {
46 Config {
47 compression: if compression {
48 Compression::Deflate
49 } else {
50 Compression::None
51 },
52 }
53 }
54
55 fn producer(compression: bool) -> (Producer, moq_net::track::Subscriber) {
56 let track = moq_net::broadcast::Info::new()
57 .produce()
58 .create_track("test", None)
59 .unwrap();
60 let consumer = track.subscribe(None);
61 (Producer::new(track, cfg(compression)), consumer)
62 }
63
64 fn consume(track: moq_net::track::Subscriber, compression: bool) -> Consumer {
65 Consumer::new(track, cfg(compression))
66 }
67
68 fn drain(mut consumer: Consumer) -> Vec<Bytes> {
70 let waiter = kio::Waiter::noop();
71 let mut out = Vec::new();
72 while let Poll::Ready(Ok(Some(payload))) = consumer.poll_next(&waiter) {
73 out.push(payload);
74 }
75 out
76 }
77
78 #[test]
83 fn a_lost_group_waits_for_its_replacement() {
84 let track = moq_net::broadcast::Info::new()
85 .produce()
86 .create_track("test", None)
87 .unwrap();
88 let mut consumer = consume(track.subscribe(None), false);
89 let waiter = kio::Waiter::noop();
90
91 let mut group = track.create_group(moq_net::group::Info { sequence: 0 }).unwrap();
93 group
94 .write_frame(moq_net::Timestamp::now(), Bytes::from_static(b"one"))
95 .unwrap();
96 assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(v))) if v == "one"));
97 assert!(consumer.poll_next(&waiter).is_pending());
98
99 group.abort(moq_net::Error::Old).unwrap();
101 assert!(
102 consumer.poll_next(&waiter).is_pending(),
103 "a lost group must not end the reader"
104 );
105
106 let mut group = track.create_group(moq_net::group::Info { sequence: 1 }).unwrap();
108 group
109 .write_frame(moq_net::Timestamp::now(), Bytes::from_static(b"two"))
110 .unwrap();
111 assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(v))) if v == "two"));
112 }
113
114 #[test]
117 fn a_stamped_update_writes_its_capture_time() {
118 let (mut producer, track) = producer(true);
119 let mut groups = producer.consume();
120 let captured = moq_net::Timestamp::from_millis(1_234).unwrap();
121 let payload = Bytes::from(vec![7u8; 4096]);
122 let size = producer
123 .update(moq_net::Timed::from(payload.clone()).at(captured))
124 .unwrap();
125
126 let waiter = kio::Waiter::noop();
127 let Poll::Ready(Ok(Some(mut group))) = groups.poll_recv_group(&waiter) else {
128 panic!("expected a group");
129 };
130 let Poll::Ready(Ok(Some(frame))) = group.poll_read_frame(&waiter) else {
131 panic!("expected a frame");
132 };
133 assert_eq!(frame.timestamp.as_micros(), captured.as_micros());
134 assert_eq!(size, frame.payload.len());
135 assert!(size < payload.len(), "the size is the compressed frame");
136 assert_eq!(drain(consume(track, true)), vec![payload]);
137 }
138
139 #[test]
140 fn one_group_per_update() {
141 let (mut producer, track) = producer(false);
142 producer.update(&b"first"[..]).unwrap();
143 producer.update(&b"second"[..]).unwrap();
144 producer.finish().unwrap();
145
146 assert_eq!(track.latest(), Some(1));
148 assert_eq!(drain(consume(track, false)), vec![Bytes::from_static(b"second")]);
149 }
150
151 #[test]
152 fn live_consumer_sees_each_update() {
153 let (mut producer, track) = producer(false);
154 let mut consumer = consume(track, false);
155 let waiter = kio::Waiter::noop();
156
157 for n in 0..3u8 {
158 producer.update(vec![n]).unwrap();
159 match consumer.poll_next(&waiter) {
160 Poll::Ready(Ok(Some(payload))) => assert_eq!(&payload[..], &[n]),
161 other => panic!("expected value, got {other:?}"),
162 }
163 }
164 }
165
166 #[test]
167 fn compressed_roundtrip() {
168 let (mut producer, track) = producer(true);
169 let payload = Bytes::from(b"the quick brown fox".repeat(64));
170 producer.update(payload.clone()).unwrap();
171 producer.finish().unwrap();
172
173 assert_eq!(drain(consume(track, true)), vec![payload]);
174 }
175
176 #[test]
179 fn compression_shrinks_the_frame() {
180 let (mut producer, track) = producer(true);
181 let payload = Bytes::from(b"the quick brown fox".repeat(64));
182 producer.update(payload.clone()).unwrap();
183 producer.finish().unwrap();
184
185 let waiter = kio::Waiter::noop();
187 let Poll::Ready(Ok(Some(mut group))) = track.ordered().poll_next_group(&waiter) else {
188 panic!("expected a group");
189 };
190 let Poll::Ready(Ok(Some(frame))) = group.poll_read_frame(&waiter) else {
191 panic!("expected a frame");
192 };
193 assert!(
194 frame.payload.len() < payload.len() / 4,
195 "compressed frame {} should be far below the raw {}",
196 frame.payload.len(),
197 payload.len()
198 );
199 }
200
201 #[test]
205 fn a_rejected_update_leaves_the_previous_value_readable() {
206 let (mut producer, track) = producer(false);
207 producer.update(&b"keep"[..]).unwrap();
208
209 let oversized = Bytes::from(vec![0u8; moq_net::group::MAX_CACHE_BYTES as usize + 1]);
210 assert!(producer.update(oversized).is_err());
211 producer.finish().unwrap();
212
213 assert_eq!(drain(consume(track, false)), vec![Bytes::from_static(b"keep")]);
215 }
216
217 #[test]
220 fn updating_after_finish_fails_on_every_clone() {
221 let (mut producer, _track) = producer(false);
222 let mut clone = producer.clone();
223
224 producer.update(&b"first"[..]).unwrap();
225 producer.finish().unwrap();
226
227 assert!(producer.update(&b"late"[..]).is_err());
228 assert!(clone.update(&b"late"[..]).is_err());
229 }
230
231 #[test]
232 fn a_finished_track_ends_the_consumer() {
233 let (mut producer, track) = producer(false);
234 producer.update(&b"only"[..]).unwrap();
235 producer.finish().unwrap();
236
237 let mut consumer = consume(track, false);
238 let waiter = kio::Waiter::noop();
239 assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(_)))));
240 assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(None))));
241 }
242}