pub mod consumer;
mod decoder;
mod encoder;
pub mod producer;
pub use consumer::Consumer;
pub use decoder::Decoder;
pub use encoder::{Config, Encoder, Pending};
pub use producer::Producer;
#[cfg(test)]
mod test {
use std::task::Poll;
use serde_json::{Value, json};
use super::*;
use crate::Compression;
fn producer(config: Config) -> (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)
}
#[test]
fn demand_follows_subscribers() {
let (producer, consumer) = producer(Config::default());
let demand = producer.demand();
let waiter = kio::Waiter::noop();
assert!(matches!(demand.poll_used(&waiter), Poll::Ready(Ok(()))));
drop(consumer);
assert!(matches!(demand.poll_unused(&waiter), Poll::Ready(Ok(()))));
assert!(demand.poll_used(&waiter).is_pending());
let _consumer = producer.consume();
assert!(matches!(demand.poll_used(&waiter), Poll::Ready(Ok(()))));
assert!(demand.poll_unused(&waiter).is_pending());
}
fn compressed() -> Config {
Config {
compression: Compression::Deflate,
}
}
fn consume(track: moq_net::track::Subscriber, compression: bool) -> Consumer<Value> {
Consumer::new(
track,
Config {
compression: if compression {
Compression::Deflate
} else {
Compression::None
},
},
)
}
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(Config::default());
for n in 0..5 {
producer.append(&json!({ "n": n })).unwrap();
}
producer.finish().unwrap();
let records = drain(consume(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(consume(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(consume(track, true)).len(), 50);
}
#[test]
fn live_consumer_sees_each_record() {
let (mut producer, track) = producer(compressed());
let mut consumer = consume(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, 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 mut track = track.ordered();
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 mut subscriber = track.subscribe(None);
let mut producer = Producer::<std::collections::BTreeMap<(u8, u8), u8>>::new(track, Config::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");
let waiter = kio::Waiter::noop();
assert!(
matches!(subscriber.poll_recv_group(&waiter), Poll::Ready(Err(_))),
"the log is missing a record, so the track must end rather than stay writable"
);
}
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_failed_write_aborts_the_track() {
let track = rejecting_track();
let mut subscriber = track.subscribe(None);
let mut producer = Producer::<Value>::new(track, Config::default());
assert!(matches!(producer.append(&json!({ "n": 1 })), Err(crate::Error::Net(_))));
let waiter = kio::Waiter::noop();
assert!(
matches!(subscriber.poll_recv_group(&waiter), Poll::Ready(Err(_))),
"a truncated log must surface an error rather than read as a completed one"
);
}
#[test]
fn a_failed_write_ends_the_track() {
let track = rejecting_track();
let mut producer = Producer::<Value>::new(track, compressed());
assert!(matches!(producer.append(&json!({ "n": 1 })), Err(crate::Error::Net(_))));
assert!(
matches!(producer.append(&json!({ "n": 2 })), Err(crate::Error::Net(_))),
"a second append must fail on the ended track rather than open another group"
);
let waiter = kio::Waiter::noop();
assert!(matches!(
producer.consume().poll_recv_group(&waiter),
Poll::Ready(Err(_))
));
}
#[test]
fn appending_after_finish_fails_without_aborting() {
let (mut producer, _track) = producer(compressed());
producer.append(&json!({ "n": 0 })).unwrap();
producer.finish().unwrap();
assert!(producer.append(&json!({ "n": 1 })).is_err());
assert_eq!(drain(consume(producer.consume(), true)), vec![json!({ "n": 0 })]);
}
#[test]
fn a_second_group_is_reported_while_the_first_is_open() {
let track = moq_net::broadcast::Info::new()
.produce()
.create_track("test", None)
.unwrap();
let subscription = moq_net::track::Subscription::default().with_max_age(std::time::Duration::from_secs(30));
let subscriber = track.subscribe(subscription);
let mut first = track.append_group().unwrap();
first
.write_frame(moq_net::Timestamp::now(), br#"{"n":0}"#.as_slice())
.unwrap();
let mut second = track.append_group().unwrap();
second
.write_frame(moq_net::Timestamp::now(), br#"{"n":1}"#.as_slice())
.unwrap();
let mut consumer = consume(subscriber, false);
let waiter = kio::Waiter::noop();
assert!(matches!(
consumer.poll_next(&waiter),
Poll::Ready(Ok(Some(value))) if value == json!({ "n": 0 })
));
assert!(matches!(
consumer.poll_next(&waiter),
Poll::Ready(Err(crate::Error::Rolled))
));
first
.write_frame(moq_net::Timestamp::now(), br#"{"n":2}"#.as_slice())
.unwrap();
assert!(matches!(
consumer.poll_next(&waiter),
Poll::Ready(Err(crate::Error::Rolled))
));
}
#[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(consume(track, true));
assert_eq!(records, vec![value.clone(), value.clone(), value.clone(), value]);
}
}