Skip to main content

moq_json/window/
mod.rs

1//! Sliding-window JSON publishing over [`moq-net`](moq_net) tracks.
2//!
3//! A window is an ordered run of records the publisher appends to the back of and drops from the
4//! front of. Unlike [`stream`](crate::stream), which preserves a log forever in one group, and
5//! [`snapshot`](crate::snapshot), which keeps only the latest value, a window keeps a bounded
6//! stretch of records and lets a reader join it at any point.
7//!
8//! # Why a log can't do this
9//!
10//! The obvious alternative is an append-only log that rolls its group and re-seeds the new one with
11//! the records it still holds. That breaks the reader: re-seeded records are indistinguishable from
12//! new ones, so a reader that was keeping up receives them twice. This mode exists to make the
13//! restatement explicit, so a reader can tell "you already have these" from "here is another one".
14//!
15//! # On the wire
16//!
17//! The first frame of every group names the retained `offset` and a decodable `records` suffix. Its
18//! optional `start` identifies the suffix when a checkpoint bound omits older records. Later frames
19//! are tagged `push` and `pop` ops, positional against the group header.
20//! Indices stop at 2^53 - 1, the largest integer represented exactly by both implementations.
21//!
22//! Trimming is therefore an op, not a group boundary. Dropping a record costs one small frame
23//! inside the shared compression window instead of a roll that would throw that window away.
24//!
25//! # Group boundaries are invisible
26//!
27//! The publisher rolls a group when the ops in it outgrow
28//! [`ProducerConfig::op_ratio`](ProducerConfig::op_ratio) times the header that opened it, exactly as
29//! [`snapshot`](crate::snapshot) rolls on its delta budget. That is purely a compression decision:
30//! there is no caller-driven cut and no age bound, and a [`Consumer`] never surfaces it. A header
31//! restating records a reader already has yields nothing, so however often the publisher rolls, the
32//! reader sees one continuous stream of [`Event`]s. [`ProducerConfig::checkpoint_records`] bounds
33//! the suffix repeated on each roll for a long-lived window.
34//!
35//! # What a reader is told
36//!
37//! A reader gets [`Event::Push`] when a record arrives, [`Event::Pop`] when a contiguous range
38//! leaves, and [`Event::Skip`] when a range was dropped before this reader saw it. A reader that
39//! keeps up sees pushes and pops; one that falls a group behind learns from the header's offset which
40//! records it will never get, rather than silently missing them.
41//!
42//! # Choosing a layer
43//!
44//! [`Producer`] and [`Consumer`] own a track. [`Encoder`] and [`Decoder`] are the same logic
45//! without it, for when something else is already in charge of the track; the encoder owns the
46//! retained window and says where the group boundaries fall, and the decoder turns frames into
47//! events.
48
49mod consumer;
50mod decoder;
51mod encoder;
52mod op;
53mod producer;
54
55pub use consumer::Consumer;
56pub use decoder::{ConsumerConfig, Decoder, Event, Group};
57pub use encoder::{Encoded, Encoder, Pending, ProducerConfig};
58pub use producer::Producer;
59
60#[cfg(test)]
61mod test {
62	use std::task::Poll;
63
64	use serde_json::{Value, json};
65
66	use super::*;
67
68	fn producer(config: ProducerConfig) -> (Producer<Value>, moq_net::track::Subscriber) {
69		let track = moq_net::broadcast::Info::new()
70			.produce()
71			.create_track("test", None)
72			.unwrap();
73		let consumer = track.subscribe(None);
74		(Producer::new(track, config), consumer)
75	}
76
77	#[test]
78	#[should_panic(expected = "checkpoint_records must be positive")]
79	fn zero_checkpoint_records_is_rejected() {
80		let config = ProducerConfig {
81			checkpoint_records: Some(0),
82			..Default::default()
83		};
84		let _ = Encoder::<Value>::new(config);
85	}
86
87	fn consumer(track: moq_net::track::Subscriber, compression: bool) -> Consumer<Value> {
88		Consumer::new(track, ConsumerConfig::default().with_compression(compression))
89	}
90
91	/// A track whose timestamp conversion rejects every frame after its group is published.
92	fn rejecting_track() -> moq_net::track::Producer {
93		let mut info = moq_net::track::Info::default();
94		info.timescale = moq_net::Timescale::new((1u64 << 62) - 1).unwrap();
95
96		moq_net::broadcast::Info::new()
97			.produce()
98			.create_track("test", Some(info))
99			.unwrap()
100	}
101
102	/// Drain every event currently available without blocking.
103	fn drain(consumer: &mut Consumer<Value>) -> Vec<Event<Value>> {
104		let waiter = kio::Waiter::noop();
105		let mut out = Vec::new();
106		while let Poll::Ready(Ok(Some(event))) = consumer.poll_next(&waiter) {
107			out.push(event);
108		}
109		out
110	}
111
112	fn rec(n: u64) -> Value {
113		json!({ "n": n })
114	}
115
116	/// A producer and a consumer that reads after every edit.
117	///
118	/// Polling as the publisher goes is what "keeping up" means: a consumer left until the end is a
119	/// whole group behind, and the default subscription abandons a group as soon as a newer one
120	/// exists, so it would resume at the newest header instead of reading the rolls in between.
121	struct Live {
122		producer: Producer<Value>,
123		consumer: Consumer<Value>,
124		events: Vec<Event<Value>>,
125	}
126
127	impl Live {
128		fn new(config: ProducerConfig) -> Self {
129			let compression = config.compression;
130			let (producer, track) = producer(config);
131			Self {
132				producer,
133				consumer: consumer(track, compression),
134				events: Vec::new(),
135			}
136		}
137
138		fn push(&mut self, n: u64) {
139			self.producer.push(&rec(n)).unwrap();
140			self.read();
141		}
142
143		fn pop(&mut self, count: u64) {
144			self.producer.pop(count).unwrap();
145			self.read();
146		}
147
148		fn read(&mut self) {
149			self.events.extend(drain(&mut self.consumer));
150		}
151
152		fn finish(self) -> Vec<Event<Value>> {
153			let Live {
154				mut producer,
155				mut consumer,
156				mut events,
157			} = self;
158			producer.finish().unwrap();
159			events.extend(drain(&mut consumer));
160			events
161		}
162
163		/// Just the indices pushed, in order.
164		fn pushed(events: &[Event<Value>]) -> Vec<u64> {
165			events
166				.iter()
167				.filter_map(|e| match e {
168					Event::Push { index, .. } => Some(*index),
169					_ => None,
170				})
171				.collect()
172		}
173	}
174
175	#[test]
176	fn push_and_pop_round_trip() {
177		let mut live = Live::new(ProducerConfig::default());
178		live.push(0);
179		live.push(1);
180		live.pop(1);
181		live.push(2);
182
183		assert_eq!(
184			live.finish(),
185			vec![
186				Event::Push {
187					index: 0,
188					value: rec(0)
189				},
190				Event::Push {
191					index: 1,
192					value: rec(1)
193				},
194				Event::Pop(0..1),
195				Event::Push {
196					index: 2,
197					value: rec(2)
198				},
199			]
200		);
201	}
202
203	#[test]
204	fn the_window_slides() {
205		let (mut producer, _track) = producer(ProducerConfig::default());
206		for n in 0..5 {
207			producer.push(&rec(n)).unwrap();
208			if n >= 2 {
209				producer.pop(1).unwrap();
210			}
211		}
212
213		// Three pops leave the two newest records, at indices 3 and 4.
214		assert_eq!(producer.range(), 3..5);
215		assert_eq!(producer.window(), vec![rec(3), rec(4)]);
216	}
217
218	#[test]
219	fn a_popped_record_is_never_restated() {
220		// Ops disabled, so every single edit is its own group restating the whole window.
221		let mut live = Live::new(ProducerConfig::default().with_op_ratio(0));
222		live.push(0);
223		live.push(1);
224		live.pop(1);
225		live.push(2);
226
227		// Every edit restates the window, yet a record already delivered is never pushed twice. That
228		// is the property an append-only log cannot provide.
229		assert_eq!(
230			live.finish(),
231			vec![
232				Event::Push {
233					index: 0,
234					value: rec(0)
235				},
236				Event::Push {
237					index: 1,
238					value: rec(1)
239				},
240				Event::Pop(0..1),
241				Event::Push {
242					index: 2,
243					value: rec(2)
244				},
245			]
246		);
247	}
248
249	#[test]
250	fn bounded_checkpoints_keep_a_following_consumer_contiguous() {
251		let config = ProducerConfig::default().with_op_ratio(0).with_checkpoint_records(2);
252		let mut live = Live::new(config);
253		for n in 0..6 {
254			live.push(n);
255		}
256
257		assert_eq!(live.producer.range(), 0..6);
258		assert_eq!(live.producer.window(), vec![rec(4), rec(5)]);
259		let events = live.finish();
260		assert_eq!(Live::pushed(&events), (0..6).collect::<Vec<_>>());
261		assert!(!events.iter().any(|event| matches!(event, Event::Skip(_))));
262	}
263
264	#[test]
265	fn a_late_consumer_skips_to_the_bounded_checkpoint() {
266		let config = ProducerConfig::default().with_op_ratio(0).with_checkpoint_records(2);
267		let mut encoder = Encoder::<Value>::new(config);
268		let mut latest = None;
269		for n in 0..5 {
270			let frame = encoder.push(&rec(n)).unwrap();
271			latest = Some(frame.payload.clone());
272			frame.commit();
273		}
274
275		let mut decoder = Decoder::<Value>::new(ConsumerConfig::default());
276		decoder.group().decode(&latest.unwrap()).unwrap();
277		assert_eq!(
278			std::iter::from_fn(|| decoder.next_event()).collect::<Vec<_>>(),
279			vec![
280				Event::Skip(0..3),
281				Event::Push {
282					index: 3,
283					value: rec(3)
284				},
285				Event::Push {
286					index: 4,
287					value: rec(4)
288				},
289			]
290		);
291	}
292
293	#[test]
294	fn pops_cross_the_omitted_checkpoint_prefix() {
295		let config = ProducerConfig::default().with_checkpoint_records(2);
296		let mut encoder = Encoder::<Value>::new(config);
297		for n in 0..5 {
298			encoder.push(&rec(n)).unwrap().commit();
299		}
300		assert_eq!(encoder.range(), 0..5);
301		assert_eq!(encoder.window(), vec![rec(3), rec(4)]);
302
303		encoder.pop(2).unwrap().unwrap().commit();
304		assert_eq!(encoder.range(), 2..5);
305		assert_eq!(encoder.window(), vec![rec(3), rec(4)]);
306
307		encoder.pop(2).unwrap().unwrap().commit();
308		assert_eq!(encoder.range(), 4..5);
309		assert_eq!(encoder.window(), vec![rec(4)]);
310	}
311
312	#[test]
313	fn a_fresh_consumer_adopts_the_offset_without_skipping_history() {
314		let track = moq_net::broadcast::Info::new()
315			.produce()
316			.create_track("test", None)
317			.unwrap();
318		let mut producer = Producer::<Value>::new(track, ProducerConfig::default().with_op_ratio(0));
319
320		for n in 0..5 {
321			producer.push(&rec(n)).unwrap();
322		}
323		producer.pop(3).unwrap();
324		let mut subscriber = producer.consume();
325		subscriber.set_groups(subscriber.latest().unwrap()..);
326		let mut fresh = consumer(subscriber, false);
327		producer.finish().unwrap();
328
329		// Joining at offset 3 must not report 3 skips for records that were never this reader's to
330		// miss: it simply starts where the window starts.
331		let events = drain(&mut fresh);
332		assert_eq!(
333			events,
334			vec![
335				Event::Push {
336					index: 3,
337					value: rec(3)
338				},
339				Event::Push {
340					index: 4,
341					value: rec(4)
342				}
343			]
344		);
345		assert!(!events.iter().any(|e| matches!(e, Event::Skip(_))));
346	}
347
348	#[test]
349	fn a_lagging_consumer_is_told_what_it_missed() {
350		// Ops are disabled, so every edit opens a new group. Feed the first two groups to the
351		// decoder, skip the middle groups as a lagging track subscriber would, then resume at the
352		// latest header.
353		let mut encoder = Encoder::<Value>::new(ProducerConfig::default().with_op_ratio(0));
354		let mut decoder = Decoder::<Value>::new(ConsumerConfig::default());
355		for n in 0..2 {
356			let frame = encoder.push(&rec(n)).unwrap();
357			let mut group = decoder.group();
358			group.decode(&frame.payload).unwrap();
359			frame.commit();
360		}
361		assert_eq!(
362			std::iter::from_fn(|| decoder.next_event()).collect::<Vec<_>>(),
363			vec![
364				Event::Push {
365					index: 0,
366					value: rec(0)
367				},
368				Event::Push {
369					index: 1,
370					value: rec(1)
371				}
372			]
373		);
374
375		let mut latest = None;
376		for n in 2..8 {
377			let frame = encoder.push(&rec(n)).unwrap();
378			frame.commit();
379
380			let frame = encoder.pop(1).unwrap().unwrap();
381			latest = Some(frame.payload.clone());
382			frame.commit();
383		}
384		let mut group = decoder.group();
385		group.decode(&latest.unwrap()).unwrap();
386
387		let events = std::iter::from_fn(|| decoder.next_event()).collect::<Vec<_>>();
388		let skipped: Vec<std::ops::Range<u64>> = events
389			.iter()
390			.filter_map(|e| match e {
391				Event::Skip(range) => Some(range.clone()),
392				_ => None,
393			})
394			.collect();
395
396		// Records 2..=5 existed but this reader will never receive them, and it is told so rather than
397		// silently jumping from 1 to 6.
398		assert!(!skipped.is_empty(), "expected skips, got {events:?}");
399		assert_eq!(skipped.first().map(|range| range.start), Some(2));
400
401		// Every index is still accounted for exactly once, in order.
402		let reported: Vec<u64> = events
403			.iter()
404			.flat_map(|e| match e {
405				Event::Push { index, .. } => vec![*index],
406				Event::Skip(range) => range.clone().collect(),
407				Event::Pop(_) => Vec::new(),
408			})
409			.collect();
410		assert!(reported.windows(2).all(|w| w[1] == w[0] + 1), "gaps in {reported:?}");
411	}
412
413	#[test]
414	fn consumer_resumes_at_a_checkpoint_after_losing_a_group() {
415		for err in [
416			moq_net::Error::Old,
417			moq_net::Error::Lagged,
418			moq_net::Error::Evicted,
419			moq_net::Error::GroupTooLarge,
420		] {
421			let track = moq_net::broadcast::Info::new()
422				.produce()
423				.create_track("test", None)
424				.unwrap();
425			let mut consumer = consumer(track.subscribe(None), false);
426			let mut encoder = Encoder::<Value>::new(ProducerConfig::default().with_op_ratio(0));
427
428			let first = encoder.push(&rec(0)).unwrap();
429			let payload = first.payload.clone();
430			first.commit();
431			let mut lost = track.append_group().unwrap();
432			lost.write_frame(moq_net::Timestamp::ZERO, payload).unwrap();
433			assert_eq!(
434				drain(&mut consumer),
435				vec![Event::Push {
436					index: 0,
437					value: rec(0)
438				}]
439			);
440			lost.abort(err).unwrap();
441
442			let second = encoder.push(&rec(1)).unwrap();
443			let payload = second.payload.clone();
444			second.commit();
445			let mut checkpoint = track.append_group().unwrap();
446			checkpoint.write_frame(moq_net::Timestamp::ZERO, payload).unwrap();
447			checkpoint.finish().unwrap();
448			track.finish().unwrap();
449
450			assert_eq!(
451				drain(&mut consumer),
452				vec![Event::Push {
453					index: 1,
454					value: rec(1)
455				}]
456			);
457		}
458	}
459
460	#[test]
461	fn compressed_round_trip_across_rolls() {
462		let mut live = Live::new(ProducerConfig::default().with_compression(true).with_op_ratio(1));
463		for n in 0..40 {
464			live.push(n);
465			if n >= 10 {
466				live.pop(1);
467			}
468		}
469
470		// A tight ratio rolls many times; every record still arrives exactly once, in order.
471		assert_eq!(Live::pushed(&live.finish()), (0..40).collect::<Vec<_>>());
472	}
473
474	#[test]
475	fn an_empty_pop_writes_nothing() {
476		let (mut producer, track) = producer(ProducerConfig::default());
477		producer.pop(5).unwrap();
478		producer.finish().unwrap();
479
480		// Nothing was ever pushed, so there is nothing to drop and no group to publish.
481		assert_eq!(track.latest(), None);
482	}
483
484	#[test]
485	fn a_rejected_edit_leaves_the_window_unchanged() {
486		let track = rejecting_track();
487		let mut subscriber = track.subscribe(None).ordered();
488		let mut producer = Producer::<Value>::new(track, ProducerConfig::default());
489
490		assert!(producer.push(&rec(1)).is_err());
491		assert_eq!(producer.range(), 0..0);
492		assert!(producer.window().is_empty());
493
494		let waiter = kio::Waiter::noop();
495		let Poll::Ready(Ok(Some(mut group))) = subscriber.poll_next_group(&waiter) else {
496			panic!("the rejected group's header was published");
497		};
498		assert!(matches!(group.poll_read_frame(&waiter), Poll::Ready(Ok(None))));
499	}
500
501	#[test]
502	fn the_handle_remains_usable_after_finish() {
503		let (mut producer, track) = producer(ProducerConfig::default());
504		producer.push(&rec(1)).unwrap();
505		producer.finish().unwrap();
506
507		assert_eq!(producer.window(), vec![rec(1)]);
508		assert_eq!(producer.range(), 0..1);
509		assert_eq!(producer.consume().latest(), track.latest());
510		producer.finish().unwrap();
511		assert!(matches!(
512			producer.push(&rec(2)),
513			Err(crate::Error::Net(moq_net::Error::Closed))
514		));
515	}
516
517	#[test]
518	fn writes_after_another_clone_finishes_are_rejected() {
519		let (mut producer, _track) = producer(ProducerConfig::default());
520		producer.push(&rec(1)).unwrap();
521		producer.clone().finish().unwrap();
522
523		assert!(matches!(
524			producer.push(&rec(2)),
525			Err(crate::Error::Net(moq_net::Error::Closed))
526		));
527		assert!(matches!(
528			producer.pop(1),
529			Err(crate::Error::Net(moq_net::Error::Closed))
530		));
531		assert_eq!(producer.window(), vec![rec(1)]);
532	}
533
534	#[test]
535	fn a_pop_is_clamped_to_the_window() {
536		let mut live = Live::new(ProducerConfig::default());
537		live.push(0);
538		live.pop(9);
539		live.push(1);
540
541		assert_eq!(
542			live.finish(),
543			vec![
544				Event::Push {
545					index: 0,
546					value: rec(0)
547				},
548				Event::Pop(0..1),
549				Event::Push {
550					index: 1,
551					value: rec(1)
552				},
553			]
554		);
555	}
556
557	#[test]
558	fn a_large_gap_is_one_skip_event() {
559		let mut decoder = Decoder::<Value>::new(ConsumerConfig::default());
560		let mut group = decoder.group();
561		group.decode(br#"{"offset":0,"records":[]}"#).unwrap();
562		let mut group = decoder.group();
563		group.decode(br#"{"offset":9007199254740991,"records":[]}"#).unwrap();
564
565		assert_eq!(decoder.next_event(), Some(Event::Skip(0..super::encoder::MAX_INDEX)));
566		assert_eq!(decoder.next_event(), None);
567	}
568
569	#[test]
570	fn indices_must_fit_the_shared_safe_integer_range() {
571		let mut decoder = Decoder::<Value>::new(ConsumerConfig::default());
572		let mut group = decoder.group();
573		assert!(group.decode(br#"{"offset":9007199254740992,"records":[]}"#).is_err());
574
575		let mut decoder = Decoder::<Value>::new(ConsumerConfig::default());
576		let mut group = decoder.group();
577		group.decode(br#"{"offset":9007199254740991,"records":[]}"#).unwrap();
578		assert!(group.decode(br#"{"push":null}"#).is_err());
579	}
580
581	#[test]
582	fn a_checkpoint_cannot_start_before_the_window() {
583		let mut decoder = Decoder::<Value>::new(ConsumerConfig::default());
584		let mut group = decoder.group();
585		assert!(group.decode(br#"{"offset":2,"start":1,"records":[]}"#).is_err());
586	}
587
588	#[test]
589	fn every_group_requires_a_header() {
590		let mut decoder = Decoder::<Value>::new(ConsumerConfig::default());
591		let mut group = decoder.group();
592		group.decode(br#"{"offset":0,"records":[]}"#).unwrap();
593		let mut group = decoder.group();
594
595		assert!(group.decode(br#"{"push":null}"#).is_err());
596	}
597
598	#[test]
599	fn a_header_is_only_valid_as_frame_zero() {
600		let mut decoder = Decoder::<Value>::new(ConsumerConfig::default());
601		let mut group = decoder.group();
602		group.decode(br#"{"offset":0,"records":[]}"#).unwrap();
603
604		assert!(group.decode(br#"{"offset":0,"records":[]}"#).is_err());
605	}
606
607	#[test]
608	fn rolling_is_invisible_to_the_consumer() {
609		// The same edits, framed two ways: one group for everything, versus a roll per edit.
610		let edits = |ratio: u32| {
611			let mut live = Live::new(ProducerConfig::default().with_op_ratio(ratio));
612			for n in 0..6 {
613				live.push(n);
614				if n >= 3 {
615					live.pop(1);
616				}
617			}
618			live.finish()
619		};
620
621		assert_eq!(edits(1_000), edits(0));
622	}
623}