use crate::Compression;
pub mod consumer;
pub mod producer;
pub use consumer::Consumer;
pub use producer::Producer;
#[derive(Debug, Clone, Default)]
#[non_exhaustive]
pub struct Config {
pub compression: Compression,
}
#[cfg(test)]
mod test {
use std::task::Poll;
use bytes::Bytes;
use super::*;
fn cfg(compression: bool) -> Config {
Config {
compression: if compression {
Compression::Deflate
} else {
Compression::None
},
}
}
fn producer(compression: bool) -> (Producer, 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, cfg(compression)), consumer)
}
fn drain(track: moq_net::track::Subscriber, compression: bool) -> Vec<Bytes> {
let mut consumer = Consumer::new(track, cfg(compression));
let waiter = kio::Waiter::noop();
let mut out = Vec::new();
while let Poll::Ready(Ok(Some(payload))) = consumer.poll_next(&waiter) {
out.push(payload);
}
out
}
fn payloads(count: u8) -> Vec<Bytes> {
(0..count).map(|n| Bytes::from(vec![n; 16])).collect()
}
#[test]
fn every_payload_survives_in_order() {
let (mut producer, track) = producer(false);
let expected = payloads(5);
for payload in &expected {
producer.append(payload.clone()).unwrap();
}
producer.finish().unwrap();
assert_eq!(track.latest(), Some(0));
assert_eq!(drain(track, false), expected);
}
#[test]
fn compressed_roundtrip_in_order() {
let (mut producer, track) = producer(true);
let expected = payloads(20);
for payload in &expected {
producer.append(payload.clone()).unwrap();
}
producer.finish().unwrap();
assert_eq!(drain(track, true), expected);
}
#[test]
fn the_shared_window_shrinks_repetitive_payloads() {
let (mut producer, track) = producer(true);
let payload = Bytes::from(b"the quick brown fox".repeat(16));
for _ in 0..8 {
producer.append(payload.clone()).unwrap();
}
producer.finish().unwrap();
let waiter = kio::Waiter::noop();
let Poll::Ready(Ok(Some(mut group))) = track.ordered().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);
assert!(
*sizes.last().unwrap() < sizes[0],
"windowed frame {} should be below the first {}",
sizes.last().unwrap(),
sizes[0]
);
}
#[test]
fn clones_append_into_one_log() {
let (mut producer, track) = producer(true);
let mut clone = producer.clone();
producer.append(&b"a"[..]).unwrap();
clone.append(&b"b"[..]).unwrap();
producer.append(&b"c"[..]).unwrap();
producer.finish().unwrap();
assert_eq!(
drain(track, true),
vec![
Bytes::from_static(b"a"),
Bytes::from_static(b"b"),
Bytes::from_static(b"c")
]
);
}
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::new(track, Config::default());
assert!(producer.append(&b"rejected"[..]).is_err());
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 an_undecodable_record_ends_the_track() {
let track = moq_net::broadcast::Info::new()
.produce()
.create_track("test", None)
.unwrap();
let mut subscriber = track.subscribe(None);
let mut producer = Producer::new(track, cfg(true));
assert!(producer.is_used());
let oversized = Bytes::from(vec![0u8; moq_flate::DEFAULT_MAX_FRAME_SIZE as usize + 1]);
assert!(matches!(
producer.append(oversized),
Err(crate::Error::Flate(moq_flate::Error::TooLarge(_)))
));
assert!(!producer.is_used());
let waiter = kio::Waiter::noop();
assert!(matches!(subscriber.poll_recv_group(&waiter), Poll::Ready(Err(_))));
assert!(producer.append(&b"after"[..]).is_err());
}
#[test]
fn a_reader_inside_the_group_sees_the_real_error() {
let track = moq_net::broadcast::Info::new()
.produce()
.create_track("test", None)
.unwrap();
let mut subscriber = track.subscribe(None);
let mut producer = Producer::new(track, Config::default());
producer.append(&b"first"[..]).unwrap();
let waiter = kio::Waiter::noop();
let Poll::Ready(Ok(Some(mut group))) = subscriber.poll_recv_group(&waiter) else {
panic!("the group was published, so a subscriber sees it");
};
assert!(matches!(group.poll_read_frame(&waiter), Poll::Ready(Ok(Some(_)))));
let oversized = Bytes::from(vec![0u8; moq_net::group::MAX_CACHE_BYTES as usize + 1]);
assert!(producer.append(oversized).is_err());
match group.poll_read_frame(&waiter) {
Poll::Ready(Err(err)) => assert!(
matches!(err, moq_net::Error::FrameTooLarge),
"the reader should see the write's own error, got {err:?}"
),
other => panic!("expected the write error, got {other:?}"),
}
}
fn replaying() -> moq_net::track::Subscription {
moq_net::track::Subscription::default().with_max_age(std::time::Duration::from_secs(30))
}
#[test]
fn a_failed_write_ends_the_track() {
let track = rejecting_track();
let mut producer = Producer::new(track, Config::default());
assert!(producer.append(&b"rejected"[..]).is_err());
assert!(
producer.append(&b"again"[..]).is_err(),
"a second append must fail on the closed track rather than open another group"
);
}
#[test]
fn a_second_group_is_a_rolled_log() {
let track = moq_net::broadcast::Info::new()
.produce()
.create_track("test", None)
.unwrap();
let subscriber = track.subscribe(replaying());
for pair in payloads(4).chunks(2) {
let mut flate = moq_flate::Encoder::new();
let mut group = track.append_group().unwrap();
for payload in pair {
group
.write_frame(moq_net::Timestamp::now(), flate.frame(payload))
.unwrap();
}
group.finish().unwrap();
}
track.finish().unwrap();
let mut consumer = Consumer::new(subscriber, cfg(true));
let waiter = kio::Waiter::noop();
assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(_)))));
assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(_)))));
assert!(matches!(
consumer.poll_next(&waiter),
Poll::Ready(Err(crate::Error::Rolled))
));
}
#[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 subscriber = track.subscribe(replaying());
let mut first = track.append_group().unwrap();
first.write_frame(moq_net::Timestamp::now(), &b"first"[..]).unwrap();
let mut second = track.append_group().unwrap();
second.write_frame(moq_net::Timestamp::now(), &b"second"[..]).unwrap();
let mut consumer = Consumer::new(subscriber, Config::default());
let waiter = kio::Waiter::noop();
assert!(matches!(
consumer.poll_next(&waiter),
Poll::Ready(Ok(Some(payload))) if payload == Bytes::from_static(b"first")
));
assert!(matches!(
consumer.poll_next(&waiter),
Poll::Ready(Err(crate::Error::Rolled))
));
first.write_frame(moq_net::Timestamp::now(), &b"more"[..]).unwrap();
assert!(matches!(
consumer.poll_next(&waiter),
Poll::Ready(Err(crate::Error::Rolled))
));
}
#[test]
fn a_late_lower_group_is_still_reported() {
let track = moq_net::broadcast::Info::new()
.produce()
.create_track("test", None)
.unwrap();
let subscriber = track.subscribe(replaying());
for sequence in [1u64, 0] {
let mut flate = moq_flate::Encoder::new();
let mut group = track.create_group(moq_net::group::Info { sequence }).unwrap();
group
.write_frame(moq_net::Timestamp::now(), flate.frame(&[sequence as u8; 8]))
.unwrap();
group.finish().unwrap();
}
track.finish().unwrap();
let mut consumer = Consumer::new(subscriber, cfg(true));
let waiter = kio::Waiter::noop();
assert!(matches!(
consumer.poll_next(&waiter),
Poll::Ready(Ok(Some(payload))) if payload == vec![1u8; 8]
));
assert!(matches!(
consumer.poll_next(&waiter),
Poll::Ready(Err(crate::Error::Rolled))
));
}
#[test]
fn appending_after_finish_fails_on_every_clone() {
let (mut producer, _track) = producer(false);
let mut clone = producer.clone();
producer.append(&b"first"[..]).unwrap();
producer.finish().unwrap();
assert!(producer.append(&b"late"[..]).is_err());
assert!(clone.append(&b"late"[..]).is_err());
}
#[test]
fn a_finished_track_ends_the_consumer() {
let (mut producer, track) = producer(false);
producer.append(&b"only"[..]).unwrap();
producer.finish().unwrap();
let mut consumer = Consumer::new(track, Config::default());
let waiter = kio::Waiter::noop();
assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(_)))));
assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(None))));
}
}