1pub mod consumer;
34mod decoder;
35mod encoder;
36pub mod producer;
37
38pub use consumer::Consumer;
39pub use decoder::Decoder;
40pub use encoder::{Config, Encoder, Pending};
41pub use producer::Producer;
42
43#[cfg(test)]
44mod test {
45 use std::task::Poll;
46
47 use serde_json::{Value, json};
48
49 use super::*;
50 use crate::Compression;
51
52 fn producer(config: Config) -> (Producer<Value>, moq_net::track::Subscriber) {
53 let track = moq_net::broadcast::Info::new()
54 .produce()
55 .create_track("test", None)
56 .unwrap();
57 let consumer = track.subscribe(None);
58 (Producer::new(track, config), consumer)
59 }
60
61 fn compressed() -> Config {
62 Config {
63 compression: Compression::Deflate,
64 }
65 }
66
67 fn consume(track: moq_net::track::Subscriber, compression: bool) -> Consumer<Value> {
68 Consumer::new(
69 track,
70 Config {
71 compression: if compression {
72 Compression::Deflate
73 } else {
74 Compression::None
75 },
76 },
77 )
78 }
79
80 fn drain(mut consumer: Consumer<Value>) -> Vec<Value> {
82 let waiter = kio::Waiter::noop();
83 let mut out = Vec::new();
84 while let Poll::Ready(Ok(Some(value))) = consumer.poll_next(&waiter) {
85 out.push(value);
86 }
87 out
88 }
89
90 #[test]
91 fn plaintext_roundtrip_in_order() {
92 let (mut producer, track) = producer(Config::default());
93 for n in 0..5 {
94 producer.append(&json!({ "n": n })).unwrap();
95 }
96 producer.finish().unwrap();
97
98 let records = drain(consume(track, false));
99 assert_eq!(records, (0..5).map(|n| json!({ "n": n })).collect::<Vec<_>>());
100 }
101
102 #[test]
103 fn compressed_roundtrip_in_order() {
104 let (mut producer, track) = producer(compressed());
105 for n in 0..20 {
106 producer.append(&json!({ "group": n, "pts": n * 2_000 })).unwrap();
107 }
108 producer.finish().unwrap();
109
110 let records = drain(consume(track, true));
111 assert_eq!(records.len(), 20);
112 assert_eq!(records[7], json!({ "group": 7, "pts": 14_000 }));
113 }
114
115 #[test]
116 fn all_records_ride_one_group() {
117 let (mut producer, track) = producer(compressed());
118 for n in 0..50 {
119 producer.append(&json!({ "n": n })).unwrap();
120 }
121 producer.finish().unwrap();
122
123 assert_eq!(track.latest(), Some(0));
125 assert_eq!(drain(consume(track, true)).len(), 50);
126 }
127
128 #[test]
129 fn live_consumer_sees_each_record() {
130 let (mut producer, track) = producer(compressed());
131 let mut consumer = consume(track, true);
132 let waiter = kio::Waiter::noop();
133
134 for n in 0..3 {
135 producer.append(&json!({ "n": n })).unwrap();
136 match consumer.poll_next(&waiter) {
137 Poll::Ready(Ok(Some(value))) => assert_eq!(value, json!({ "n": n })),
138 other => panic!("expected record, got {other:?}"),
139 }
140 }
141 assert!(matches!(consumer.poll_next(&waiter), Poll::Pending));
142 producer.finish().unwrap();
143 }
144
145 #[test]
146 fn shared_window_shrinks_repetitive_records() {
147 let (mut producer, track) = producer(compressed());
148 for n in 0..8 {
149 producer.append(&json!({ "group": n, "pts": n * 2_000 })).unwrap();
150 }
151 producer.finish().unwrap();
152
153 let waiter = kio::Waiter::noop();
154 let mut track = track.ordered();
155 let Poll::Ready(Ok(Some(mut group))) = track.poll_next_group(&waiter) else {
156 panic!("expected a group");
157 };
158 let mut sizes = Vec::new();
159 while let Poll::Ready(Ok(Some(frame))) = group.poll_read_frame(&waiter) {
160 sizes.push(frame.payload.len());
161 }
162 assert_eq!(sizes.len(), 8);
163 let raw = serde_json::to_vec(&json!({ "group": 7, "pts": 14_000 })).unwrap().len();
164 assert!(
165 *sizes.last().unwrap() < raw / 2,
166 "windowed record {} should be far below its raw size {raw}",
167 sizes.last().unwrap()
168 );
169 }
170
171 #[test]
175 fn a_rejected_record_does_not_open_a_group() {
176 let track = moq_net::broadcast::Info::new()
178 .produce()
179 .create_track("test", None)
180 .unwrap();
181 let mut subscriber = track.subscribe(None);
182 let mut producer = Producer::<std::collections::BTreeMap<(u8, u8), u8>>::new(track, Config::default());
183
184 let mut bad = std::collections::BTreeMap::new();
185 bad.insert((1, 2), 3);
186 assert!(producer.append(&bad).is_err());
187
188 assert_eq!(subscriber.latest(), None, "a rejected record opened a group");
189
190 let waiter = kio::Waiter::noop();
191 assert!(
192 matches!(subscriber.poll_recv_group(&waiter), Poll::Ready(Err(_))),
193 "the log is missing a record, so the track must end rather than stay writable"
194 );
195 }
196
197 fn rejecting_track() -> moq_net::track::Producer {
202 let mut info = moq_net::track::Info::default();
203 info.timescale = moq_net::Timescale::new((1u64 << 62) - 1).unwrap();
204
205 moq_net::broadcast::Info::new()
206 .produce()
207 .create_track("test", Some(info))
208 .unwrap()
209 }
210
211 #[test]
215 fn a_failed_write_aborts_the_track() {
216 let track = rejecting_track();
217 let mut subscriber = track.subscribe(None);
218 let mut producer = Producer::<Value>::new(track, Config::default());
219
220 assert!(matches!(producer.append(&json!({ "n": 1 })), Err(crate::Error::Net(_))));
221
222 let waiter = kio::Waiter::noop();
223 assert!(
224 matches!(subscriber.poll_recv_group(&waiter), Poll::Ready(Err(_))),
225 "a truncated log must surface an error rather than read as a completed one"
226 );
227 }
228
229 #[test]
233 fn a_failed_write_ends_the_track() {
234 let track = rejecting_track();
235 let mut producer = Producer::<Value>::new(track, compressed());
236
237 assert!(matches!(producer.append(&json!({ "n": 1 })), Err(crate::Error::Net(_))));
238
239 assert!(
242 matches!(producer.append(&json!({ "n": 2 })), Err(crate::Error::Net(_))),
243 "a second append must fail on the ended track rather than open another group"
244 );
245
246 let waiter = kio::Waiter::noop();
249 assert!(matches!(
250 producer.consume().poll_recv_group(&waiter),
251 Poll::Ready(Err(_))
252 ));
253 }
254
255 #[test]
258 fn appending_after_finish_fails_without_aborting() {
259 let (mut producer, _track) = producer(compressed());
260 producer.append(&json!({ "n": 0 })).unwrap();
261 producer.finish().unwrap();
262
263 assert!(producer.append(&json!({ "n": 1 })).is_err());
264 assert_eq!(drain(consume(producer.consume(), true)), vec![json!({ "n": 0 })]);
265 }
266
267 #[test]
272 fn a_second_group_is_reported_while_the_first_is_open() {
273 let track = moq_net::broadcast::Info::new()
274 .produce()
275 .create_track("test", None)
276 .unwrap();
277
278 let subscription = moq_net::track::Subscription::default().with_max_age(std::time::Duration::from_secs(30));
281 let subscriber = track.subscribe(subscription);
282
283 let mut first = track.append_group().unwrap();
285 first
286 .write_frame(moq_net::Timestamp::now(), br#"{"n":0}"#.as_slice())
287 .unwrap();
288 let mut second = track.append_group().unwrap();
289 second
290 .write_frame(moq_net::Timestamp::now(), br#"{"n":1}"#.as_slice())
291 .unwrap();
292
293 let mut consumer = consume(subscriber, false);
294 let waiter = kio::Waiter::noop();
295
296 assert!(matches!(
297 consumer.poll_next(&waiter),
298 Poll::Ready(Ok(Some(value))) if value == json!({ "n": 0 })
299 ));
300 assert!(matches!(
301 consumer.poll_next(&waiter),
302 Poll::Ready(Err(crate::Error::Rolled))
303 ));
304
305 first
307 .write_frame(moq_net::Timestamp::now(), br#"{"n":2}"#.as_slice())
308 .unwrap();
309 assert!(matches!(
310 consumer.poll_next(&waiter),
311 Poll::Ready(Err(crate::Error::Rolled))
312 ));
313 }
314
315 #[test]
316 fn embedded_newlines_survive() {
317 let (mut producer, track) = producer(compressed());
320 let value = json!({ "s": "line1\nline2\ttab", "u": "a\u{000a}b" });
321 for _ in 0..4 {
322 producer.append(&value).unwrap();
323 }
324 producer.finish().unwrap();
325
326 let records = drain(consume(track, true));
327 assert_eq!(records, vec![value.clone(), value.clone(), value.clone(), value]);
328 }
329}