1use crate::Compression;
21
22pub mod consumer;
23pub mod producer;
24
25pub use consumer::Consumer;
26pub use producer::Producer;
27
28#[derive(Debug, Clone, Default)]
33#[non_exhaustive]
34pub struct Config {
35 pub compression: Compression,
41}
42
43#[cfg(test)]
44mod test {
45 use std::task::Poll;
46
47 use bytes::Bytes;
48
49 use super::*;
50
51 fn cfg(compression: bool) -> Config {
52 Config {
53 compression: if compression {
54 Compression::Deflate
55 } else {
56 Compression::None
57 },
58 }
59 }
60
61 fn producer(compression: bool) -> (Producer, moq_net::track::Subscriber) {
62 let track = moq_net::broadcast::Info::new()
63 .produce()
64 .create_track("test", None)
65 .unwrap();
66 let consumer = track.subscribe(None);
67 (Producer::new(track, cfg(compression)), consumer)
68 }
69
70 fn drain(track: moq_net::track::Subscriber, compression: bool) -> Vec<Bytes> {
72 let mut consumer = Consumer::new(track, cfg(compression));
73 let waiter = kio::Waiter::noop();
74 let mut out = Vec::new();
75 while let Poll::Ready(Ok(Some(payload))) = consumer.poll_next(&waiter) {
76 out.push(payload);
77 }
78 out
79 }
80
81 fn payloads(count: u8) -> Vec<Bytes> {
82 (0..count).map(|n| Bytes::from(vec![n; 16])).collect()
83 }
84
85 #[test]
86 fn every_payload_survives_in_order() {
87 let (mut producer, track) = producer(false);
88 let expected = payloads(5);
89 for payload in &expected {
90 producer.append(payload.clone()).unwrap();
91 }
92 producer.finish().unwrap();
93
94 assert_eq!(track.latest(), Some(0));
96 assert_eq!(drain(track, false), expected);
97 }
98
99 #[test]
101 fn a_stamped_append_writes_its_capture_time() {
102 let (mut producer, _track) = producer(true);
103 let mut groups = producer.consume();
104 let captured = moq_net::Timestamp::from_millis(1_234).unwrap();
105 let first = producer
106 .append(moq_net::Timed::from(vec![1u8; 4096]).at(captured))
107 .unwrap();
108 producer.append(vec![2u8; 16]).unwrap();
109
110 let waiter = kio::Waiter::noop();
111 let Poll::Ready(Ok(Some(mut group))) = groups.poll_recv_group(&waiter) else {
112 panic!("expected a group");
113 };
114 let Poll::Ready(Ok(Some(frame))) = group.poll_read_frame(&waiter) else {
115 panic!("expected a frame");
116 };
117 assert_eq!(frame.timestamp.as_micros(), captured.as_micros());
118 assert_eq!(first, frame.payload.len(), "the size is the compressed frame");
119 let Poll::Ready(Ok(Some(frame))) = group.poll_read_frame(&waiter) else {
120 panic!("expected a frame");
121 };
122 assert_ne!(frame.timestamp.as_micros(), captured.as_micros());
123 }
124
125 #[test]
126 fn compressed_roundtrip_in_order() {
127 let (mut producer, track) = producer(true);
128 let expected = payloads(20);
129 for payload in &expected {
130 producer.append(payload.clone()).unwrap();
131 }
132 producer.finish().unwrap();
133
134 assert_eq!(drain(track, true), expected);
135 }
136
137 #[test]
139 fn the_shared_window_shrinks_repetitive_payloads() {
140 let (mut producer, track) = producer(true);
141 let payload = Bytes::from(b"the quick brown fox".repeat(16));
142 for _ in 0..8 {
143 producer.append(payload.clone()).unwrap();
144 }
145 producer.finish().unwrap();
146
147 let waiter = kio::Waiter::noop();
148 let Poll::Ready(Ok(Some(mut group))) = track.ordered().poll_next_group(&waiter) else {
149 panic!("expected a group");
150 };
151
152 let mut sizes = Vec::new();
153 while let Poll::Ready(Ok(Some(frame))) = group.poll_read_frame(&waiter) {
154 sizes.push(frame.payload.len());
155 }
156
157 assert_eq!(sizes.len(), 8);
158 assert!(
159 *sizes.last().unwrap() < sizes[0],
160 "windowed frame {} should be below the first {}",
161 sizes.last().unwrap(),
162 sizes[0]
163 );
164 }
165
166 #[test]
168 fn clones_append_into_one_log() {
169 let (mut producer, track) = producer(true);
170 let mut clone = producer.clone();
171
172 producer.append(&b"a"[..]).unwrap();
173 clone.append(&b"b"[..]).unwrap();
174 producer.append(&b"c"[..]).unwrap();
175 producer.finish().unwrap();
176
177 assert_eq!(
178 drain(track, true),
179 vec![
180 Bytes::from_static(b"a"),
181 Bytes::from_static(b"b"),
182 Bytes::from_static(b"c")
183 ]
184 );
185 }
186
187 fn rejecting_track() -> moq_net::track::Producer {
192 let mut info = moq_net::track::Info::default();
193 info.timescale = moq_net::Timescale::new((1u64 << 62) - 1).unwrap();
194
195 moq_net::broadcast::Info::new()
196 .produce()
197 .create_track("test", Some(info))
198 .unwrap()
199 }
200
201 #[test]
205 fn a_failed_write_aborts_the_track() {
206 let track = rejecting_track();
207 let mut subscriber = track.subscribe(None);
208 let mut producer = Producer::new(track, Config::default());
209
210 assert!(producer.append(&b"rejected"[..]).is_err());
211
212 let waiter = kio::Waiter::noop();
215 assert!(
216 matches!(subscriber.poll_recv_group(&waiter), Poll::Ready(Err(_))),
217 "a truncated log must surface an error rather than read as a completed one"
218 );
219 }
220
221 #[test]
225 fn an_undecodable_record_ends_the_track() {
226 let track = moq_net::broadcast::Info::new()
227 .produce()
228 .create_track("test", None)
229 .unwrap();
230 let mut subscriber = track.subscribe(None);
231 let mut producer = Producer::new(track, cfg(true));
232
233 assert!(producer.is_used());
234 let oversized = Bytes::from(vec![0u8; moq_flate::DEFAULT_MAX_FRAME_SIZE as usize + 1]);
235 assert!(matches!(
236 producer.append(oversized),
237 Err(crate::Error::Flate(moq_flate::Error::TooLarge(_)))
238 ));
239
240 assert!(!producer.is_used());
241
242 let waiter = kio::Waiter::noop();
244 assert!(matches!(subscriber.poll_recv_group(&waiter), Poll::Ready(Err(_))));
245 assert!(producer.append(&b"after"[..]).is_err());
246 }
247
248 #[test]
252 fn a_reader_inside_the_group_sees_the_real_error() {
253 let track = moq_net::broadcast::Info::new()
254 .produce()
255 .create_track("test", None)
256 .unwrap();
257 let mut subscriber = track.subscribe(None);
258 let mut producer = Producer::new(track, Config::default());
259
260 producer.append(&b"first"[..]).unwrap();
261
262 let waiter = kio::Waiter::noop();
264 let Poll::Ready(Ok(Some(mut group))) = subscriber.poll_recv_group(&waiter) else {
265 panic!("the group was published, so a subscriber sees it");
266 };
267 assert!(matches!(group.poll_read_frame(&waiter), Poll::Ready(Ok(Some(_)))));
268
269 let oversized = Bytes::from(vec![0u8; moq_net::group::MAX_CACHE_BYTES as usize + 1]);
271 assert!(producer.append(oversized).is_err());
272
273 match group.poll_read_frame(&waiter) {
274 Poll::Ready(Err(err)) => assert!(
275 matches!(err, moq_net::Error::FrameTooLarge),
276 "the reader should see the write's own error, got {err:?}"
277 ),
278 other => panic!("expected the write error, got {other:?}"),
279 }
280 }
281
282 fn replaying() -> moq_net::track::Subscription {
286 moq_net::track::Subscription::default().with_max_age(std::time::Duration::from_secs(30))
287 }
288
289 #[test]
291 fn a_failed_write_ends_the_track() {
292 let track = rejecting_track();
293 let mut producer = Producer::new(track, Config::default());
294
295 assert!(producer.append(&b"rejected"[..]).is_err());
296 assert!(
297 producer.append(&b"again"[..]).is_err(),
298 "a second append must fail on the closed track rather than open another group"
299 );
300 }
301
302 #[test]
306 fn a_second_group_is_a_rolled_log() {
307 let track = moq_net::broadcast::Info::new()
308 .produce()
309 .create_track("test", None)
310 .unwrap();
311 let subscriber = track.subscribe(replaying());
312
313 for pair in payloads(4).chunks(2) {
314 let mut flate = moq_flate::Encoder::new();
316 let mut group = track.append_group().unwrap();
317 for payload in pair {
318 group
319 .write_frame(moq_net::Timestamp::now(), flate.frame(payload))
320 .unwrap();
321 }
322 group.finish().unwrap();
323 }
324 track.finish().unwrap();
325
326 let mut consumer = Consumer::new(subscriber, cfg(true));
327 let waiter = kio::Waiter::noop();
328
329 assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(_)))));
331 assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(_)))));
332
333 assert!(matches!(
335 consumer.poll_next(&waiter),
336 Poll::Ready(Err(crate::Error::Rolled))
337 ));
338 }
339
340 #[test]
344 fn a_second_group_is_reported_while_the_first_is_open() {
345 let track = moq_net::broadcast::Info::new()
346 .produce()
347 .create_track("test", None)
348 .unwrap();
349 let subscriber = track.subscribe(replaying());
350
351 let mut first = track.append_group().unwrap();
353 first.write_frame(moq_net::Timestamp::now(), &b"first"[..]).unwrap();
354 let mut second = track.append_group().unwrap();
355 second.write_frame(moq_net::Timestamp::now(), &b"second"[..]).unwrap();
356
357 let mut consumer = Consumer::new(subscriber, Config::default());
358 let waiter = kio::Waiter::noop();
359
360 assert!(matches!(
361 consumer.poll_next(&waiter),
362 Poll::Ready(Ok(Some(payload))) if payload == Bytes::from_static(b"first")
363 ));
364 assert!(matches!(
365 consumer.poll_next(&waiter),
366 Poll::Ready(Err(crate::Error::Rolled))
367 ));
368
369 first.write_frame(moq_net::Timestamp::now(), &b"more"[..]).unwrap();
371 assert!(matches!(
372 consumer.poll_next(&waiter),
373 Poll::Ready(Err(crate::Error::Rolled))
374 ));
375 }
376
377 #[test]
381 fn a_late_lower_group_is_still_reported() {
382 let track = moq_net::broadcast::Info::new()
383 .produce()
384 .create_track("test", None)
385 .unwrap();
386 let subscriber = track.subscribe(replaying());
387
388 for sequence in [1u64, 0] {
390 let mut flate = moq_flate::Encoder::new();
391 let mut group = track.create_group(moq_net::group::Info { sequence }).unwrap();
392 group
393 .write_frame(moq_net::Timestamp::now(), flate.frame(&[sequence as u8; 8]))
394 .unwrap();
395 group.finish().unwrap();
396 }
397 track.finish().unwrap();
398
399 let mut consumer = Consumer::new(subscriber, cfg(true));
400 let waiter = kio::Waiter::noop();
401
402 assert!(matches!(
403 consumer.poll_next(&waiter),
404 Poll::Ready(Ok(Some(payload))) if payload == vec![1u8; 8]
405 ));
406 assert!(matches!(
407 consumer.poll_next(&waiter),
408 Poll::Ready(Err(crate::Error::Rolled))
409 ));
410 }
411
412 #[test]
415 fn appending_after_finish_fails_on_every_clone() {
416 let (mut producer, _track) = producer(false);
417 let mut clone = producer.clone();
418
419 producer.append(&b"first"[..]).unwrap();
420 producer.finish().unwrap();
421
422 assert!(producer.append(&b"late"[..]).is_err());
423 assert!(clone.append(&b"late"[..]).is_err());
424 }
425
426 #[test]
427 fn a_finished_track_ends_the_consumer() {
428 let (mut producer, track) = producer(false);
429 producer.append(&b"only"[..]).unwrap();
430 producer.finish().unwrap();
431
432 let mut consumer = Consumer::new(track, Config::default());
433 let waiter = kio::Waiter::noop();
434 assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(_)))));
435 assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(None))));
436 }
437}