Skip to main content

moq_binary/snapshot/
mod.rs

1//! Lossy latest-value binary publishing over [`moq-net`](moq_net) tracks.
2//!
3//! One opaque value updated over time, for consumers that only care about the current state (a
4//! poster image, a serialized state blob). This mode is **lossy** by design: a consumer yields only
5//! the most recent value. A late joiner (or a consumer that falls behind) jumps straight to the
6//! newest group, and older groups are dropped entirely. For an ordered log where every payload is
7//! preserved, use [`stream`](crate::stream) instead.
8//!
9//! On the wire each value is one group holding one frame, so a group is self-contained and a
10//! consumer never needs an older one. With [`Config::compression`] set to
11//! [`Compression::Deflate`], that frame is its own raw DEFLATE stream;
12//! there is no window to share across a single-frame group.
13
14use crate::Compression;
15
16pub mod consumer;
17pub mod producer;
18
19pub use consumer::Consumer;
20pub use producer::Producer;
21
22/// Codec options for a snapshot track.
23///
24/// Build from [`Default`] and override fields (the struct is `#[non_exhaustive]`, so new options
25/// stay additive).
26#[derive(Debug, Clone, Default)]
27#[non_exhaustive]
28pub struct Config {
29	/// Compress each value as its own raw DEFLATE stream.
30	///
31	/// A snapshot group holds a single self-contained value, so there is no window to share: each
32	/// value is compressed alone. [`Compression::None`] (the default) writes the bytes through
33	/// untouched. A [`Consumer`] must set the same [`compression`](Self::compression).
34	pub compression: Compression,
35}
36
37#[cfg(test)]
38mod test {
39	use std::task::Poll;
40
41	use bytes::Bytes;
42
43	use super::*;
44
45	fn cfg(compression: bool) -> Config {
46		Config {
47			compression: if compression {
48				Compression::Deflate
49			} else {
50				Compression::None
51			},
52		}
53	}
54
55	fn producer(compression: bool) -> (Producer, moq_net::track::Subscriber) {
56		let track = moq_net::broadcast::Info::new()
57			.produce()
58			.create_track("test", None)
59			.unwrap();
60		let consumer = track.subscribe(None);
61		(Producer::new(track, cfg(compression)), consumer)
62	}
63
64	fn consume(track: moq_net::track::Subscriber, compression: bool) -> Consumer {
65		Consumer::new(track, cfg(compression))
66	}
67
68	/// Drain every value currently available without blocking.
69	fn drain(mut consumer: Consumer) -> Vec<Bytes> {
70		let waiter = kio::Waiter::noop();
71		let mut out = Vec::new();
72		while let Poll::Ready(Ok(Some(payload))) = consumer.poll_next(&waiter) {
73			out.push(payload);
74		}
75		out
76	}
77
78	/// A snapshot group the transport can no longer serve -- `Old` when the relay reclaims a
79	/// superseded group, `Evicted` under memory pressure, `Lagged` past the drift budget -- is not
80	/// fatal. A snapshot reader only wants the newest value, so it drops the group and takes the
81	/// replacement rather than ending the reader.
82	#[test]
83	fn a_lost_group_waits_for_its_replacement() {
84		let track = moq_net::broadcast::Info::new()
85			.produce()
86			.create_track("test", None)
87			.unwrap();
88		let mut consumer = consume(track.subscribe(None), false);
89		let waiter = kio::Waiter::noop();
90
91		// Group 0 delivers a value, then stays open with the reader parked on its next frame.
92		let mut group = track.create_group(moq_net::group::Info { sequence: 0 }).unwrap();
93		group
94			.write_frame(moq_net::Timestamp::now(), Bytes::from_static(b"one"))
95			.unwrap();
96		assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(v))) if v == "one"));
97		assert!(consumer.poll_next(&waiter).is_pending());
98
99		// The relay reclaims the group out from under the reader.
100		group.abort(moq_net::Error::Old).unwrap();
101		assert!(
102			consumer.poll_next(&waiter).is_pending(),
103			"a lost group must not end the reader"
104		);
105
106		// The replacement arrives and the reader picks up where the value now lives.
107		let mut group = track.create_group(moq_net::group::Info { sequence: 1 }).unwrap();
108		group
109			.write_frame(moq_net::Timestamp::now(), Bytes::from_static(b"two"))
110			.unwrap();
111		assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(v))) if v == "two"));
112	}
113
114	#[test]
115	fn one_group_per_update() {
116		let (mut producer, track) = producer(false);
117		producer.update(&b"first"[..]).unwrap();
118		producer.update(&b"second"[..]).unwrap();
119		producer.finish().unwrap();
120
121		// Two updates => two groups. A consumer that joins after both only sees the latest.
122		assert_eq!(track.latest(), Some(1));
123		assert_eq!(drain(consume(track, false)), vec![Bytes::from_static(b"second")]);
124	}
125
126	#[test]
127	fn live_consumer_sees_each_update() {
128		let (mut producer, track) = producer(false);
129		let mut consumer = consume(track, false);
130		let waiter = kio::Waiter::noop();
131
132		for n in 0..3u8 {
133			producer.update(vec![n]).unwrap();
134			match consumer.poll_next(&waiter) {
135				Poll::Ready(Ok(Some(payload))) => assert_eq!(&payload[..], &[n]),
136				other => panic!("expected value, got {other:?}"),
137			}
138		}
139	}
140
141	#[test]
142	fn compressed_roundtrip() {
143		let (mut producer, track) = producer(true);
144		let payload = Bytes::from(b"the quick brown fox".repeat(64));
145		producer.update(payload.clone()).unwrap();
146		producer.finish().unwrap();
147
148		assert_eq!(drain(consume(track, true)), vec![payload]);
149	}
150
151	/// The compressed frame on the wire is much smaller than the value it carries, so a consumer
152	/// that ignored the flag would read garbage rather than the payload.
153	#[test]
154	fn compression_shrinks_the_frame() {
155		let (mut producer, track) = producer(true);
156		let payload = Bytes::from(b"the quick brown fox".repeat(64));
157		producer.update(payload.clone()).unwrap();
158		producer.finish().unwrap();
159
160		// Read the raw frame, bypassing the consumer's decompression.
161		let waiter = kio::Waiter::noop();
162		let Poll::Ready(Ok(Some(mut group))) = track.ordered().poll_next_group(&waiter) else {
163			panic!("expected a group");
164		};
165		let Poll::Ready(Ok(Some(frame))) = group.poll_read_frame(&waiter) else {
166			panic!("expected a frame");
167		};
168		assert!(
169			frame.payload.len() < payload.len() / 4,
170			"compressed frame {} should be far below the raw {}",
171			frame.payload.len(),
172			payload.len()
173		);
174	}
175
176	/// `append_group` publishes immediately, so rejecting the frame inside `write_frame` would leave
177	/// an empty newest group behind. A snapshot consumer jumps to the newest, so the previous value
178	/// would vanish even though the update reported an error.
179	#[test]
180	fn a_rejected_update_leaves_the_previous_value_readable() {
181		let (mut producer, track) = producer(false);
182		producer.update(&b"keep"[..]).unwrap();
183
184		let oversized = Bytes::from(vec![0u8; moq_net::group::MAX_CACHE_BYTES as usize + 1]);
185		assert!(producer.update(oversized).is_err());
186		producer.finish().unwrap();
187
188		// A reader arriving now still finds the last good value, not an empty superseding group.
189		assert_eq!(drain(consume(track, false)), vec![Bytes::from_static(b"keep")]);
190	}
191
192	/// `finish` closes the underlying track, so a later update fails rather than being silently
193	/// accepted, and that holds for every clone since they share one track.
194	#[test]
195	fn updating_after_finish_fails_on_every_clone() {
196		let (mut producer, _track) = producer(false);
197		let mut clone = producer.clone();
198
199		producer.update(&b"first"[..]).unwrap();
200		producer.finish().unwrap();
201
202		assert!(producer.update(&b"late"[..]).is_err());
203		assert!(clone.update(&b"late"[..]).is_err());
204	}
205
206	#[test]
207	fn a_finished_track_ends_the_consumer() {
208		let (mut producer, track) = producer(false);
209		producer.update(&b"only"[..]).unwrap();
210		producer.finish().unwrap();
211
212		let mut consumer = consume(track, false);
213		let waiter = kio::Waiter::noop();
214		assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(_)))));
215		assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(None))));
216	}
217}