moq_binary/snapshot/
mod.rs1use crate::Compression;
15
16pub mod consumer;
17pub mod producer;
18
19pub use consumer::Consumer;
20pub use producer::Producer;
21
22#[derive(Debug, Clone, Default)]
27#[non_exhaustive]
28pub struct Config {
29 pub compression: Compression,
35}
36
37#[cfg(test)]
38mod test {
39 use std::task::Poll;
40
41 use bytes::Bytes;
42
43 use super::*;
44
45 fn cfg(compression: bool) -> Config {
46 Config {
47 compression: if compression {
48 Compression::Deflate
49 } else {
50 Compression::None
51 },
52 }
53 }
54
55 fn producer(compression: bool) -> (Producer, moq_net::track::Subscriber) {
56 let track = moq_net::broadcast::Info::new()
57 .produce()
58 .create_track("test", None)
59 .unwrap();
60 let consumer = track.subscribe(None);
61 (Producer::new(track, cfg(compression)), consumer)
62 }
63
64 fn consume(track: moq_net::track::Subscriber, compression: bool) -> Consumer {
65 Consumer::new(track, cfg(compression))
66 }
67
68 fn drain(mut consumer: Consumer) -> Vec<Bytes> {
70 let waiter = kio::Waiter::noop();
71 let mut out = Vec::new();
72 while let Poll::Ready(Ok(Some(payload))) = consumer.poll_next(&waiter) {
73 out.push(payload);
74 }
75 out
76 }
77
78 #[test]
83 fn a_lost_group_waits_for_its_replacement() {
84 let track = moq_net::broadcast::Info::new()
85 .produce()
86 .create_track("test", None)
87 .unwrap();
88 let mut consumer = consume(track.subscribe(None), false);
89 let waiter = kio::Waiter::noop();
90
91 let mut group = track.create_group(moq_net::group::Info { sequence: 0 }).unwrap();
93 group
94 .write_frame(moq_net::Timestamp::now(), Bytes::from_static(b"one"))
95 .unwrap();
96 assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(v))) if v == "one"));
97 assert!(consumer.poll_next(&waiter).is_pending());
98
99 group.abort(moq_net::Error::Old).unwrap();
101 assert!(
102 consumer.poll_next(&waiter).is_pending(),
103 "a lost group must not end the reader"
104 );
105
106 let mut group = track.create_group(moq_net::group::Info { sequence: 1 }).unwrap();
108 group
109 .write_frame(moq_net::Timestamp::now(), Bytes::from_static(b"two"))
110 .unwrap();
111 assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(v))) if v == "two"));
112 }
113
114 #[test]
115 fn one_group_per_update() {
116 let (mut producer, track) = producer(false);
117 producer.update(&b"first"[..]).unwrap();
118 producer.update(&b"second"[..]).unwrap();
119 producer.finish().unwrap();
120
121 assert_eq!(track.latest(), Some(1));
123 assert_eq!(drain(consume(track, false)), vec![Bytes::from_static(b"second")]);
124 }
125
126 #[test]
127 fn live_consumer_sees_each_update() {
128 let (mut producer, track) = producer(false);
129 let mut consumer = consume(track, false);
130 let waiter = kio::Waiter::noop();
131
132 for n in 0..3u8 {
133 producer.update(vec![n]).unwrap();
134 match consumer.poll_next(&waiter) {
135 Poll::Ready(Ok(Some(payload))) => assert_eq!(&payload[..], &[n]),
136 other => panic!("expected value, got {other:?}"),
137 }
138 }
139 }
140
141 #[test]
142 fn compressed_roundtrip() {
143 let (mut producer, track) = producer(true);
144 let payload = Bytes::from(b"the quick brown fox".repeat(64));
145 producer.update(payload.clone()).unwrap();
146 producer.finish().unwrap();
147
148 assert_eq!(drain(consume(track, true)), vec![payload]);
149 }
150
151 #[test]
154 fn compression_shrinks_the_frame() {
155 let (mut producer, track) = producer(true);
156 let payload = Bytes::from(b"the quick brown fox".repeat(64));
157 producer.update(payload.clone()).unwrap();
158 producer.finish().unwrap();
159
160 let waiter = kio::Waiter::noop();
162 let Poll::Ready(Ok(Some(mut group))) = track.ordered().poll_next_group(&waiter) else {
163 panic!("expected a group");
164 };
165 let Poll::Ready(Ok(Some(frame))) = group.poll_read_frame(&waiter) else {
166 panic!("expected a frame");
167 };
168 assert!(
169 frame.payload.len() < payload.len() / 4,
170 "compressed frame {} should be far below the raw {}",
171 frame.payload.len(),
172 payload.len()
173 );
174 }
175
176 #[test]
180 fn a_rejected_update_leaves_the_previous_value_readable() {
181 let (mut producer, track) = producer(false);
182 producer.update(&b"keep"[..]).unwrap();
183
184 let oversized = Bytes::from(vec![0u8; moq_net::group::MAX_CACHE_BYTES as usize + 1]);
185 assert!(producer.update(oversized).is_err());
186 producer.finish().unwrap();
187
188 assert_eq!(drain(consume(track, false)), vec![Bytes::from_static(b"keep")]);
190 }
191
192 #[test]
195 fn updating_after_finish_fails_on_every_clone() {
196 let (mut producer, _track) = producer(false);
197 let mut clone = producer.clone();
198
199 producer.update(&b"first"[..]).unwrap();
200 producer.finish().unwrap();
201
202 assert!(producer.update(&b"late"[..]).is_err());
203 assert!(clone.update(&b"late"[..]).is_err());
204 }
205
206 #[test]
207 fn a_finished_track_ends_the_consumer() {
208 let (mut producer, track) = producer(false);
209 producer.update(&b"only"[..]).unwrap();
210 producer.finish().unwrap();
211
212 let mut consumer = consume(track, false);
213 let waiter = kio::Waiter::noop();
214 assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(_)))));
215 assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(None))));
216 }
217}