mod consumer;
mod decoder;
mod encoder;
mod producer;
pub use consumer::Consumer;
pub use decoder::{ConsumerConfig, Decoder};
pub use encoder::{Encoded, Encoder, Pending, ProducerConfig};
pub use producer::{Guard, Producer};
#[cfg(test)]
mod test {
use std::task::Poll;
use bytes::Bytes;
use serde_json::{Value, json};
use super::encoder::MAX_DELTA_FRAMES;
use super::*;
fn cfg(delta_ratio: u32) -> ProducerConfig {
ProducerConfig::default().with_delta_ratio(delta_ratio)
}
fn cfg_deflate(delta_ratio: u32) -> ProducerConfig {
cfg(delta_ratio).with_compression(true)
}
fn deflate_consumer(track: moq_net::track::Subscriber) -> Consumer<Value> {
Consumer::new(track, ConsumerConfig::default().with_compression(true))
}
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 drain(track: moq_net::track::Subscriber) -> Vec<Value> {
drain_with(Consumer::<Value>::new(track, ConsumerConfig::default()))
}
fn drain_with(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 deltas_off_snapshot_per_group() {
let (mut producer, track) = producer(cfg(0));
producer.update(&json!({ "a": 1 })).unwrap();
producer.update(&json!({ "a": 2 })).unwrap();
producer.finish().unwrap();
assert_eq!(track.latest(), Some(1));
assert_eq!(drain(track), vec![json!({ "a": 2 })]);
}
#[test]
fn live_consumer_sees_each_update() {
let (mut producer, track) = producer(ProducerConfig::default());
let mut consumer = Consumer::<Value>::new(track, ConsumerConfig::default());
let waiter = kio::Waiter::noop();
for n in 1..=3 {
producer.update(&json!({ "a": n })).unwrap();
match consumer.poll_next(&waiter) {
Poll::Ready(Ok(Some(value))) => assert_eq!(value, json!({ "a": n })),
other => panic!("expected value, got {other:?}"),
}
}
}
#[test]
fn unchanged_value_writes_nothing() {
let (mut producer, track) = producer(ProducerConfig::default());
producer.update(&json!({ "a": 1 })).unwrap();
producer.update(&json!({ "a": 1 })).unwrap();
producer.finish().unwrap();
assert_eq!(track.latest(), Some(0));
assert_eq!(drain(track), vec![json!({ "a": 1 })]);
}
#[test]
fn deltas_share_one_group() {
let (mut producer, track) = producer(cfg(100));
producer.update(&json!({ "a": 1, "b": 1 })).unwrap();
producer.update(&json!({ "a": 1, "b": 2 })).unwrap();
producer.update(&json!({ "a": 1, "b": 3 })).unwrap();
producer.finish().unwrap();
assert_eq!(track.latest(), Some(0));
let values = drain(track);
assert_eq!(values.last().unwrap(), &json!({ "a": 1, "b": 3 }));
}
#[test]
fn tight_ratio_rolls_snapshots() {
let (mut producer, track) = producer(cfg(1));
producer.update(&json!({ "a": 1 })).unwrap(); producer.update(&json!({ "a": 2 })).unwrap(); producer.update(&json!({ "a": 3 })).unwrap(); producer.update(&json!({ "a": 4 })).unwrap(); producer.finish().unwrap();
assert_eq!(track.latest(), Some(1));
}
#[test]
fn deltas_stay_within_ratio_times_snapshot() {
let (mut producer, track) = producer(cfg(8));
for n in 0..=10 {
producer.update(&json!({ "n": n })).unwrap();
}
producer.finish().unwrap();
assert_eq!(track.latest(), Some(1));
assert_eq!(drain(track).last().unwrap(), &json!({ "n": 10 }));
}
#[test]
fn array_change_is_delta() {
let (mut producer, track) = producer(cfg(100));
producer.update(&json!({ "list": [1, 2] })).unwrap();
producer.update(&json!({ "list": [1, 2, 3] })).unwrap();
producer.finish().unwrap();
assert_eq!(track.latest(), Some(0));
assert_eq!(drain(track).last().unwrap(), &json!({ "list": [1, 2, 3] }));
}
#[test]
fn frame_cap_rolls_snapshot() {
let (mut producer, track) = producer(cfg(1_000_000));
for i in 0..=MAX_DELTA_FRAMES {
producer.update(&json!({ "n": i })).unwrap();
}
producer.finish().unwrap();
assert_eq!(track.latest(), Some(1));
assert_eq!(drain(track).last().unwrap(), &json!({ "n": MAX_DELTA_FRAMES }));
}
#[test]
fn late_joiner_reconstructs_from_deltas() {
let (mut producer, track) = producer(cfg(100));
producer.update(&json!({ "a": 1, "b": 1 })).unwrap();
producer.update(&json!({ "a": 1, "b": 2 })).unwrap();
producer.update(&json!({ "a": 5, "b": 2 })).unwrap();
producer.finish().unwrap();
assert_eq!(drain(track).last().unwrap(), &json!({ "a": 5, "b": 2 }));
}
#[test]
fn lock_composes_independent_owners() {
#[derive(serde::Serialize, serde::Deserialize, Default, PartialEq, Debug)]
struct Doc {
#[serde(skip_serializing_if = "Option::is_none")]
video: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
scte35: Option<u32>,
}
let track = moq_net::broadcast::Info::new()
.produce()
.create_track("test", None)
.unwrap();
let consumer = track.subscribe(None);
let mut producer = Producer::<Doc>::new(track, ProducerConfig::default());
producer.lock().video = Some("v1".to_string());
producer.lock().scte35 = Some(42);
let _ = producer.lock();
producer.finish().unwrap();
let mut consumer = Consumer::<Doc>::new(consumer, ConsumerConfig::default());
let waiter = kio::Waiter::noop();
let mut last = None;
while let Poll::Ready(Ok(Some(value))) = consumer.poll_next(&waiter) {
last = Some(value);
}
assert_eq!(
last.unwrap(),
Doc {
video: Some("v1".to_string()),
scte35: Some(42),
}
);
}
#[test]
fn commit_reports_a_publish_failure() {
#[derive(serde::Serialize, serde::Deserialize, Default, PartialEq, Debug)]
struct Doc {
a: u32,
}
let track = moq_net::broadcast::Info::new()
.produce()
.create_track("test", None)
.unwrap();
let mut producer = Producer::<Doc>::new(track, ProducerConfig::default());
producer.finish().unwrap();
let mut guard = producer.lock();
guard.a = 1;
assert!(matches!(guard.commit(), Err(crate::Error::Net(_))));
}
#[test]
fn a_failed_publish_keeps_the_value_for_lock() {
#[derive(serde::Serialize, serde::Deserialize, Default, PartialEq, Debug)]
struct Doc {
#[serde(skip_serializing_if = "Option::is_none")]
video: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
scte35: Option<u32>,
}
let track = moq_net::broadcast::Info::new()
.produce()
.create_track("test", None)
.unwrap();
let mut producer = Producer::<Doc>::new(track, ProducerConfig::default());
producer.lock().video = Some("v1".to_string());
producer.finish().unwrap();
let mut guard = producer.lock();
guard.scte35 = Some(42);
assert!(guard.commit().is_err());
let guard = producer.lock();
assert_eq!(
guard.video,
Some("v1".to_string()),
"a rejected frame erased the baseline"
);
}
#[test]
fn update_after_finish_errors_instead_of_panicking() {
let (mut producer, _track) = producer(cfg(100));
producer.update(&json!({ "a": 1 })).unwrap();
producer.update(&json!({ "a": 2 })).unwrap(); producer.finish().unwrap();
assert!(matches!(producer.update(&json!({ "a": 3 })), Err(crate::Error::Net(_))));
}
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_snapshot_does_not_strand_an_empty_group() {
let track = rejecting_track();
let mut subscriber = track.subscribe(None);
let mut producer = Producer::<Value>::new(track, cfg(0));
assert!(matches!(producer.update(&json!({ "a": 1 })), Err(crate::Error::Net(_))));
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 commit_publishes_once() {
#[derive(serde::Serialize, serde::Deserialize, Default, PartialEq, Debug)]
struct Doc {
a: u32,
}
let track = moq_net::broadcast::Info::new()
.produce()
.create_track("test", None)
.unwrap();
let consumer = track.subscribe(None);
let mut producer = Producer::<Doc>::new(track, cfg(0));
let mut guard = producer.lock();
guard.a = 1;
guard.commit().unwrap();
producer.finish().unwrap();
assert_eq!(consumer.latest(), Some(0));
}
#[test]
fn newer_group_supersedes_in_progress_reconstruction() {
let (mut producer, track) = producer(cfg(1));
let observer = producer.consume();
let mut consumer = Consumer::<Value>::new(track, ConsumerConfig::default());
let waiter = kio::Waiter::noop();
producer.update(&json!({ "a": 1 })).unwrap(); match consumer.poll_next(&waiter) {
Poll::Ready(Ok(Some(value))) => assert_eq!(value, json!({ "a": 1 })),
other => panic!("expected first value, got {other:?}"),
}
producer.update(&json!({ "a": 2 })).unwrap(); producer.update(&json!({ "a": 3 })).unwrap(); producer.update(&json!({ "a": 4 })).unwrap(); producer.finish().unwrap();
assert_eq!(observer.latest(), Some(1));
let mut last = None;
while let Poll::Ready(Ok(Some(value))) = consumer.poll_next(&waiter) {
last = Some(value);
}
assert_eq!(last.unwrap(), json!({ "a": 4 }));
}
#[test]
fn open_group_pends_after_track_finish() {
let mut track = moq_net::broadcast::Info::new()
.produce()
.create_track("test", None)
.unwrap();
let mut group = track.append_group().unwrap();
let consumer_track = track.subscribe(None);
track.finish().unwrap();
let mut consumer = Consumer::<Value>::new(consumer_track, ConsumerConfig::default());
let waiter = kio::Waiter::noop();
assert!(matches!(consumer.poll_next(&waiter), Poll::Pending));
group
.write_frame(
moq_net::Timestamp::ZERO,
Bytes::from(serde_json::to_vec(&json!({ "a": 1 })).unwrap()),
)
.unwrap();
group.finish().unwrap();
match consumer.poll_next(&waiter) {
Poll::Ready(Ok(Some(value))) => assert_eq!(value, json!({ "a": 1 })),
other => panic!("expected the catalog value, got {other:?}"),
}
}
#[test]
fn late_joiner_collapses_backlog_to_latest() {
let (mut producer, track) = producer(cfg(100));
for n in 0..=20 {
producer.update(&json!({ "n": n })).unwrap();
}
producer.finish().unwrap();
assert_eq!(track.latest(), Some(0));
let values = drain(track);
assert_eq!(
values,
vec![json!({ "n": 20 })],
"backlog should collapse to the latest value"
);
}
#[test]
fn compressed_late_joiner_collapses_backlog_to_latest() {
let (mut producer, track) = producer(cfg_deflate(100));
for n in 0..=20 {
producer.update(&json!({ "n": n })).unwrap();
}
producer.finish().unwrap();
assert_eq!(track.latest(), Some(0));
let values = drain_with(deflate_consumer(track));
assert_eq!(
values,
vec![json!({ "n": 20 })],
"compressed backlog should collapse to the latest"
);
}
#[test]
fn compressed_snapshot_per_group_roundtrips() {
let (mut producer, track) = producer(cfg_deflate(0));
producer.update(&json!({ "a": 1 })).unwrap();
producer.update(&json!({ "a": 2 })).unwrap();
producer.finish().unwrap();
assert_eq!(track.latest(), Some(1));
let values = drain_with(deflate_consumer(track));
assert_eq!(values, vec![json!({ "a": 2 })]);
}
#[test]
fn compressed_deltas_share_one_group() {
let (mut producer, track) = producer(cfg_deflate(100));
producer.update(&json!({ "a": 1, "b": 1 })).unwrap();
producer.update(&json!({ "a": 1, "b": 2 })).unwrap();
producer.update(&json!({ "a": 1, "b": 3 })).unwrap();
producer.finish().unwrap();
assert_eq!(track.latest(), Some(0));
let values = drain_with(deflate_consumer(track));
assert_eq!(values.last().unwrap(), &json!({ "a": 1, "b": 3 }));
}
#[test]
fn compressed_late_joiner_reconstructs_from_deltas() {
let (mut producer, track) = producer(cfg_deflate(100));
producer.update(&json!({ "a": 1, "b": 1 })).unwrap();
producer.update(&json!({ "a": 1, "b": 2 })).unwrap();
producer.update(&json!({ "a": 5, "b": 2 })).unwrap();
producer.finish().unwrap();
let values = drain_with(deflate_consumer(track));
assert_eq!(values.last().unwrap(), &json!({ "a": 5, "b": 2 }));
}
#[test]
fn compressed_deltas_roll_on_compressed_budget() {
let (mut producer, track) = producer(cfg_deflate(2));
for n in 0..=40 {
producer.update(&json!({ "n": n })).unwrap();
}
producer.finish().unwrap();
assert!(
track.latest().unwrap() > 0,
"a tight ratio should roll at least one compressed group"
);
assert_eq!(drain_with(deflate_consumer(track)).last().unwrap(), &json!({ "n": 40 }));
}
#[test]
fn compression_shrinks_wire_frames() {
let value = json!({ "renditions": ["video".repeat(50), "video".repeat(50), "video".repeat(50)] });
let plaintext_bytes = wire_frame_len(cfg(0), &value);
let compressed_bytes = wire_frame_len(cfg_deflate(0), &value);
assert!(
compressed_bytes < plaintext_bytes,
"compressed frame {compressed_bytes} should be smaller than plaintext {plaintext_bytes}"
);
}
#[test]
fn compressed_deltas_reuse_window() {
let (mut producer, mut track) = producer(cfg_deflate(100));
let phrase = "Media over QUIC delivers real-time latency at massive scale";
producer.update(&json!({ "note": phrase })).unwrap();
producer.update(&json!({ "note": phrase, "echo": phrase })).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 frames = Vec::new();
while let Poll::Ready(Ok(Some(frame))) = group.poll_read_frame(&waiter) {
frames.push(frame.payload);
}
assert_eq!(frames.len(), 2, "snapshot + one delta in a single group");
let raw_delta = serde_json::to_vec(&json!({ "echo": phrase })).unwrap();
assert!(
frames[1].len() < raw_delta.len() / 2,
"windowed delta {} should be far below the raw patch {}",
frames[1].len(),
raw_delta.len()
);
}
#[test]
fn rejected_field_names_its_path() {
#[derive(serde::Deserialize)]
#[allow(dead_code)]
struct Inner {
count: u8,
}
#[derive(serde::Deserialize)]
#[allow(dead_code)]
struct Outer {
inner: Inner,
}
let (mut producer, track) = producer(cfg(0));
producer.update(&json!({ "inner": { "count": 300 } })).unwrap();
let mut consumer = Consumer::<Outer>::new(track, ConsumerConfig::default());
let Poll::Ready(Err(err)) = consumer.poll_next(&kio::Waiter::noop()) else {
panic!("expected a deserialize error");
};
assert!(err.to_string().starts_with("json: inner.count: "), "{err}");
}
#[test]
fn rejected_root_omits_the_path() {
let (mut producer, track) = producer(cfg(0));
producer.update(&json!("not a map")).unwrap();
let mut consumer = Consumer::<std::collections::BTreeMap<String, u8>>::new(track, ConsumerConfig::default());
let Poll::Ready(Err(err)) = consumer.poll_next(&kio::Waiter::noop()) else {
panic!("expected a deserialize error");
};
assert_eq!(
err.to_string(),
"json: invalid type: string \"not a map\", expected a map"
);
}
fn wire_frame_len(config: ProducerConfig, value: &Value) -> usize {
let (mut producer, mut track) = producer(config);
producer.update(value).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 Poll::Ready(Ok(Some(frame))) = group.poll_read_frame(&waiter) else {
panic!("expected a frame");
};
frame.payload.len()
}
}