Skip to main content

moq_json/stream/
mod.rs

1//! Append-log JSON publishing over [`moq-net`](moq_net) tracks.
2//!
3//! The counterpart to [`snapshot`](crate::snapshot) mode: instead of one JSON value updated over
4//! time, a stream is an ordered log of self-contained records. Every [`Producer::append`] writes one
5//! JSON object as one frame, and a [`Consumer`] yields every record in order.
6//!
7//! The whole log rides a **single group** that is never rolled: with
8//! [`Config::compression`] set to [`crate::Compression::Deflate`], that one
9//! group is one DEFLATE window, so every record
10//! compresses against all the earlier ones. There is deliberately no group rolling (and so no
11//! catch-up machinery): the only reason to roll would be moq-net's per-group frame cap, which
12//! isn't worth working around here. A caller that wants to bound the record rate throttles at
13//! the source (e.g. the timeline's segment cadence); a consumer that finds a gap can fetch or
14//! extrapolate.
15//!
16//! A record that cannot be encoded or written therefore ends the track rather than continuing in a
17//! second group: a log missing a record is not lossless, and a gap dressed up as a complete log is
18//! worse than a visible failure. A publisher with more to say opens a new track.
19//!
20//! That single group is what bounds the log's history. moq-net caps a group's cached bytes and
21//! frame count, and a consumer always starts at frame 0, so a write that would outgrow the
22//! budget aborts the group with [`moq_net::Error::GroupTooLarge`] rather than dropping a prefix
23//! some readers missed. (With compression the retained suffix would be undecodable anyway,
24//! since its DEFLATE window depends on the dropped prefix.) The live stream is therefore
25//! bounded history by design; deep history is served from a recording.
26//!
27//! # Choosing a layer
28//!
29//! [`Producer`] and [`Consumer`] own a track. [`Encoder`] and [`Decoder`] are the same logic
30//! without it, for when something else is already in charge of the track; they carry the shared
31//! DEFLATE window and nothing else, since a log has no group boundaries to report.
32
33pub 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	/// Demand follows the track's subscribers, so a producer can idle while nobody is watching.
62	#[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	/// Drain every record currently available without blocking.
98	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	/// Each record keeps its own capture time, and the returned size is the encoded frame.
120	#[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		// Never rolled: a single group holds the whole log.
161		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	/// A record the encoder rejects must not have published a group first: a live consumer would
209	/// advance into it and wait there even though nothing was ever appended. It still ends the
210	/// track, since the log is missing the record either way.
211	#[test]
212	fn a_rejected_record_does_not_open_a_group() {
213		// A map with non-string keys can't be represented as JSON, so serialization fails.
214		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	/// A track whose timescale is extreme enough that converting a wall-clock timestamp into it
235	/// overflows, so `write_frame` rejects every frame. That stands in for any post-`append_group`
236	/// write failure (the reported one is a frame over moq-net's 32 MB per-group cache) without
237	/// allocating 32 MB to provoke it.
238	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	/// A failed write must reach the consumer, not just the caller. A clean close drains a reader to
249	/// `None`, which is exactly what a completed log looks like, so a truncated log would be
250	/// indistinguishable from a whole one.
251	#[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	/// The track ends with the group, so nothing opens a second one and splits the log. The retry
267	/// reports the ended track rather than the [`Error::Desync`](crate::Error::Desync) the dropped
268	/// record left on the encoder, which says nothing about why the log stopped.
269	#[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		// The retry reports the abort rather than the `Error::Desync` the dropped record left on the
277		// encoder, which says nothing about why the log stopped.
278		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		// A subscriber taken after the abort still exists; it surfaces the failure on its first read,
284		// which is how a late reader learns the log is truncated.
285		let waiter = kio::Waiter::noop();
286		assert!(matches!(
287			producer.consume().poll_recv_group(&waiter),
288			Poll::Ready(Err(_))
289		));
290	}
291
292	/// A completed log is still readable, so finishing must not end the track the way an abort does.
293	/// The append that follows fails on the closed track without turning it into a failure.
294	#[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	/// A stream is one group. A publisher that opens a second lost whatever would have completed the
305	/// first, so the read reports that rather than handing back the remainder as a continuous log.
306	/// A boundary-only check would never look at the track again while the first group is open, so
307	/// this parks forever without the eager check. Written by hand because this producer never rolls.
308	#[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		// Ask for a replay window, so the first group is delivered rather than skipped by the
316		// subscriber's default max-age budget once a newer group exists.
317		let subscription = moq_net::track::Subscription::default().with_max_age(std::time::Duration::from_secs(30));
318		let subscriber = track.subscribe(subscription);
319
320		// Both groups stay open, the way a publisher writing to two at once leaves them.
321		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		// Sticky: a later read must not report the rest of the first group as a whole log.
343		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		// Each record is its own frame (one JSON object), and JSON escapes control characters, so a
355		// string value containing a newline round-trips cleanly.
356		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}