1mod consumer;
50mod decoder;
51mod encoder;
52mod op;
53mod producer;
54
55pub use consumer::Consumer;
56pub use decoder::{ConsumerConfig, Decoder, Event, Group};
57pub use encoder::{Encoded, Encoder, Pending, ProducerConfig};
58pub use producer::Producer;
59
60#[cfg(test)]
61mod test {
62 use std::task::Poll;
63
64 use serde_json::{Value, json};
65
66 use super::*;
67
68 fn producer(config: ProducerConfig) -> (Producer<Value>, moq_net::track::Subscriber) {
69 let track = moq_net::broadcast::Info::new()
70 .produce()
71 .create_track("test", None)
72 .unwrap();
73 let consumer = track.subscribe(None);
74 (Producer::new(track, config), consumer)
75 }
76
77 #[test]
78 #[should_panic(expected = "checkpoint_records must be positive")]
79 fn zero_checkpoint_records_is_rejected() {
80 let config = ProducerConfig {
81 checkpoint_records: Some(0),
82 ..Default::default()
83 };
84 let _ = Encoder::<Value>::new(config);
85 }
86
87 fn consumer(track: moq_net::track::Subscriber, compression: bool) -> Consumer<Value> {
88 Consumer::new(track, ConsumerConfig::default().with_compression(compression))
89 }
90
91 fn rejecting_track() -> moq_net::track::Producer {
93 let mut info = moq_net::track::Info::default();
94 info.timescale = moq_net::Timescale::new((1u64 << 62) - 1).unwrap();
95
96 moq_net::broadcast::Info::new()
97 .produce()
98 .create_track("test", Some(info))
99 .unwrap()
100 }
101
102 fn drain(consumer: &mut Consumer<Value>) -> Vec<Event<Value>> {
104 let waiter = kio::Waiter::noop();
105 let mut out = Vec::new();
106 while let Poll::Ready(Ok(Some(event))) = consumer.poll_next(&waiter) {
107 out.push(event);
108 }
109 out
110 }
111
112 fn rec(n: u64) -> Value {
113 json!({ "n": n })
114 }
115
116 struct Live {
122 producer: Producer<Value>,
123 consumer: Consumer<Value>,
124 events: Vec<Event<Value>>,
125 }
126
127 impl Live {
128 fn new(config: ProducerConfig) -> Self {
129 let compression = config.compression;
130 let (producer, track) = producer(config);
131 Self {
132 producer,
133 consumer: consumer(track, compression),
134 events: Vec::new(),
135 }
136 }
137
138 fn push(&mut self, n: u64) {
139 self.producer.push(&rec(n)).unwrap();
140 self.read();
141 }
142
143 fn pop(&mut self, count: u64) {
144 self.producer.pop(count).unwrap();
145 self.read();
146 }
147
148 fn read(&mut self) {
149 self.events.extend(drain(&mut self.consumer));
150 }
151
152 fn finish(self) -> Vec<Event<Value>> {
153 let Live {
154 mut producer,
155 mut consumer,
156 mut events,
157 } = self;
158 producer.finish().unwrap();
159 events.extend(drain(&mut consumer));
160 events
161 }
162
163 fn pushed(events: &[Event<Value>]) -> Vec<u64> {
165 events
166 .iter()
167 .filter_map(|e| match e {
168 Event::Push { index, .. } => Some(*index),
169 _ => None,
170 })
171 .collect()
172 }
173 }
174
175 #[test]
176 fn push_and_pop_round_trip() {
177 let mut live = Live::new(ProducerConfig::default());
178 live.push(0);
179 live.push(1);
180 live.pop(1);
181 live.push(2);
182
183 assert_eq!(
184 live.finish(),
185 vec![
186 Event::Push {
187 index: 0,
188 value: rec(0)
189 },
190 Event::Push {
191 index: 1,
192 value: rec(1)
193 },
194 Event::Pop(0..1),
195 Event::Push {
196 index: 2,
197 value: rec(2)
198 },
199 ]
200 );
201 }
202
203 #[test]
204 fn the_window_slides() {
205 let (mut producer, _track) = producer(ProducerConfig::default());
206 for n in 0..5 {
207 producer.push(&rec(n)).unwrap();
208 if n >= 2 {
209 producer.pop(1).unwrap();
210 }
211 }
212
213 assert_eq!(producer.range(), 3..5);
215 assert_eq!(producer.window(), vec![rec(3), rec(4)]);
216 }
217
218 #[test]
219 fn a_popped_record_is_never_restated() {
220 let mut live = Live::new(ProducerConfig::default().with_op_ratio(0));
222 live.push(0);
223 live.push(1);
224 live.pop(1);
225 live.push(2);
226
227 assert_eq!(
230 live.finish(),
231 vec![
232 Event::Push {
233 index: 0,
234 value: rec(0)
235 },
236 Event::Push {
237 index: 1,
238 value: rec(1)
239 },
240 Event::Pop(0..1),
241 Event::Push {
242 index: 2,
243 value: rec(2)
244 },
245 ]
246 );
247 }
248
249 #[test]
250 fn bounded_checkpoints_keep_a_following_consumer_contiguous() {
251 let config = ProducerConfig::default().with_op_ratio(0).with_checkpoint_records(2);
252 let mut live = Live::new(config);
253 for n in 0..6 {
254 live.push(n);
255 }
256
257 assert_eq!(live.producer.range(), 0..6);
258 assert_eq!(live.producer.window(), vec![rec(4), rec(5)]);
259 let events = live.finish();
260 assert_eq!(Live::pushed(&events), (0..6).collect::<Vec<_>>());
261 assert!(!events.iter().any(|event| matches!(event, Event::Skip(_))));
262 }
263
264 #[test]
265 fn a_late_consumer_skips_to_the_bounded_checkpoint() {
266 let config = ProducerConfig::default().with_op_ratio(0).with_checkpoint_records(2);
267 let mut encoder = Encoder::<Value>::new(config);
268 let mut latest = None;
269 for n in 0..5 {
270 let frame = encoder.push(&rec(n)).unwrap();
271 latest = Some(frame.payload.clone());
272 frame.commit();
273 }
274
275 let mut decoder = Decoder::<Value>::new(ConsumerConfig::default());
276 decoder.group().decode(&latest.unwrap()).unwrap();
277 assert_eq!(
278 std::iter::from_fn(|| decoder.next_event()).collect::<Vec<_>>(),
279 vec![
280 Event::Skip(0..3),
281 Event::Push {
282 index: 3,
283 value: rec(3)
284 },
285 Event::Push {
286 index: 4,
287 value: rec(4)
288 },
289 ]
290 );
291 }
292
293 #[test]
294 fn pops_cross_the_omitted_checkpoint_prefix() {
295 let config = ProducerConfig::default().with_checkpoint_records(2);
296 let mut encoder = Encoder::<Value>::new(config);
297 for n in 0..5 {
298 encoder.push(&rec(n)).unwrap().commit();
299 }
300 assert_eq!(encoder.range(), 0..5);
301 assert_eq!(encoder.window(), vec![rec(3), rec(4)]);
302
303 encoder.pop(2).unwrap().unwrap().commit();
304 assert_eq!(encoder.range(), 2..5);
305 assert_eq!(encoder.window(), vec![rec(3), rec(4)]);
306
307 encoder.pop(2).unwrap().unwrap().commit();
308 assert_eq!(encoder.range(), 4..5);
309 assert_eq!(encoder.window(), vec![rec(4)]);
310 }
311
312 #[test]
313 fn a_fresh_consumer_adopts_the_offset_without_skipping_history() {
314 let track = moq_net::broadcast::Info::new()
315 .produce()
316 .create_track("test", None)
317 .unwrap();
318 let mut producer = Producer::<Value>::new(track, ProducerConfig::default().with_op_ratio(0));
319
320 for n in 0..5 {
321 producer.push(&rec(n)).unwrap();
322 }
323 producer.pop(3).unwrap();
324 let mut subscriber = producer.consume();
325 subscriber.set_groups(subscriber.latest().unwrap()..);
326 let mut fresh = consumer(subscriber, false);
327 producer.finish().unwrap();
328
329 let events = drain(&mut fresh);
332 assert_eq!(
333 events,
334 vec![
335 Event::Push {
336 index: 3,
337 value: rec(3)
338 },
339 Event::Push {
340 index: 4,
341 value: rec(4)
342 }
343 ]
344 );
345 assert!(!events.iter().any(|e| matches!(e, Event::Skip(_))));
346 }
347
348 #[test]
349 fn a_lagging_consumer_is_told_what_it_missed() {
350 let mut encoder = Encoder::<Value>::new(ProducerConfig::default().with_op_ratio(0));
354 let mut decoder = Decoder::<Value>::new(ConsumerConfig::default());
355 for n in 0..2 {
356 let frame = encoder.push(&rec(n)).unwrap();
357 let mut group = decoder.group();
358 group.decode(&frame.payload).unwrap();
359 frame.commit();
360 }
361 assert_eq!(
362 std::iter::from_fn(|| decoder.next_event()).collect::<Vec<_>>(),
363 vec![
364 Event::Push {
365 index: 0,
366 value: rec(0)
367 },
368 Event::Push {
369 index: 1,
370 value: rec(1)
371 }
372 ]
373 );
374
375 let mut latest = None;
376 for n in 2..8 {
377 let frame = encoder.push(&rec(n)).unwrap();
378 frame.commit();
379
380 let frame = encoder.pop(1).unwrap().unwrap();
381 latest = Some(frame.payload.clone());
382 frame.commit();
383 }
384 let mut group = decoder.group();
385 group.decode(&latest.unwrap()).unwrap();
386
387 let events = std::iter::from_fn(|| decoder.next_event()).collect::<Vec<_>>();
388 let skipped: Vec<std::ops::Range<u64>> = events
389 .iter()
390 .filter_map(|e| match e {
391 Event::Skip(range) => Some(range.clone()),
392 _ => None,
393 })
394 .collect();
395
396 assert!(!skipped.is_empty(), "expected skips, got {events:?}");
399 assert_eq!(skipped.first().map(|range| range.start), Some(2));
400
401 let reported: Vec<u64> = events
403 .iter()
404 .flat_map(|e| match e {
405 Event::Push { index, .. } => vec![*index],
406 Event::Skip(range) => range.clone().collect(),
407 Event::Pop(_) => Vec::new(),
408 })
409 .collect();
410 assert!(reported.windows(2).all(|w| w[1] == w[0] + 1), "gaps in {reported:?}");
411 }
412
413 #[test]
414 fn consumer_resumes_at_a_checkpoint_after_losing_a_group() {
415 for err in [
416 moq_net::Error::Old,
417 moq_net::Error::Lagged,
418 moq_net::Error::Evicted,
419 moq_net::Error::GroupTooLarge,
420 ] {
421 let track = moq_net::broadcast::Info::new()
422 .produce()
423 .create_track("test", None)
424 .unwrap();
425 let mut consumer = consumer(track.subscribe(None), false);
426 let mut encoder = Encoder::<Value>::new(ProducerConfig::default().with_op_ratio(0));
427
428 let first = encoder.push(&rec(0)).unwrap();
429 let payload = first.payload.clone();
430 first.commit();
431 let mut lost = track.append_group().unwrap();
432 lost.write_frame(moq_net::Timestamp::ZERO, payload).unwrap();
433 assert_eq!(
434 drain(&mut consumer),
435 vec![Event::Push {
436 index: 0,
437 value: rec(0)
438 }]
439 );
440 lost.abort(err).unwrap();
441
442 let second = encoder.push(&rec(1)).unwrap();
443 let payload = second.payload.clone();
444 second.commit();
445 let mut checkpoint = track.append_group().unwrap();
446 checkpoint.write_frame(moq_net::Timestamp::ZERO, payload).unwrap();
447 checkpoint.finish().unwrap();
448 track.finish().unwrap();
449
450 assert_eq!(
451 drain(&mut consumer),
452 vec![Event::Push {
453 index: 1,
454 value: rec(1)
455 }]
456 );
457 }
458 }
459
460 #[test]
461 fn compressed_round_trip_across_rolls() {
462 let mut live = Live::new(ProducerConfig::default().with_compression(true).with_op_ratio(1));
463 for n in 0..40 {
464 live.push(n);
465 if n >= 10 {
466 live.pop(1);
467 }
468 }
469
470 assert_eq!(Live::pushed(&live.finish()), (0..40).collect::<Vec<_>>());
472 }
473
474 #[test]
475 fn an_empty_pop_writes_nothing() {
476 let (mut producer, track) = producer(ProducerConfig::default());
477 producer.pop(5).unwrap();
478 producer.finish().unwrap();
479
480 assert_eq!(track.latest(), None);
482 }
483
484 #[test]
485 fn a_rejected_edit_leaves_the_window_unchanged() {
486 let track = rejecting_track();
487 let mut subscriber = track.subscribe(None).ordered();
488 let mut producer = Producer::<Value>::new(track, ProducerConfig::default());
489
490 assert!(producer.push(&rec(1)).is_err());
491 assert_eq!(producer.range(), 0..0);
492 assert!(producer.window().is_empty());
493
494 let waiter = kio::Waiter::noop();
495 let Poll::Ready(Ok(Some(mut group))) = subscriber.poll_next_group(&waiter) else {
496 panic!("the rejected group's header was published");
497 };
498 assert!(matches!(group.poll_read_frame(&waiter), Poll::Ready(Ok(None))));
499 }
500
501 #[test]
502 fn the_handle_remains_usable_after_finish() {
503 let (mut producer, track) = producer(ProducerConfig::default());
504 producer.push(&rec(1)).unwrap();
505 producer.finish().unwrap();
506
507 assert_eq!(producer.window(), vec![rec(1)]);
508 assert_eq!(producer.range(), 0..1);
509 assert_eq!(producer.consume().latest(), track.latest());
510 producer.finish().unwrap();
511 assert!(matches!(
512 producer.push(&rec(2)),
513 Err(crate::Error::Net(moq_net::Error::Closed))
514 ));
515 }
516
517 #[test]
518 fn writes_after_another_clone_finishes_are_rejected() {
519 let (mut producer, _track) = producer(ProducerConfig::default());
520 producer.push(&rec(1)).unwrap();
521 producer.clone().finish().unwrap();
522
523 assert!(matches!(
524 producer.push(&rec(2)),
525 Err(crate::Error::Net(moq_net::Error::Closed))
526 ));
527 assert!(matches!(
528 producer.pop(1),
529 Err(crate::Error::Net(moq_net::Error::Closed))
530 ));
531 assert_eq!(producer.window(), vec![rec(1)]);
532 }
533
534 #[test]
535 fn a_pop_is_clamped_to_the_window() {
536 let mut live = Live::new(ProducerConfig::default());
537 live.push(0);
538 live.pop(9);
539 live.push(1);
540
541 assert_eq!(
542 live.finish(),
543 vec![
544 Event::Push {
545 index: 0,
546 value: rec(0)
547 },
548 Event::Pop(0..1),
549 Event::Push {
550 index: 1,
551 value: rec(1)
552 },
553 ]
554 );
555 }
556
557 #[test]
558 fn a_large_gap_is_one_skip_event() {
559 let mut decoder = Decoder::<Value>::new(ConsumerConfig::default());
560 let mut group = decoder.group();
561 group.decode(br#"{"offset":0,"records":[]}"#).unwrap();
562 let mut group = decoder.group();
563 group.decode(br#"{"offset":9007199254740991,"records":[]}"#).unwrap();
564
565 assert_eq!(decoder.next_event(), Some(Event::Skip(0..super::encoder::MAX_INDEX)));
566 assert_eq!(decoder.next_event(), None);
567 }
568
569 #[test]
570 fn indices_must_fit_the_shared_safe_integer_range() {
571 let mut decoder = Decoder::<Value>::new(ConsumerConfig::default());
572 let mut group = decoder.group();
573 assert!(group.decode(br#"{"offset":9007199254740992,"records":[]}"#).is_err());
574
575 let mut decoder = Decoder::<Value>::new(ConsumerConfig::default());
576 let mut group = decoder.group();
577 group.decode(br#"{"offset":9007199254740991,"records":[]}"#).unwrap();
578 assert!(group.decode(br#"{"push":null}"#).is_err());
579 }
580
581 #[test]
582 fn a_checkpoint_cannot_start_before_the_window() {
583 let mut decoder = Decoder::<Value>::new(ConsumerConfig::default());
584 let mut group = decoder.group();
585 assert!(group.decode(br#"{"offset":2,"start":1,"records":[]}"#).is_err());
586 }
587
588 #[test]
589 fn every_group_requires_a_header() {
590 let mut decoder = Decoder::<Value>::new(ConsumerConfig::default());
591 let mut group = decoder.group();
592 group.decode(br#"{"offset":0,"records":[]}"#).unwrap();
593 let mut group = decoder.group();
594
595 assert!(group.decode(br#"{"push":null}"#).is_err());
596 }
597
598 #[test]
599 fn a_header_is_only_valid_as_frame_zero() {
600 let mut decoder = Decoder::<Value>::new(ConsumerConfig::default());
601 let mut group = decoder.group();
602 group.decode(br#"{"offset":0,"records":[]}"#).unwrap();
603
604 assert!(group.decode(br#"{"offset":0,"records":[]}"#).is_err());
605 }
606
607 #[test]
608 fn rolling_is_invisible_to_the_consumer() {
609 let edits = |ratio: u32| {
611 let mut live = Live::new(ProducerConfig::default().with_op_ratio(ratio));
612 for n in 0..6 {
613 live.push(n);
614 if n >= 3 {
615 live.pop(1);
616 }
617 }
618 live.finish()
619 };
620
621 assert_eq!(edits(1_000), edits(0));
622 }
623}