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	/// A stamped payload is written at its capture time, and the returned size is the encoded frame
115	/// on the wire rather than the payload handed in.
116	#[test]
117	fn a_stamped_update_writes_its_capture_time() {
118		let (mut producer, track) = producer(true);
119		let mut groups = producer.consume();
120		let captured = moq_net::Timestamp::from_millis(1_234).unwrap();
121		let payload = Bytes::from(vec![7u8; 4096]);
122		let size = producer
123			.update(moq_net::Timed::from(payload.clone()).at(captured))
124			.unwrap();
125
126		let waiter = kio::Waiter::noop();
127		let Poll::Ready(Ok(Some(mut group))) = groups.poll_recv_group(&waiter) else {
128			panic!("expected a group");
129		};
130		let Poll::Ready(Ok(Some(frame))) = group.poll_read_frame(&waiter) else {
131			panic!("expected a frame");
132		};
133		assert_eq!(frame.timestamp.as_micros(), captured.as_micros());
134		assert_eq!(size, frame.payload.len());
135		assert!(size < payload.len(), "the size is the compressed frame");
136		assert_eq!(drain(consume(track, true)), vec![payload]);
137	}
138
139	#[test]
140	fn one_group_per_update() {
141		let (mut producer, track) = producer(false);
142		producer.update(&b"first"[..]).unwrap();
143		producer.update(&b"second"[..]).unwrap();
144		producer.finish().unwrap();
145
146		// Two updates => two groups. A consumer that joins after both only sees the latest.
147		assert_eq!(track.latest(), Some(1));
148		assert_eq!(drain(consume(track, false)), vec![Bytes::from_static(b"second")]);
149	}
150
151	#[test]
152	fn live_consumer_sees_each_update() {
153		let (mut producer, track) = producer(false);
154		let mut consumer = consume(track, false);
155		let waiter = kio::Waiter::noop();
156
157		for n in 0..3u8 {
158			producer.update(vec![n]).unwrap();
159			match consumer.poll_next(&waiter) {
160				Poll::Ready(Ok(Some(payload))) => assert_eq!(&payload[..], &[n]),
161				other => panic!("expected value, got {other:?}"),
162			}
163		}
164	}
165
166	#[test]
167	fn compressed_roundtrip() {
168		let (mut producer, track) = producer(true);
169		let payload = Bytes::from(b"the quick brown fox".repeat(64));
170		producer.update(payload.clone()).unwrap();
171		producer.finish().unwrap();
172
173		assert_eq!(drain(consume(track, true)), vec![payload]);
174	}
175
176	/// The compressed frame on the wire is much smaller than the value it carries, so a consumer
177	/// that ignored the flag would read garbage rather than the payload.
178	#[test]
179	fn compression_shrinks_the_frame() {
180		let (mut producer, track) = producer(true);
181		let payload = Bytes::from(b"the quick brown fox".repeat(64));
182		producer.update(payload.clone()).unwrap();
183		producer.finish().unwrap();
184
185		// Read the raw frame, bypassing the consumer's decompression.
186		let waiter = kio::Waiter::noop();
187		let Poll::Ready(Ok(Some(mut group))) = track.ordered().poll_next_group(&waiter) else {
188			panic!("expected a group");
189		};
190		let Poll::Ready(Ok(Some(frame))) = group.poll_read_frame(&waiter) else {
191			panic!("expected a frame");
192		};
193		assert!(
194			frame.payload.len() < payload.len() / 4,
195			"compressed frame {} should be far below the raw {}",
196			frame.payload.len(),
197			payload.len()
198		);
199	}
200
201	/// `append_group` publishes immediately, so rejecting the frame inside `write_frame` would leave
202	/// an empty newest group behind. A snapshot consumer jumps to the newest, so the previous value
203	/// would vanish even though the update reported an error.
204	#[test]
205	fn a_rejected_update_leaves_the_previous_value_readable() {
206		let (mut producer, track) = producer(false);
207		producer.update(&b"keep"[..]).unwrap();
208
209		let oversized = Bytes::from(vec![0u8; moq_net::group::MAX_CACHE_BYTES as usize + 1]);
210		assert!(producer.update(oversized).is_err());
211		producer.finish().unwrap();
212
213		// A reader arriving now still finds the last good value, not an empty superseding group.
214		assert_eq!(drain(consume(track, false)), vec![Bytes::from_static(b"keep")]);
215	}
216
217	/// `finish` closes the underlying track, so a later update fails rather than being silently
218	/// accepted, and that holds for every clone since they share one track.
219	#[test]
220	fn updating_after_finish_fails_on_every_clone() {
221		let (mut producer, _track) = producer(false);
222		let mut clone = producer.clone();
223
224		producer.update(&b"first"[..]).unwrap();
225		producer.finish().unwrap();
226
227		assert!(producer.update(&b"late"[..]).is_err());
228		assert!(clone.update(&b"late"[..]).is_err());
229	}
230
231	#[test]
232	fn a_finished_track_ends_the_consumer() {
233		let (mut producer, track) = producer(false);
234		producer.update(&b"only"[..]).unwrap();
235		producer.finish().unwrap();
236
237		let mut consumer = consume(track, false);
238		let waiter = kio::Waiter::noop();
239		assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(_)))));
240		assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(None))));
241	}
242}