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 consume(track: moq_net::track::Subscriber, compression: bool) -> Consumer {
Consumer::new(track, cfg(compression))
}
fn drain(mut consumer: Consumer) -> Vec<Bytes> {
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
}
#[test]
fn a_lost_group_waits_for_its_replacement() {
let track = moq_net::broadcast::Info::new()
.produce()
.create_track("test", None)
.unwrap();
let mut consumer = consume(track.subscribe(None), false);
let waiter = kio::Waiter::noop();
let mut group = track.create_group(moq_net::group::Info { sequence: 0 }).unwrap();
group
.write_frame(moq_net::Timestamp::now(), Bytes::from_static(b"one"))
.unwrap();
assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(v))) if v == "one"));
assert!(consumer.poll_next(&waiter).is_pending());
group.abort(moq_net::Error::Old).unwrap();
assert!(
consumer.poll_next(&waiter).is_pending(),
"a lost group must not end the reader"
);
let mut group = track.create_group(moq_net::group::Info { sequence: 1 }).unwrap();
group
.write_frame(moq_net::Timestamp::now(), Bytes::from_static(b"two"))
.unwrap();
assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(v))) if v == "two"));
}
#[test]
fn one_group_per_update() {
let (mut producer, track) = producer(false);
producer.update(&b"first"[..]).unwrap();
producer.update(&b"second"[..]).unwrap();
producer.finish().unwrap();
assert_eq!(track.latest(), Some(1));
assert_eq!(drain(consume(track, false)), vec![Bytes::from_static(b"second")]);
}
#[test]
fn live_consumer_sees_each_update() {
let (mut producer, track) = producer(false);
let mut consumer = consume(track, false);
let waiter = kio::Waiter::noop();
for n in 0..3u8 {
producer.update(vec![n]).unwrap();
match consumer.poll_next(&waiter) {
Poll::Ready(Ok(Some(payload))) => assert_eq!(&payload[..], &[n]),
other => panic!("expected value, got {other:?}"),
}
}
}
#[test]
fn compressed_roundtrip() {
let (mut producer, track) = producer(true);
let payload = Bytes::from(b"the quick brown fox".repeat(64));
producer.update(payload.clone()).unwrap();
producer.finish().unwrap();
assert_eq!(drain(consume(track, true)), vec![payload]);
}
#[test]
fn compression_shrinks_the_frame() {
let (mut producer, track) = producer(true);
let payload = Bytes::from(b"the quick brown fox".repeat(64));
producer.update(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 Poll::Ready(Ok(Some(frame))) = group.poll_read_frame(&waiter) else {
panic!("expected a frame");
};
assert!(
frame.payload.len() < payload.len() / 4,
"compressed frame {} should be far below the raw {}",
frame.payload.len(),
payload.len()
);
}
#[test]
fn a_rejected_update_leaves_the_previous_value_readable() {
let (mut producer, track) = producer(false);
producer.update(&b"keep"[..]).unwrap();
let oversized = Bytes::from(vec![0u8; moq_net::group::MAX_CACHE_BYTES as usize + 1]);
assert!(producer.update(oversized).is_err());
producer.finish().unwrap();
assert_eq!(drain(consume(track, false)), vec![Bytes::from_static(b"keep")]);
}
#[test]
fn updating_after_finish_fails_on_every_clone() {
let (mut producer, _track) = producer(false);
let mut clone = producer.clone();
producer.update(&b"first"[..]).unwrap();
producer.finish().unwrap();
assert!(producer.update(&b"late"[..]).is_err());
assert!(clone.update(&b"late"[..]).is_err());
}
#[test]
fn a_finished_track_ends_the_consumer() {
let (mut producer, track) = producer(false);
producer.update(&b"only"[..]).unwrap();
producer.finish().unwrap();
let mut consumer = consume(track, false);
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))));
}
}