mod consumer;
mod decoder;
mod encoder;
mod producer;
pub use consumer::Consumer;
pub use decoder::{ConsumerConfig, Decoder};
pub use encoder::{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 compressed() -> ProducerConfig {
ProducerConfig::default().with_compression(true)
}
fn consumer(track: moq_net::track::Subscriber, compression: bool) -> Consumer<Value> {
Consumer::new(track, ConsumerConfig::default().with_compression(compression))
}
fn drain(mut consumer: Consumer<Value>) -> Vec<Value> {
let waiter = kio::Waiter::noop();
let mut out = Vec::new();
while let Poll::Ready(Ok(Some(value))) = consumer.poll_next(&waiter) {
out.push(value);
}
out
}
#[test]
fn plaintext_roundtrip_in_order() {
let (mut producer, track) = producer(ProducerConfig::default());
for n in 0..5 {
producer.append(&json!({ "n": n })).unwrap();
}
producer.finish().unwrap();
let records = drain(consumer(track, false));
assert_eq!(records, (0..5).map(|n| json!({ "n": n })).collect::<Vec<_>>());
}
#[test]
fn compressed_roundtrip_in_order() {
let (mut producer, track) = producer(compressed());
for n in 0..20 {
producer.append(&json!({ "group": n, "pts": n * 2_000 })).unwrap();
}
producer.finish().unwrap();
let records = drain(consumer(track, true));
assert_eq!(records.len(), 20);
assert_eq!(records[7], json!({ "group": 7, "pts": 14_000 }));
}
#[test]
fn all_records_ride_one_group() {
let (mut producer, track) = producer(compressed());
for n in 0..50 {
producer.append(&json!({ "n": n })).unwrap();
}
producer.finish().unwrap();
assert_eq!(track.latest(), Some(0));
assert_eq!(drain(consumer(track, true)).len(), 50);
}
#[test]
fn live_consumer_sees_each_record() {
let (mut producer, track) = producer(compressed());
let mut consumer = consumer(track, true);
let waiter = kio::Waiter::noop();
for n in 0..3 {
producer.append(&json!({ "n": n })).unwrap();
match consumer.poll_next(&waiter) {
Poll::Ready(Ok(Some(value))) => assert_eq!(value, json!({ "n": n })),
other => panic!("expected record, got {other:?}"),
}
}
assert!(matches!(consumer.poll_next(&waiter), Poll::Pending));
producer.finish().unwrap();
}
#[test]
fn shared_window_shrinks_repetitive_records() {
let (mut producer, mut track) = producer(compressed());
for n in 0..8 {
producer.append(&json!({ "group": n, "pts": n * 2_000 })).unwrap();
}
producer.finish().unwrap();
let waiter = kio::Waiter::noop();
let Poll::Ready(Ok(Some(mut group))) = track.poll_next_group(&waiter) else {
panic!("expected a group");
};
let mut sizes = Vec::new();
while let Poll::Ready(Ok(Some(frame))) = group.poll_read_frame(&waiter) {
sizes.push(frame.payload.len());
}
assert_eq!(sizes.len(), 8);
let raw = serde_json::to_vec(&json!({ "group": 7, "pts": 14_000 })).unwrap().len();
assert!(
*sizes.last().unwrap() < raw / 2,
"windowed record {} should be far below its raw size {raw}",
sizes.last().unwrap()
);
}
#[test]
fn a_rejected_record_does_not_open_a_group() {
let track = moq_net::broadcast::Info::new()
.produce()
.create_track("test", None)
.unwrap();
let subscriber = track.subscribe(None);
let mut producer = Producer::<std::collections::BTreeMap<(u8, u8), u8>>::new(track, ProducerConfig::default());
let mut bad = std::collections::BTreeMap::new();
bad.insert((1, 2), 3);
assert!(producer.append(&bad).is_err());
assert_eq!(subscriber.latest(), None, "a rejected record opened a group");
}
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()
}
#[test]
fn a_rejected_record_does_not_strand_an_empty_group() {
let track = rejecting_track();
let mut subscriber = track.subscribe(None);
let mut producer = Producer::<Value>::new(track, ProducerConfig::default());
assert!(producer.append(&json!({ "n": 1 })).is_err());
let waiter = kio::Waiter::noop();
let Poll::Ready(Ok(Some(mut group))) = subscriber.poll_next_group(&waiter) else {
panic!("the group was published, so a subscriber sees it");
};
assert!(
matches!(group.poll_read_frame(&waiter), Poll::Ready(Ok(None))),
"the empty group must be closed, not left open for a subscriber to wait in"
);
}
#[test]
fn a_rejected_record_leaves_the_encoder_able_to_retry() {
let track = rejecting_track();
let mut producer = Producer::<Value>::new(track, ProducerConfig::default().with_compression(true));
assert!(matches!(producer.append(&json!({ "n": 1 })), Err(crate::Error::Net(_))));
assert!(
matches!(producer.append(&json!({ "n": 2 })), Err(crate::Error::Net(_))),
"the encoder latched a desync instead of retrying into a fresh group"
);
}
#[test]
fn embedded_newlines_survive() {
let (mut producer, track) = producer(compressed());
let value = json!({ "s": "line1\nline2\ttab", "u": "a\u{000a}b" });
for _ in 0..4 {
producer.append(&value).unwrap();
}
producer.finish().unwrap();
let records = drain(consumer(track, true));
assert_eq!(records, vec![value.clone(), value.clone(), value.clone(), value]);
}
}