mod consumer;
mod decoder;
mod encoder;
mod op;
mod producer;
pub use consumer::Consumer;
pub use decoder::{ConsumerConfig, Decoder, Event, Group};
pub use encoder::{Encoded, Encoder, Pending, ProducerConfig};
pub use producer::Producer;
#[cfg(test)]
mod test {
use std::task::Poll;
use serde_json::{Value, json};
use super::*;
fn producer(config: ProducerConfig) -> (Producer<Value>, moq_net::track::Subscriber) {
let track = moq_net::broadcast::Info::new()
.produce()
.create_track("test", None)
.unwrap();
let consumer = track.subscribe(None);
(Producer::new(track, config), consumer)
}
fn consumer(track: moq_net::track::Subscriber, compression: bool) -> Consumer<Value> {
Consumer::new(track, ConsumerConfig::default().with_compression(compression))
}
fn rejecting_track() -> moq_net::track::Producer {
let mut info = moq_net::track::Info::default();
info.timescale = moq_net::Timescale::new((1u64 << 62) - 1).unwrap();
moq_net::broadcast::Info::new()
.produce()
.create_track("test", Some(info))
.unwrap()
}
fn drain(consumer: &mut Consumer<Value>) -> Vec<Event<Value>> {
let waiter = kio::Waiter::noop();
let mut out = Vec::new();
while let Poll::Ready(Ok(Some(event))) = consumer.poll_next(&waiter) {
out.push(event);
}
out
}
fn rec(n: u64) -> Value {
json!({ "n": n })
}
struct Live {
producer: Producer<Value>,
consumer: Consumer<Value>,
events: Vec<Event<Value>>,
}
impl Live {
fn new(config: ProducerConfig) -> Self {
let compression = config.compression;
let (producer, track) = producer(config);
Self {
producer,
consumer: consumer(track, compression),
events: Vec::new(),
}
}
fn push(&mut self, n: u64) {
self.producer.push(&rec(n)).unwrap();
self.read();
}
fn pop(&mut self, count: u64) {
self.producer.pop(count).unwrap();
self.read();
}
fn read(&mut self) {
self.events.extend(drain(&mut self.consumer));
}
fn finish(self) -> Vec<Event<Value>> {
let Live {
producer,
mut consumer,
mut events,
} = self;
producer.finish().unwrap();
events.extend(drain(&mut consumer));
events
}
fn pushed(events: &[Event<Value>]) -> Vec<u64> {
events
.iter()
.filter_map(|e| match e {
Event::Push { index, .. } => Some(*index),
_ => None,
})
.collect()
}
}
#[test]
fn push_and_pop_round_trip() {
let mut live = Live::new(ProducerConfig::default());
live.push(0);
live.push(1);
live.pop(1);
live.push(2);
assert_eq!(
live.finish(),
vec![
Event::Push {
index: 0,
value: rec(0)
},
Event::Push {
index: 1,
value: rec(1)
},
Event::Pop(0..1),
Event::Push {
index: 2,
value: rec(2)
},
]
);
}
#[test]
fn the_window_slides() {
let (mut producer, _track) = producer(ProducerConfig::default());
for n in 0..5 {
producer.push(&rec(n)).unwrap();
if n >= 2 {
producer.pop(1).unwrap();
}
}
assert_eq!(producer.range(), 3..5);
assert_eq!(producer.window(), vec![rec(3), rec(4)]);
}
#[test]
fn a_popped_record_is_never_restated() {
let mut live = Live::new(ProducerConfig::default().with_op_ratio(0));
live.push(0);
live.push(1);
live.pop(1);
live.push(2);
assert_eq!(
live.finish(),
vec![
Event::Push {
index: 0,
value: rec(0)
},
Event::Push {
index: 1,
value: rec(1)
},
Event::Pop(0..1),
Event::Push {
index: 2,
value: rec(2)
},
]
);
}
#[test]
fn a_fresh_consumer_adopts_the_offset_without_skipping_history() {
let track = moq_net::broadcast::Info::new()
.produce()
.create_track("test", None)
.unwrap();
let mut producer = Producer::<Value>::new(track, ProducerConfig::default().with_op_ratio(0));
for n in 0..5 {
producer.push(&rec(n)).unwrap();
}
producer.pop(3).unwrap();
let mut subscriber = producer.consume();
subscriber.start_at(subscriber.latest().unwrap());
let mut fresh = consumer(subscriber, false);
producer.finish().unwrap();
let events = drain(&mut fresh);
assert_eq!(
events,
vec![
Event::Push {
index: 3,
value: rec(3)
},
Event::Push {
index: 4,
value: rec(4)
}
]
);
assert!(!events.iter().any(|e| matches!(e, Event::Skip(_))));
}
#[test]
fn a_lagging_consumer_is_told_what_it_missed() {
let mut encoder = Encoder::<Value>::new(ProducerConfig::default().with_op_ratio(0));
let mut decoder = Decoder::<Value>::new(ConsumerConfig::default());
for n in 0..2 {
let frame = encoder.push(&rec(n)).unwrap();
let mut group = decoder.group();
group.decode(&frame.payload).unwrap();
frame.commit();
}
assert_eq!(
std::iter::from_fn(|| decoder.next_event()).collect::<Vec<_>>(),
vec![
Event::Push {
index: 0,
value: rec(0)
},
Event::Push {
index: 1,
value: rec(1)
}
]
);
let mut latest = None;
for n in 2..8 {
let frame = encoder.push(&rec(n)).unwrap();
frame.commit();
let frame = encoder.pop(1).unwrap().unwrap();
latest = Some(frame.payload.clone());
frame.commit();
}
let mut group = decoder.group();
group.decode(&latest.unwrap()).unwrap();
let events = std::iter::from_fn(|| decoder.next_event()).collect::<Vec<_>>();
let skipped: Vec<std::ops::Range<u64>> = events
.iter()
.filter_map(|e| match e {
Event::Skip(range) => Some(range.clone()),
_ => None,
})
.collect();
assert!(!skipped.is_empty(), "expected skips, got {events:?}");
assert_eq!(skipped.first().map(|range| range.start), Some(2));
let reported: Vec<u64> = events
.iter()
.flat_map(|e| match e {
Event::Push { index, .. } => vec![*index],
Event::Skip(range) => range.clone().collect(),
Event::Pop(_) => Vec::new(),
})
.collect();
assert!(reported.windows(2).all(|w| w[1] == w[0] + 1), "gaps in {reported:?}");
}
#[test]
fn compressed_round_trip_across_rolls() {
let mut live = Live::new(ProducerConfig::default().with_compression(true).with_op_ratio(1));
for n in 0..40 {
live.push(n);
if n >= 10 {
live.pop(1);
}
}
assert_eq!(Live::pushed(&live.finish()), (0..40).collect::<Vec<_>>());
}
#[test]
fn an_empty_pop_writes_nothing() {
let (mut producer, track) = producer(ProducerConfig::default());
producer.pop(5).unwrap();
producer.finish().unwrap();
assert_eq!(track.latest(), None);
}
#[test]
fn a_rejected_edit_leaves_the_window_unchanged() {
let track = rejecting_track();
let mut subscriber = track.subscribe(None);
let mut producer = Producer::<Value>::new(track, ProducerConfig::default());
assert!(producer.push(&rec(1)).is_err());
assert_eq!(producer.range(), 0..0);
assert!(producer.window().is_empty());
let waiter = kio::Waiter::noop();
let Poll::Ready(Ok(Some(mut group))) = subscriber.poll_next_group(&waiter) else {
panic!("the rejected group's header was published");
};
assert!(matches!(group.poll_read_frame(&waiter), Poll::Ready(Ok(None))));
}
#[test]
fn writes_after_another_clone_finishes_are_rejected() {
let (mut producer, _track) = producer(ProducerConfig::default());
producer.push(&rec(1)).unwrap();
producer.clone().finish().unwrap();
assert!(matches!(
producer.push(&rec(2)),
Err(crate::Error::Net(moq_net::Error::Closed))
));
assert!(matches!(
producer.pop(1),
Err(crate::Error::Net(moq_net::Error::Closed))
));
assert_eq!(producer.window(), vec![rec(1)]);
}
#[test]
fn a_pop_is_clamped_to_the_window() {
let mut live = Live::new(ProducerConfig::default());
live.push(0);
live.pop(9);
live.push(1);
assert_eq!(
live.finish(),
vec![
Event::Push {
index: 0,
value: rec(0)
},
Event::Pop(0..1),
Event::Push {
index: 1,
value: rec(1)
},
]
);
}
#[test]
fn a_large_gap_is_one_skip_event() {
let mut decoder = Decoder::<Value>::new(ConsumerConfig::default());
let mut group = decoder.group();
group.decode(br#"{"offset":0,"records":[]}"#).unwrap();
let mut group = decoder.group();
group.decode(br#"{"offset":9007199254740991,"records":[]}"#).unwrap();
assert_eq!(decoder.next_event(), Some(Event::Skip(0..super::encoder::MAX_INDEX)));
assert_eq!(decoder.next_event(), None);
}
#[test]
fn indices_must_fit_the_shared_safe_integer_range() {
let mut decoder = Decoder::<Value>::new(ConsumerConfig::default());
let mut group = decoder.group();
assert!(group.decode(br#"{"offset":9007199254740992,"records":[]}"#).is_err());
let mut decoder = Decoder::<Value>::new(ConsumerConfig::default());
let mut group = decoder.group();
group.decode(br#"{"offset":9007199254740991,"records":[]}"#).unwrap();
assert!(group.decode(br#"{"push":null}"#).is_err());
}
#[test]
fn every_group_requires_a_header() {
let mut decoder = Decoder::<Value>::new(ConsumerConfig::default());
let mut group = decoder.group();
group.decode(br#"{"offset":0,"records":[]}"#).unwrap();
let mut group = decoder.group();
assert!(group.decode(br#"{"push":null}"#).is_err());
}
#[test]
fn a_header_is_only_valid_as_frame_zero() {
let mut decoder = Decoder::<Value>::new(ConsumerConfig::default());
let mut group = decoder.group();
group.decode(br#"{"offset":0,"records":[]}"#).unwrap();
assert!(group.decode(br#"{"offset":0,"records":[]}"#).is_err());
}
#[test]
fn rolling_is_invisible_to_the_consumer() {
let edits = |ratio: u32| {
let mut live = Live::new(ProducerConfig::default().with_op_ratio(ratio));
for n in 0..6 {
live.push(n);
if n >= 3 {
live.pop(1);
}
}
live.finish()
};
assert_eq!(edits(1_000), edits(0));
}
}