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 #[test]
63 fn demand_follows_subscribers() {
64 let (producer, consumer) = producer(Config::default());
65 let demand = producer.demand();
66 let waiter = kio::Waiter::noop();
67 assert!(matches!(demand.poll_used(&waiter), Poll::Ready(Ok(()))));
68
69 drop(consumer);
70 assert!(matches!(demand.poll_unused(&waiter), Poll::Ready(Ok(()))));
71 assert!(demand.poll_used(&waiter).is_pending());
72
73 let _consumer = producer.consume();
74 assert!(matches!(demand.poll_used(&waiter), Poll::Ready(Ok(()))));
75 assert!(demand.poll_unused(&waiter).is_pending());
76 }
77
78 fn compressed() -> Config {
79 Config {
80 compression: Compression::Deflate,
81 }
82 }
83
84 fn consume(track: moq_net::track::Subscriber, compression: bool) -> Consumer<Value> {
85 Consumer::new(
86 track,
87 Config {
88 compression: if compression {
89 Compression::Deflate
90 } else {
91 Compression::None
92 },
93 },
94 )
95 }
96
97 fn drain(mut consumer: Consumer<Value>) -> Vec<Value> {
99 let waiter = kio::Waiter::noop();
100 let mut out = Vec::new();
101 while let Poll::Ready(Ok(Some(value))) = consumer.poll_next(&waiter) {
102 out.push(value);
103 }
104 out
105 }
106
107 #[test]
108 fn plaintext_roundtrip_in_order() {
109 let (mut producer, track) = producer(Config::default());
110 for n in 0..5 {
111 producer.append(&json!({ "n": n })).unwrap();
112 }
113 producer.finish().unwrap();
114
115 let records = drain(consume(track, false));
116 assert_eq!(records, (0..5).map(|n| json!({ "n": n })).collect::<Vec<_>>());
117 }
118
119 #[test]
121 fn a_stamped_append_writes_its_capture_time() {
122 let (mut producer, _track) = producer(compressed());
123 let mut groups = producer.consume();
124 let captured = moq_net::Timestamp::from_millis(1_234).unwrap();
125 let record = json!({ "n": 1 });
126 let size = producer.append(moq_net::Timed::from(&record).at(captured)).unwrap();
127
128 let waiter = kio::Waiter::noop();
129 let Poll::Ready(Ok(Some(mut group))) = groups.poll_recv_group(&waiter) else {
130 panic!("expected a group");
131 };
132 let Poll::Ready(Ok(Some(frame))) = group.poll_read_frame(&waiter) else {
133 panic!("expected a frame");
134 };
135 assert_eq!(frame.timestamp.as_micros(), captured.as_micros());
136 assert_eq!(size, frame.payload.len());
137 }
138
139 #[test]
140 fn compressed_roundtrip_in_order() {
141 let (mut producer, track) = producer(compressed());
142 for n in 0..20 {
143 producer.append(&json!({ "group": n, "pts": n * 2_000 })).unwrap();
144 }
145 producer.finish().unwrap();
146
147 let records = drain(consume(track, true));
148 assert_eq!(records.len(), 20);
149 assert_eq!(records[7], json!({ "group": 7, "pts": 14_000 }));
150 }
151
152 #[test]
153 fn all_records_ride_one_group() {
154 let (mut producer, track) = producer(compressed());
155 for n in 0..50 {
156 producer.append(&json!({ "n": n })).unwrap();
157 }
158 producer.finish().unwrap();
159
160 assert_eq!(track.latest(), Some(0));
162 assert_eq!(drain(consume(track, true)).len(), 50);
163 }
164
165 #[test]
166 fn live_consumer_sees_each_record() {
167 let (mut producer, track) = producer(compressed());
168 let mut consumer = consume(track, true);
169 let waiter = kio::Waiter::noop();
170
171 for n in 0..3 {
172 producer.append(&json!({ "n": n })).unwrap();
173 match consumer.poll_next(&waiter) {
174 Poll::Ready(Ok(Some(value))) => assert_eq!(value, json!({ "n": n })),
175 other => panic!("expected record, got {other:?}"),
176 }
177 }
178 assert!(matches!(consumer.poll_next(&waiter), Poll::Pending));
179 producer.finish().unwrap();
180 }
181
182 #[test]
183 fn shared_window_shrinks_repetitive_records() {
184 let (mut producer, track) = producer(compressed());
185 for n in 0..8 {
186 producer.append(&json!({ "group": n, "pts": n * 2_000 })).unwrap();
187 }
188 producer.finish().unwrap();
189
190 let waiter = kio::Waiter::noop();
191 let mut track = track.ordered();
192 let Poll::Ready(Ok(Some(mut group))) = track.poll_next_group(&waiter) else {
193 panic!("expected a group");
194 };
195 let mut sizes = Vec::new();
196 while let Poll::Ready(Ok(Some(frame))) = group.poll_read_frame(&waiter) {
197 sizes.push(frame.payload.len());
198 }
199 assert_eq!(sizes.len(), 8);
200 let raw = serde_json::to_vec(&json!({ "group": 7, "pts": 14_000 })).unwrap().len();
201 assert!(
202 *sizes.last().unwrap() < raw / 2,
203 "windowed record {} should be far below its raw size {raw}",
204 sizes.last().unwrap()
205 );
206 }
207
208 #[test]
212 fn a_rejected_record_does_not_open_a_group() {
213 let track = moq_net::broadcast::Info::new()
215 .produce()
216 .create_track("test", None)
217 .unwrap();
218 let mut subscriber = track.subscribe(None);
219 let mut producer = Producer::<std::collections::BTreeMap<(u8, u8), u8>>::new(track, Config::default());
220
221 let mut bad = std::collections::BTreeMap::new();
222 bad.insert((1, 2), 3);
223 assert!(producer.append(&bad).is_err());
224
225 assert_eq!(subscriber.latest(), None, "a rejected record opened a group");
226
227 let waiter = kio::Waiter::noop();
228 assert!(
229 matches!(subscriber.poll_recv_group(&waiter), Poll::Ready(Err(_))),
230 "the log is missing a record, so the track must end rather than stay writable"
231 );
232 }
233
234 fn rejecting_track() -> moq_net::track::Producer {
239 let mut info = moq_net::track::Info::default();
240 info.timescale = moq_net::Timescale::new((1u64 << 62) - 1).unwrap();
241
242 moq_net::broadcast::Info::new()
243 .produce()
244 .create_track("test", Some(info))
245 .unwrap()
246 }
247
248 #[test]
252 fn a_failed_write_aborts_the_track() {
253 let track = rejecting_track();
254 let mut subscriber = track.subscribe(None);
255 let mut producer = Producer::<Value>::new(track, Config::default());
256
257 assert!(matches!(producer.append(&json!({ "n": 1 })), Err(crate::Error::Net(_))));
258
259 let waiter = kio::Waiter::noop();
260 assert!(
261 matches!(subscriber.poll_recv_group(&waiter), Poll::Ready(Err(_))),
262 "a truncated log must surface an error rather than read as a completed one"
263 );
264 }
265
266 #[test]
270 fn a_failed_write_ends_the_track() {
271 let track = rejecting_track();
272 let mut producer = Producer::<Value>::new(track, compressed());
273
274 assert!(matches!(producer.append(&json!({ "n": 1 })), Err(crate::Error::Net(_))));
275
276 assert!(
279 matches!(producer.append(&json!({ "n": 2 })), Err(crate::Error::Net(_))),
280 "a second append must fail on the ended track rather than open another group"
281 );
282
283 let waiter = kio::Waiter::noop();
286 assert!(matches!(
287 producer.consume().poll_recv_group(&waiter),
288 Poll::Ready(Err(_))
289 ));
290 }
291
292 #[test]
295 fn appending_after_finish_fails_without_aborting() {
296 let (mut producer, _track) = producer(compressed());
297 producer.append(&json!({ "n": 0 })).unwrap();
298 producer.finish().unwrap();
299
300 assert!(producer.append(&json!({ "n": 1 })).is_err());
301 assert_eq!(drain(consume(producer.consume(), true)), vec![json!({ "n": 0 })]);
302 }
303
304 #[test]
309 fn a_second_group_is_reported_while_the_first_is_open() {
310 let track = moq_net::broadcast::Info::new()
311 .produce()
312 .create_track("test", None)
313 .unwrap();
314
315 let subscription = moq_net::track::Subscription::default().with_max_age(std::time::Duration::from_secs(30));
318 let subscriber = track.subscribe(subscription);
319
320 let mut first = track.append_group().unwrap();
322 first
323 .write_frame(moq_net::Timestamp::now(), br#"{"n":0}"#.as_slice())
324 .unwrap();
325 let mut second = track.append_group().unwrap();
326 second
327 .write_frame(moq_net::Timestamp::now(), br#"{"n":1}"#.as_slice())
328 .unwrap();
329
330 let mut consumer = consume(subscriber, false);
331 let waiter = kio::Waiter::noop();
332
333 assert!(matches!(
334 consumer.poll_next(&waiter),
335 Poll::Ready(Ok(Some(value))) if value == json!({ "n": 0 })
336 ));
337 assert!(matches!(
338 consumer.poll_next(&waiter),
339 Poll::Ready(Err(crate::Error::Rolled))
340 ));
341
342 first
344 .write_frame(moq_net::Timestamp::now(), br#"{"n":2}"#.as_slice())
345 .unwrap();
346 assert!(matches!(
347 consumer.poll_next(&waiter),
348 Poll::Ready(Err(crate::Error::Rolled))
349 ));
350 }
351
352 #[test]
353 fn embedded_newlines_survive() {
354 let (mut producer, track) = producer(compressed());
357 let value = json!({ "s": "line1\nline2\ttab", "u": "a\u{000a}b" });
358 for _ in 0..4 {
359 producer.append(&value).unwrap();
360 }
361 producer.finish().unwrap();
362
363 let records = drain(consume(track, true));
364 assert_eq!(records, vec![value.clone(), value.clone(), value.clone(), value]);
365 }
366}