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	/// Each record keeps its own capture time, and a bare payload is stamped when written.
100	#[test]
101	fn a_stamped_append_writes_its_capture_time() {
102		let (mut producer, _track) = producer(true);
103		let mut groups = producer.consume();
104		let captured = moq_net::Timestamp::from_millis(1_234).unwrap();
105		let first = producer
106			.append(moq_net::Timed::from(vec![1u8; 4096]).at(captured))
107			.unwrap();
108		producer.append(vec![2u8; 16]).unwrap();
109
110		let waiter = kio::Waiter::noop();
111		let Poll::Ready(Ok(Some(mut group))) = groups.poll_recv_group(&waiter) else {
112			panic!("expected a group");
113		};
114		let Poll::Ready(Ok(Some(frame))) = group.poll_read_frame(&waiter) else {
115			panic!("expected a frame");
116		};
117		assert_eq!(frame.timestamp.as_micros(), captured.as_micros());
118		assert_eq!(first, frame.payload.len(), "the size is the compressed frame");
119		let Poll::Ready(Ok(Some(frame))) = group.poll_read_frame(&waiter) else {
120			panic!("expected a frame");
121		};
122		assert_ne!(frame.timestamp.as_micros(), captured.as_micros());
123	}
124
125	#[test]
126	fn compressed_roundtrip_in_order() {
127		let (mut producer, track) = producer(true);
128		let expected = payloads(20);
129		for payload in &expected {
130			producer.append(payload.clone()).unwrap();
131		}
132		producer.finish().unwrap();
133
134		assert_eq!(drain(track, true), expected);
135	}
136
137	/// The window spans the whole group, so a repeated payload costs almost nothing after the first.
138	#[test]
139	fn the_shared_window_shrinks_repetitive_payloads() {
140		let (mut producer, track) = producer(true);
141		let payload = Bytes::from(b"the quick brown fox".repeat(16));
142		for _ in 0..8 {
143			producer.append(payload.clone()).unwrap();
144		}
145		producer.finish().unwrap();
146
147		let waiter = kio::Waiter::noop();
148		let Poll::Ready(Ok(Some(mut group))) = track.ordered().poll_next_group(&waiter) else {
149			panic!("expected a group");
150		};
151
152		let mut sizes = Vec::new();
153		while let Poll::Ready(Ok(Some(frame))) = group.poll_read_frame(&waiter) {
154			sizes.push(frame.payload.len());
155		}
156
157		assert_eq!(sizes.len(), 8);
158		assert!(
159			*sizes.last().unwrap() < sizes[0],
160			"windowed frame {} should be below the first {}",
161			sizes.last().unwrap(),
162			sizes[0]
163		);
164	}
165
166	/// Clones share one track and one window, so two owners interleave into a single ordered log.
167	#[test]
168	fn clones_append_into_one_log() {
169		let (mut producer, track) = producer(true);
170		let mut clone = producer.clone();
171
172		producer.append(&b"a"[..]).unwrap();
173		clone.append(&b"b"[..]).unwrap();
174		producer.append(&b"c"[..]).unwrap();
175		producer.finish().unwrap();
176
177		assert_eq!(
178			drain(track, true),
179			vec![
180				Bytes::from_static(b"a"),
181				Bytes::from_static(b"b"),
182				Bytes::from_static(b"c")
183			]
184		);
185	}
186
187	/// A track whose timescale is extreme enough that converting a wall-clock timestamp into it
188	/// overflows, so `write_frame` rejects every frame. That stands in for any post-`append_group`
189	/// write failure (the real one is a frame over moq-net's 32 MB per-group cache) without
190	/// allocating 32 MB to provoke it. Borrowed from moq-json's stream tests.
191	fn rejecting_track() -> moq_net::track::Producer {
192		let mut info = moq_net::track::Info::default();
193		info.timescale = moq_net::Timescale::new((1u64 << 62) - 1).unwrap();
194
195		moq_net::broadcast::Info::new()
196			.produce()
197			.create_track("test", Some(info))
198			.unwrap()
199	}
200
201	/// A failed write must reach the consumer, not just the caller. A clean close drains a reader to
202	/// `None`, which is exactly what a completed log looks like, so a truncated log would be
203	/// indistinguishable from a whole one.
204	#[test]
205	fn a_failed_write_aborts_the_track() {
206		let track = rejecting_track();
207		let mut subscriber = track.subscribe(None);
208		let mut producer = Producer::new(track, Config::default());
209
210		assert!(producer.append(&b"rejected"[..]).is_err());
211
212		// Aborting the group alone is not enough: the group is dropped from the cache and the reader
213		// still sees a clean end, which is exactly what a completed log looks like.
214		let waiter = kio::Waiter::noop();
215		assert!(
216			matches!(subscriber.poll_recv_group(&waiter), Poll::Ready(Err(_))),
217			"a truncated log must surface an error rather than read as a completed one"
218		);
219	}
220
221	/// A record the consumer could never decode is as lost as one the track rejects, so it takes the
222	/// track with it rather than leaving a live log missing a record. Guards the guard: an early
223	/// return here would bypass the terminal path and let a later append continue the gap.
224	#[test]
225	fn an_undecodable_record_ends_the_track() {
226		let track = moq_net::broadcast::Info::new()
227			.produce()
228			.create_track("test", None)
229			.unwrap();
230		let mut subscriber = track.subscribe(None);
231		let mut producer = Producer::new(track, cfg(true));
232
233		assert!(producer.is_used());
234		let oversized = Bytes::from(vec![0u8; moq_flate::DEFAULT_MAX_FRAME_SIZE as usize + 1]);
235		assert!(matches!(
236			producer.append(oversized),
237			Err(crate::Error::Flate(moq_flate::Error::TooLarge(_)))
238		));
239
240		assert!(!producer.is_used());
241
242		// Nothing was published, and the track is terminal rather than merely skipping the record.
243		let waiter = kio::Waiter::noop();
244		assert!(matches!(subscriber.poll_recv_group(&waiter), Poll::Ready(Err(_))));
245		assert!(producer.append(&b"after"[..]).is_err());
246	}
247
248	/// A reader already inside the group keeps its own handle, which `track::Producer::abort`
249	/// deliberately leaves independent. So the group has to carry the same error, or that reader is
250	/// told the producer went away rather than why the log stopped.
251	#[test]
252	fn a_reader_inside_the_group_sees_the_real_error() {
253		let track = moq_net::broadcast::Info::new()
254			.produce()
255			.create_track("test", None)
256			.unwrap();
257		let mut subscriber = track.subscribe(None);
258		let mut producer = Producer::new(track, Config::default());
259
260		producer.append(&b"first"[..]).unwrap();
261
262		// Pull the group, the way a live reader would, so we hold a handle of our own.
263		let waiter = kio::Waiter::noop();
264		let Poll::Ready(Ok(Some(mut group))) = subscriber.poll_recv_group(&waiter) else {
265			panic!("the group was published, so a subscriber sees it");
266		};
267		assert!(matches!(group.poll_read_frame(&waiter), Poll::Ready(Ok(Some(_)))));
268
269		// A payload the group cannot hold, so the write fails and ends the log.
270		let oversized = Bytes::from(vec![0u8; moq_net::group::MAX_CACHE_BYTES as usize + 1]);
271		assert!(producer.append(oversized).is_err());
272
273		match group.poll_read_frame(&waiter) {
274			Poll::Ready(Err(err)) => assert!(
275				matches!(err, moq_net::Error::FrameTooLarge),
276				"the reader should see the write's own error, got {err:?}"
277			),
278			other => panic!("expected the write error, got {other:?}"),
279		}
280	}
281
282	/// These tests walk a finished multi-group track, so they need a subscriber that tolerates a
283	/// backlog. The default budget is [`Duration::ZERO`], which abandons any group a newer one has
284	/// already superseded.
285	fn replaying() -> moq_net::track::Subscription {
286		moq_net::track::Subscription::default().with_max_age(std::time::Duration::from_secs(30))
287	}
288
289	/// The track ends with the group, so nothing opens a second one and splits the log.
290	#[test]
291	fn a_failed_write_ends_the_track() {
292		let track = rejecting_track();
293		let mut producer = Producer::new(track, Config::default());
294
295		assert!(producer.append(&b"rejected"[..]).is_err());
296		assert!(
297			producer.append(&b"again"[..]).is_err(),
298			"a second append must fail on the closed track rather than open another group"
299		);
300	}
301
302	/// A stream is one group. A publisher that rolls to a second one lost whatever would have
303	/// completed the first, so the read reports that rather than handing back the remainder as a
304	/// continuous log. Written by hand because this producer never rolls.
305	#[test]
306	fn a_second_group_is_a_rolled_log() {
307		let track = moq_net::broadcast::Info::new()
308			.produce()
309			.create_track("test", None)
310			.unwrap();
311		let subscriber = track.subscribe(replaying());
312
313		for pair in payloads(4).chunks(2) {
314			// Each group is its own DEFLATE stream, which is what a recovery roll would produce.
315			let mut flate = moq_flate::Encoder::new();
316			let mut group = track.append_group().unwrap();
317			for payload in pair {
318				group
319					.write_frame(moq_net::Timestamp::now(), flate.frame(payload))
320					.unwrap();
321			}
322			group.finish().unwrap();
323		}
324		track.finish().unwrap();
325
326		let mut consumer = Consumer::new(subscriber, cfg(true));
327		let waiter = kio::Waiter::noop();
328
329		// The log's one group reads normally.
330		assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(_)))));
331		assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(_)))));
332
333		// The second group is a gap, not a continuation.
334		assert!(matches!(
335			consumer.poll_next(&waiter),
336			Poll::Ready(Err(crate::Error::Rolled))
337		));
338	}
339
340	/// A boundary-only check would never look at the track again while a group is open, so a
341	/// publisher that opens a second group and leaves the first one running parks the read forever
342	/// on a log that already lost payloads. The payloads already in hand are still delivered.
343	#[test]
344	fn a_second_group_is_reported_while_the_first_is_open() {
345		let track = moq_net::broadcast::Info::new()
346			.produce()
347			.create_track("test", None)
348			.unwrap();
349		let subscriber = track.subscribe(replaying());
350
351		// Both groups stay open, the way a publisher writing to two at once leaves them.
352		let mut first = track.append_group().unwrap();
353		first.write_frame(moq_net::Timestamp::now(), &b"first"[..]).unwrap();
354		let mut second = track.append_group().unwrap();
355		second.write_frame(moq_net::Timestamp::now(), &b"second"[..]).unwrap();
356
357		let mut consumer = Consumer::new(subscriber, Config::default());
358		let waiter = kio::Waiter::noop();
359
360		assert!(matches!(
361			consumer.poll_next(&waiter),
362			Poll::Ready(Ok(Some(payload))) if payload == Bytes::from_static(b"first")
363		));
364		assert!(matches!(
365			consumer.poll_next(&waiter),
366			Poll::Ready(Err(crate::Error::Rolled))
367		));
368
369		// Sticky: a later read must not report the rest of the first group as a whole log.
370		first.write_frame(moq_net::Timestamp::now(), &b"more"[..]).unwrap();
371		assert!(matches!(
372			consumer.poll_next(&waiter),
373			Poll::Ready(Err(crate::Error::Rolled))
374		));
375	}
376
377	/// Groups are separate QUIC streams, so a second one can land before the first. Reading in
378	/// arrival order is what catches that: the monotonic `next_group` would skip the late lower
379	/// sequence and end the log cleanly, reporting a truncated log as a whole one.
380	#[test]
381	fn a_late_lower_group_is_still_reported() {
382		let track = moq_net::broadcast::Info::new()
383			.produce()
384			.create_track("test", None)
385			.unwrap();
386		let subscriber = track.subscribe(replaying());
387
388		// Publish sequence 1 before sequence 0, the way reordering delivers them.
389		for sequence in [1u64, 0] {
390			let mut flate = moq_flate::Encoder::new();
391			let mut group = track.create_group(moq_net::group::Info { sequence }).unwrap();
392			group
393				.write_frame(moq_net::Timestamp::now(), flate.frame(&[sequence as u8; 8]))
394				.unwrap();
395			group.finish().unwrap();
396		}
397		track.finish().unwrap();
398
399		let mut consumer = Consumer::new(subscriber, cfg(true));
400		let waiter = kio::Waiter::noop();
401
402		assert!(matches!(
403			consumer.poll_next(&waiter),
404			Poll::Ready(Ok(Some(payload))) if payload == vec![1u8; 8]
405		));
406		assert!(matches!(
407			consumer.poll_next(&waiter),
408			Poll::Ready(Err(crate::Error::Rolled))
409		));
410	}
411
412	/// `finish` closes the underlying track, so a later append fails rather than being silently
413	/// accepted, and that holds for every clone since they share one track.
414	#[test]
415	fn appending_after_finish_fails_on_every_clone() {
416		let (mut producer, _track) = producer(false);
417		let mut clone = producer.clone();
418
419		producer.append(&b"first"[..]).unwrap();
420		producer.finish().unwrap();
421
422		assert!(producer.append(&b"late"[..]).is_err());
423		assert!(clone.append(&b"late"[..]).is_err());
424	}
425
426	#[test]
427	fn a_finished_track_ends_the_consumer() {
428		let (mut producer, track) = producer(false);
429		producer.append(&b"only"[..]).unwrap();
430		producer.finish().unwrap();
431
432		let mut consumer = Consumer::new(track, Config::default());
433		let waiter = kio::Waiter::noop();
434		assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(_)))));
435		assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(None))));
436	}
437}