Skip to main content

moq_binary/stream/
mod.rs

1//! Lossless append-log binary publishing over [`moq-net`](moq_net) tracks.
2//!
3//! An ordered log of opaque payloads, for consumers that care about every one (an event log, a
4//! sequence of samples). Nothing is ever superseded: a consumer yields each payload in the order it
5//! was appended. For a latest-value document, use [`snapshot`](crate::snapshot) instead.
6//!
7//! Retention is bounded, which is the limit of "lossless" here. The group's cache is finite, so a
8//! write that would outgrow it aborts the group with [`moq_net::Error::GroupTooLarge`] rather than
9//! dropping a prefix some readers missed: a partial log presented as a whole one is what this
10//! mode exists to prevent. Keep a log inside what the group retains, and split anything unbounded
11//! across successive tracks.
12//!
13//! On the wire the log is a single group that is never rolled, one payload per frame. A payload
14//! that cannot be written ends the track rather than opening a second group: a log missing a record
15//! is not lossless, and a gap dressed up as a complete log is worse than a visible failure. With
16//! [`Config::compression`] set to [`Compression::Deflate`], that one
17//! group is one sync-flushed DEFLATE stream, so each
18//! payload compresses against the earlier ones and a run of similar payloads shrinks sharply.
19
20use crate::Compression;
21
22pub mod consumer;
23pub mod producer;
24
25pub use consumer::Consumer;
26pub use producer::Producer;
27
28/// Codec options for a stream track.
29///
30/// Build from [`Default`] and override fields (the struct is `#[non_exhaustive]`, so new options
31/// stay additive).
32#[derive(Debug, Clone, Default)]
33#[non_exhaustive]
34pub struct Config {
35	/// Compress the group as one sync-flushed DEFLATE stream, so each payload reuses the earlier
36	/// ones as context.
37	///
38	/// [`Compression::None`] (the default) writes the bytes through untouched. A [`Consumer`] must
39	/// set the same [`compression`](Self::compression).
40	pub compression: Compression,
41}
42
43#[cfg(test)]
44mod test {
45	use std::task::Poll;
46
47	use bytes::Bytes;
48
49	use super::*;
50
51	fn cfg(compression: bool) -> Config {
52		Config {
53			compression: if compression {
54				Compression::Deflate
55			} else {
56				Compression::None
57			},
58		}
59	}
60
61	fn producer(compression: bool) -> (Producer, moq_net::track::Subscriber) {
62		let track = moq_net::broadcast::Info::new()
63			.produce()
64			.create_track("test", None)
65			.unwrap();
66		let consumer = track.subscribe(None);
67		(Producer::new(track, cfg(compression)), consumer)
68	}
69
70	/// Drain every payload currently available without blocking.
71	fn drain(track: moq_net::track::Subscriber, compression: bool) -> Vec<Bytes> {
72		let mut consumer = Consumer::new(track, cfg(compression));
73		let waiter = kio::Waiter::noop();
74		let mut out = Vec::new();
75		while let Poll::Ready(Ok(Some(payload))) = consumer.poll_next(&waiter) {
76			out.push(payload);
77		}
78		out
79	}
80
81	fn payloads(count: u8) -> Vec<Bytes> {
82		(0..count).map(|n| Bytes::from(vec![n; 16])).collect()
83	}
84
85	#[test]
86	fn every_payload_survives_in_order() {
87		let (mut producer, track) = producer(false);
88		let expected = payloads(5);
89		for payload in &expected {
90			producer.append(payload.clone()).unwrap();
91		}
92		producer.finish().unwrap();
93
94		// One group holds the whole log, unlike snapshot's group-per-value.
95		assert_eq!(track.latest(), Some(0));
96		assert_eq!(drain(track, false), expected);
97	}
98
99	#[test]
100	fn compressed_roundtrip_in_order() {
101		let (mut producer, track) = producer(true);
102		let expected = payloads(20);
103		for payload in &expected {
104			producer.append(payload.clone()).unwrap();
105		}
106		producer.finish().unwrap();
107
108		assert_eq!(drain(track, true), expected);
109	}
110
111	/// The window spans the whole group, so a repeated payload costs almost nothing after the first.
112	#[test]
113	fn the_shared_window_shrinks_repetitive_payloads() {
114		let (mut producer, track) = producer(true);
115		let payload = Bytes::from(b"the quick brown fox".repeat(16));
116		for _ in 0..8 {
117			producer.append(payload.clone()).unwrap();
118		}
119		producer.finish().unwrap();
120
121		let waiter = kio::Waiter::noop();
122		let Poll::Ready(Ok(Some(mut group))) = track.ordered().poll_next_group(&waiter) else {
123			panic!("expected a group");
124		};
125
126		let mut sizes = Vec::new();
127		while let Poll::Ready(Ok(Some(frame))) = group.poll_read_frame(&waiter) {
128			sizes.push(frame.payload.len());
129		}
130
131		assert_eq!(sizes.len(), 8);
132		assert!(
133			*sizes.last().unwrap() < sizes[0],
134			"windowed frame {} should be below the first {}",
135			sizes.last().unwrap(),
136			sizes[0]
137		);
138	}
139
140	/// Clones share one track and one window, so two owners interleave into a single ordered log.
141	#[test]
142	fn clones_append_into_one_log() {
143		let (mut producer, track) = producer(true);
144		let mut clone = producer.clone();
145
146		producer.append(&b"a"[..]).unwrap();
147		clone.append(&b"b"[..]).unwrap();
148		producer.append(&b"c"[..]).unwrap();
149		producer.finish().unwrap();
150
151		assert_eq!(
152			drain(track, true),
153			vec![
154				Bytes::from_static(b"a"),
155				Bytes::from_static(b"b"),
156				Bytes::from_static(b"c")
157			]
158		);
159	}
160
161	/// A track whose timescale is extreme enough that converting a wall-clock timestamp into it
162	/// overflows, so `write_frame` rejects every frame. That stands in for any post-`append_group`
163	/// write failure (the real one is a frame over moq-net's 32 MB per-group cache) without
164	/// allocating 32 MB to provoke it. Borrowed from moq-json's stream tests.
165	fn rejecting_track() -> moq_net::track::Producer {
166		let mut info = moq_net::track::Info::default();
167		info.timescale = moq_net::Timescale::new((1u64 << 62) - 1).unwrap();
168
169		moq_net::broadcast::Info::new()
170			.produce()
171			.create_track("test", Some(info))
172			.unwrap()
173	}
174
175	/// A failed write must reach the consumer, not just the caller. A clean close drains a reader to
176	/// `None`, which is exactly what a completed log looks like, so a truncated log would be
177	/// indistinguishable from a whole one.
178	#[test]
179	fn a_failed_write_aborts_the_track() {
180		let track = rejecting_track();
181		let mut subscriber = track.subscribe(None);
182		let mut producer = Producer::new(track, Config::default());
183
184		assert!(producer.append(&b"rejected"[..]).is_err());
185
186		// Aborting the group alone is not enough: the group is dropped from the cache and the reader
187		// still sees a clean end, which is exactly what a completed log looks like.
188		let waiter = kio::Waiter::noop();
189		assert!(
190			matches!(subscriber.poll_recv_group(&waiter), Poll::Ready(Err(_))),
191			"a truncated log must surface an error rather than read as a completed one"
192		);
193	}
194
195	/// A record the consumer could never decode is as lost as one the track rejects, so it takes the
196	/// track with it rather than leaving a live log missing a record. Guards the guard: an early
197	/// return here would bypass the terminal path and let a later append continue the gap.
198	#[test]
199	fn an_undecodable_record_ends_the_track() {
200		let track = moq_net::broadcast::Info::new()
201			.produce()
202			.create_track("test", None)
203			.unwrap();
204		let mut subscriber = track.subscribe(None);
205		let mut producer = Producer::new(track, cfg(true));
206
207		assert!(producer.is_used());
208		let oversized = Bytes::from(vec![0u8; moq_flate::DEFAULT_MAX_FRAME_SIZE as usize + 1]);
209		assert!(matches!(
210			producer.append(oversized),
211			Err(crate::Error::Flate(moq_flate::Error::TooLarge(_)))
212		));
213
214		assert!(!producer.is_used());
215
216		// Nothing was published, and the track is terminal rather than merely skipping the record.
217		let waiter = kio::Waiter::noop();
218		assert!(matches!(subscriber.poll_recv_group(&waiter), Poll::Ready(Err(_))));
219		assert!(producer.append(&b"after"[..]).is_err());
220	}
221
222	/// A reader already inside the group keeps its own handle, which `track::Producer::abort`
223	/// deliberately leaves independent. So the group has to carry the same error, or that reader is
224	/// told the producer went away rather than why the log stopped.
225	#[test]
226	fn a_reader_inside_the_group_sees_the_real_error() {
227		let track = moq_net::broadcast::Info::new()
228			.produce()
229			.create_track("test", None)
230			.unwrap();
231		let mut subscriber = track.subscribe(None);
232		let mut producer = Producer::new(track, Config::default());
233
234		producer.append(&b"first"[..]).unwrap();
235
236		// Pull the group, the way a live reader would, so we hold a handle of our own.
237		let waiter = kio::Waiter::noop();
238		let Poll::Ready(Ok(Some(mut group))) = subscriber.poll_recv_group(&waiter) else {
239			panic!("the group was published, so a subscriber sees it");
240		};
241		assert!(matches!(group.poll_read_frame(&waiter), Poll::Ready(Ok(Some(_)))));
242
243		// A payload the group cannot hold, so the write fails and ends the log.
244		let oversized = Bytes::from(vec![0u8; moq_net::group::MAX_CACHE_BYTES as usize + 1]);
245		assert!(producer.append(oversized).is_err());
246
247		match group.poll_read_frame(&waiter) {
248			Poll::Ready(Err(err)) => assert!(
249				matches!(err, moq_net::Error::FrameTooLarge),
250				"the reader should see the write's own error, got {err:?}"
251			),
252			other => panic!("expected the write error, got {other:?}"),
253		}
254	}
255
256	/// These tests walk a finished multi-group track, so they need a subscriber that tolerates a
257	/// backlog. The default budget is [`Duration::ZERO`], which abandons any group a newer one has
258	/// already superseded.
259	fn replaying() -> moq_net::track::Subscription {
260		moq_net::track::Subscription::default().with_max_age(std::time::Duration::from_secs(30))
261	}
262
263	/// The track ends with the group, so nothing opens a second one and splits the log.
264	#[test]
265	fn a_failed_write_ends_the_track() {
266		let track = rejecting_track();
267		let mut producer = Producer::new(track, Config::default());
268
269		assert!(producer.append(&b"rejected"[..]).is_err());
270		assert!(
271			producer.append(&b"again"[..]).is_err(),
272			"a second append must fail on the closed track rather than open another group"
273		);
274	}
275
276	/// A stream is one group. A publisher that rolls to a second one lost whatever would have
277	/// completed the first, so the read reports that rather than handing back the remainder as a
278	/// continuous log. Written by hand because this producer never rolls.
279	#[test]
280	fn a_second_group_is_a_rolled_log() {
281		let track = moq_net::broadcast::Info::new()
282			.produce()
283			.create_track("test", None)
284			.unwrap();
285		let subscriber = track.subscribe(replaying());
286
287		for pair in payloads(4).chunks(2) {
288			// Each group is its own DEFLATE stream, which is what a recovery roll would produce.
289			let mut flate = moq_flate::Encoder::new();
290			let mut group = track.append_group().unwrap();
291			for payload in pair {
292				group
293					.write_frame(moq_net::Timestamp::now(), flate.frame(payload))
294					.unwrap();
295			}
296			group.finish().unwrap();
297		}
298		track.finish().unwrap();
299
300		let mut consumer = Consumer::new(subscriber, cfg(true));
301		let waiter = kio::Waiter::noop();
302
303		// The log's one group reads normally.
304		assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(_)))));
305		assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(_)))));
306
307		// The second group is a gap, not a continuation.
308		assert!(matches!(
309			consumer.poll_next(&waiter),
310			Poll::Ready(Err(crate::Error::Rolled))
311		));
312	}
313
314	/// A boundary-only check would never look at the track again while a group is open, so a
315	/// publisher that opens a second group and leaves the first one running parks the read forever
316	/// on a log that already lost payloads. The payloads already in hand are still delivered.
317	#[test]
318	fn a_second_group_is_reported_while_the_first_is_open() {
319		let track = moq_net::broadcast::Info::new()
320			.produce()
321			.create_track("test", None)
322			.unwrap();
323		let subscriber = track.subscribe(replaying());
324
325		// Both groups stay open, the way a publisher writing to two at once leaves them.
326		let mut first = track.append_group().unwrap();
327		first.write_frame(moq_net::Timestamp::now(), &b"first"[..]).unwrap();
328		let mut second = track.append_group().unwrap();
329		second.write_frame(moq_net::Timestamp::now(), &b"second"[..]).unwrap();
330
331		let mut consumer = Consumer::new(subscriber, Config::default());
332		let waiter = kio::Waiter::noop();
333
334		assert!(matches!(
335			consumer.poll_next(&waiter),
336			Poll::Ready(Ok(Some(payload))) if payload == Bytes::from_static(b"first")
337		));
338		assert!(matches!(
339			consumer.poll_next(&waiter),
340			Poll::Ready(Err(crate::Error::Rolled))
341		));
342
343		// Sticky: a later read must not report the rest of the first group as a whole log.
344		first.write_frame(moq_net::Timestamp::now(), &b"more"[..]).unwrap();
345		assert!(matches!(
346			consumer.poll_next(&waiter),
347			Poll::Ready(Err(crate::Error::Rolled))
348		));
349	}
350
351	/// Groups are separate QUIC streams, so a second one can land before the first. Reading in
352	/// arrival order is what catches that: the monotonic `next_group` would skip the late lower
353	/// sequence and end the log cleanly, reporting a truncated log as a whole one.
354	#[test]
355	fn a_late_lower_group_is_still_reported() {
356		let track = moq_net::broadcast::Info::new()
357			.produce()
358			.create_track("test", None)
359			.unwrap();
360		let subscriber = track.subscribe(replaying());
361
362		// Publish sequence 1 before sequence 0, the way reordering delivers them.
363		for sequence in [1u64, 0] {
364			let mut flate = moq_flate::Encoder::new();
365			let mut group = track.create_group(moq_net::group::Info { sequence }).unwrap();
366			group
367				.write_frame(moq_net::Timestamp::now(), flate.frame(&[sequence as u8; 8]))
368				.unwrap();
369			group.finish().unwrap();
370		}
371		track.finish().unwrap();
372
373		let mut consumer = Consumer::new(subscriber, cfg(true));
374		let waiter = kio::Waiter::noop();
375
376		assert!(matches!(
377			consumer.poll_next(&waiter),
378			Poll::Ready(Ok(Some(payload))) if payload == vec![1u8; 8]
379		));
380		assert!(matches!(
381			consumer.poll_next(&waiter),
382			Poll::Ready(Err(crate::Error::Rolled))
383		));
384	}
385
386	/// `finish` closes the underlying track, so a later append fails rather than being silently
387	/// accepted, and that holds for every clone since they share one track.
388	#[test]
389	fn appending_after_finish_fails_on_every_clone() {
390		let (mut producer, _track) = producer(false);
391		let mut clone = producer.clone();
392
393		producer.append(&b"first"[..]).unwrap();
394		producer.finish().unwrap();
395
396		assert!(producer.append(&b"late"[..]).is_err());
397		assert!(clone.append(&b"late"[..]).is_err());
398	}
399
400	#[test]
401	fn a_finished_track_ends_the_consumer() {
402		let (mut producer, track) = producer(false);
403		producer.append(&b"only"[..]).unwrap();
404		producer.finish().unwrap();
405
406		let mut consumer = Consumer::new(track, Config::default());
407		let waiter = kio::Waiter::noop();
408		assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(_)))));
409		assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(None))));
410	}
411}