1use crate::Compression;
21
22pub mod consumer;
23pub mod producer;
24
25pub use consumer::Consumer;
26pub use producer::Producer;
27
28#[derive(Debug, Clone, Default)]
33#[non_exhaustive]
34pub struct Config {
35 pub compression: Compression,
41}
42
43#[cfg(test)]
44mod test {
45 use std::task::Poll;
46
47 use bytes::Bytes;
48
49 use super::*;
50
51 fn cfg(compression: bool) -> Config {
52 Config {
53 compression: if compression {
54 Compression::Deflate
55 } else {
56 Compression::None
57 },
58 }
59 }
60
61 fn producer(compression: bool) -> (Producer, moq_net::track::Subscriber) {
62 let track = moq_net::broadcast::Info::new()
63 .produce()
64 .create_track("test", None)
65 .unwrap();
66 let consumer = track.subscribe(None);
67 (Producer::new(track, cfg(compression)), consumer)
68 }
69
70 fn drain(track: moq_net::track::Subscriber, compression: bool) -> Vec<Bytes> {
72 let mut consumer = Consumer::new(track, cfg(compression));
73 let waiter = kio::Waiter::noop();
74 let mut out = Vec::new();
75 while let Poll::Ready(Ok(Some(payload))) = consumer.poll_next(&waiter) {
76 out.push(payload);
77 }
78 out
79 }
80
81 fn payloads(count: u8) -> Vec<Bytes> {
82 (0..count).map(|n| Bytes::from(vec![n; 16])).collect()
83 }
84
85 #[test]
86 fn every_payload_survives_in_order() {
87 let (mut producer, track) = producer(false);
88 let expected = payloads(5);
89 for payload in &expected {
90 producer.append(payload.clone()).unwrap();
91 }
92 producer.finish().unwrap();
93
94 assert_eq!(track.latest(), Some(0));
96 assert_eq!(drain(track, false), expected);
97 }
98
99 #[test]
100 fn compressed_roundtrip_in_order() {
101 let (mut producer, track) = producer(true);
102 let expected = payloads(20);
103 for payload in &expected {
104 producer.append(payload.clone()).unwrap();
105 }
106 producer.finish().unwrap();
107
108 assert_eq!(drain(track, true), expected);
109 }
110
111 #[test]
113 fn the_shared_window_shrinks_repetitive_payloads() {
114 let (mut producer, track) = producer(true);
115 let payload = Bytes::from(b"the quick brown fox".repeat(16));
116 for _ in 0..8 {
117 producer.append(payload.clone()).unwrap();
118 }
119 producer.finish().unwrap();
120
121 let waiter = kio::Waiter::noop();
122 let Poll::Ready(Ok(Some(mut group))) = track.ordered().poll_next_group(&waiter) else {
123 panic!("expected a group");
124 };
125
126 let mut sizes = Vec::new();
127 while let Poll::Ready(Ok(Some(frame))) = group.poll_read_frame(&waiter) {
128 sizes.push(frame.payload.len());
129 }
130
131 assert_eq!(sizes.len(), 8);
132 assert!(
133 *sizes.last().unwrap() < sizes[0],
134 "windowed frame {} should be below the first {}",
135 sizes.last().unwrap(),
136 sizes[0]
137 );
138 }
139
140 #[test]
142 fn clones_append_into_one_log() {
143 let (mut producer, track) = producer(true);
144 let mut clone = producer.clone();
145
146 producer.append(&b"a"[..]).unwrap();
147 clone.append(&b"b"[..]).unwrap();
148 producer.append(&b"c"[..]).unwrap();
149 producer.finish().unwrap();
150
151 assert_eq!(
152 drain(track, true),
153 vec![
154 Bytes::from_static(b"a"),
155 Bytes::from_static(b"b"),
156 Bytes::from_static(b"c")
157 ]
158 );
159 }
160
161 fn rejecting_track() -> moq_net::track::Producer {
166 let mut info = moq_net::track::Info::default();
167 info.timescale = moq_net::Timescale::new((1u64 << 62) - 1).unwrap();
168
169 moq_net::broadcast::Info::new()
170 .produce()
171 .create_track("test", Some(info))
172 .unwrap()
173 }
174
175 #[test]
179 fn a_failed_write_aborts_the_track() {
180 let track = rejecting_track();
181 let mut subscriber = track.subscribe(None);
182 let mut producer = Producer::new(track, Config::default());
183
184 assert!(producer.append(&b"rejected"[..]).is_err());
185
186 let waiter = kio::Waiter::noop();
189 assert!(
190 matches!(subscriber.poll_recv_group(&waiter), Poll::Ready(Err(_))),
191 "a truncated log must surface an error rather than read as a completed one"
192 );
193 }
194
195 #[test]
199 fn an_undecodable_record_ends_the_track() {
200 let track = moq_net::broadcast::Info::new()
201 .produce()
202 .create_track("test", None)
203 .unwrap();
204 let mut subscriber = track.subscribe(None);
205 let mut producer = Producer::new(track, cfg(true));
206
207 assert!(producer.is_used());
208 let oversized = Bytes::from(vec![0u8; moq_flate::DEFAULT_MAX_FRAME_SIZE as usize + 1]);
209 assert!(matches!(
210 producer.append(oversized),
211 Err(crate::Error::Flate(moq_flate::Error::TooLarge(_)))
212 ));
213
214 assert!(!producer.is_used());
215
216 let waiter = kio::Waiter::noop();
218 assert!(matches!(subscriber.poll_recv_group(&waiter), Poll::Ready(Err(_))));
219 assert!(producer.append(&b"after"[..]).is_err());
220 }
221
222 #[test]
226 fn a_reader_inside_the_group_sees_the_real_error() {
227 let track = moq_net::broadcast::Info::new()
228 .produce()
229 .create_track("test", None)
230 .unwrap();
231 let mut subscriber = track.subscribe(None);
232 let mut producer = Producer::new(track, Config::default());
233
234 producer.append(&b"first"[..]).unwrap();
235
236 let waiter = kio::Waiter::noop();
238 let Poll::Ready(Ok(Some(mut group))) = subscriber.poll_recv_group(&waiter) else {
239 panic!("the group was published, so a subscriber sees it");
240 };
241 assert!(matches!(group.poll_read_frame(&waiter), Poll::Ready(Ok(Some(_)))));
242
243 let oversized = Bytes::from(vec![0u8; moq_net::group::MAX_CACHE_BYTES as usize + 1]);
245 assert!(producer.append(oversized).is_err());
246
247 match group.poll_read_frame(&waiter) {
248 Poll::Ready(Err(err)) => assert!(
249 matches!(err, moq_net::Error::FrameTooLarge),
250 "the reader should see the write's own error, got {err:?}"
251 ),
252 other => panic!("expected the write error, got {other:?}"),
253 }
254 }
255
256 fn replaying() -> moq_net::track::Subscription {
260 moq_net::track::Subscription::default().with_max_age(std::time::Duration::from_secs(30))
261 }
262
263 #[test]
265 fn a_failed_write_ends_the_track() {
266 let track = rejecting_track();
267 let mut producer = Producer::new(track, Config::default());
268
269 assert!(producer.append(&b"rejected"[..]).is_err());
270 assert!(
271 producer.append(&b"again"[..]).is_err(),
272 "a second append must fail on the closed track rather than open another group"
273 );
274 }
275
276 #[test]
280 fn a_second_group_is_a_rolled_log() {
281 let track = moq_net::broadcast::Info::new()
282 .produce()
283 .create_track("test", None)
284 .unwrap();
285 let subscriber = track.subscribe(replaying());
286
287 for pair in payloads(4).chunks(2) {
288 let mut flate = moq_flate::Encoder::new();
290 let mut group = track.append_group().unwrap();
291 for payload in pair {
292 group
293 .write_frame(moq_net::Timestamp::now(), flate.frame(payload))
294 .unwrap();
295 }
296 group.finish().unwrap();
297 }
298 track.finish().unwrap();
299
300 let mut consumer = Consumer::new(subscriber, cfg(true));
301 let waiter = kio::Waiter::noop();
302
303 assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(_)))));
305 assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(_)))));
306
307 assert!(matches!(
309 consumer.poll_next(&waiter),
310 Poll::Ready(Err(crate::Error::Rolled))
311 ));
312 }
313
314 #[test]
318 fn a_second_group_is_reported_while_the_first_is_open() {
319 let track = moq_net::broadcast::Info::new()
320 .produce()
321 .create_track("test", None)
322 .unwrap();
323 let subscriber = track.subscribe(replaying());
324
325 let mut first = track.append_group().unwrap();
327 first.write_frame(moq_net::Timestamp::now(), &b"first"[..]).unwrap();
328 let mut second = track.append_group().unwrap();
329 second.write_frame(moq_net::Timestamp::now(), &b"second"[..]).unwrap();
330
331 let mut consumer = Consumer::new(subscriber, Config::default());
332 let waiter = kio::Waiter::noop();
333
334 assert!(matches!(
335 consumer.poll_next(&waiter),
336 Poll::Ready(Ok(Some(payload))) if payload == Bytes::from_static(b"first")
337 ));
338 assert!(matches!(
339 consumer.poll_next(&waiter),
340 Poll::Ready(Err(crate::Error::Rolled))
341 ));
342
343 first.write_frame(moq_net::Timestamp::now(), &b"more"[..]).unwrap();
345 assert!(matches!(
346 consumer.poll_next(&waiter),
347 Poll::Ready(Err(crate::Error::Rolled))
348 ));
349 }
350
351 #[test]
355 fn a_late_lower_group_is_still_reported() {
356 let track = moq_net::broadcast::Info::new()
357 .produce()
358 .create_track("test", None)
359 .unwrap();
360 let subscriber = track.subscribe(replaying());
361
362 for sequence in [1u64, 0] {
364 let mut flate = moq_flate::Encoder::new();
365 let mut group = track.create_group(moq_net::group::Info { sequence }).unwrap();
366 group
367 .write_frame(moq_net::Timestamp::now(), flate.frame(&[sequence as u8; 8]))
368 .unwrap();
369 group.finish().unwrap();
370 }
371 track.finish().unwrap();
372
373 let mut consumer = Consumer::new(subscriber, cfg(true));
374 let waiter = kio::Waiter::noop();
375
376 assert!(matches!(
377 consumer.poll_next(&waiter),
378 Poll::Ready(Ok(Some(payload))) if payload == vec![1u8; 8]
379 ));
380 assert!(matches!(
381 consumer.poll_next(&waiter),
382 Poll::Ready(Err(crate::Error::Rolled))
383 ));
384 }
385
386 #[test]
389 fn appending_after_finish_fails_on_every_clone() {
390 let (mut producer, _track) = producer(false);
391 let mut clone = producer.clone();
392
393 producer.append(&b"first"[..]).unwrap();
394 producer.finish().unwrap();
395
396 assert!(producer.append(&b"late"[..]).is_err());
397 assert!(clone.append(&b"late"[..]).is_err());
398 }
399
400 #[test]
401 fn a_finished_track_ends_the_consumer() {
402 let (mut producer, track) = producer(false);
403 producer.append(&b"only"[..]).unwrap();
404 producer.finish().unwrap();
405
406 let mut consumer = Consumer::new(track, Config::default());
407 let waiter = kio::Waiter::noop();
408 assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(_)))));
409 assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(None))));
410 }
411}