moq-binary 0.1.3

Binary publishing over MoQ tracks: latest-value snapshots, or append-log streams.
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
//! Lossless append-log binary publishing over [`moq-net`](moq_net) tracks.
//!
//! An ordered log of opaque payloads, for consumers that care about every one (an event log, a
//! sequence of samples). Nothing is ever superseded: a consumer yields each payload in the order it
//! was appended. For a latest-value document, use [`snapshot`](crate::snapshot) instead.
//!
//! Retention is bounded, which is the limit of "lossless" here. The group's cache is finite, so a
//! write that would outgrow it aborts the group with [`moq_net::Error::GroupTooLarge`] rather than
//! dropping a prefix some readers missed: a partial log presented as a whole one is what this
//! mode exists to prevent. Keep a log inside what the group retains, and split anything unbounded
//! across successive tracks.
//!
//! On the wire the log is a single group that is never rolled, one payload per frame. A payload
//! that cannot be written ends the track rather than opening a second group: a log missing a record
//! is not lossless, and a gap dressed up as a complete log is worse than a visible failure. With
//! [`Config::compression`] set to [`Compression::Deflate`], that one
//! group is one sync-flushed DEFLATE stream, so each
//! payload compresses against the earlier ones and a run of similar payloads shrinks sharply.

use crate::Compression;

pub mod consumer;
pub mod producer;

pub use consumer::Consumer;
pub use producer::Producer;

/// Codec options for a stream track.
///
/// Build from [`Default`] and override fields (the struct is `#[non_exhaustive]`, so new options
/// stay additive).
#[derive(Debug, Clone, Default)]
#[non_exhaustive]
pub struct Config {
	/// Compress the group as one sync-flushed DEFLATE stream, so each payload reuses the earlier
	/// ones as context.
	///
	/// [`Compression::None`] (the default) writes the bytes through untouched. A [`Consumer`] must
	/// set the same [`compression`](Self::compression).
	pub compression: Compression,
}

#[cfg(test)]
mod test {
	use std::task::Poll;

	use bytes::Bytes;

	use super::*;

	fn cfg(compression: bool) -> Config {
		Config {
			compression: if compression {
				Compression::Deflate
			} else {
				Compression::None
			},
		}
	}

	fn producer(compression: bool) -> (Producer, moq_net::track::Subscriber) {
		let track = moq_net::broadcast::Info::new()
			.produce()
			.create_track("test", None)
			.unwrap();
		let consumer = track.subscribe(None);
		(Producer::new(track, cfg(compression)), consumer)
	}

	/// Drain every payload currently available without blocking.
	fn drain(track: moq_net::track::Subscriber, compression: bool) -> Vec<Bytes> {
		let mut consumer = Consumer::new(track, cfg(compression));
		let waiter = kio::Waiter::noop();
		let mut out = Vec::new();
		while let Poll::Ready(Ok(Some(payload))) = consumer.poll_next(&waiter) {
			out.push(payload);
		}
		out
	}

	fn payloads(count: u8) -> Vec<Bytes> {
		(0..count).map(|n| Bytes::from(vec![n; 16])).collect()
	}

	#[test]
	fn every_payload_survives_in_order() {
		let (mut producer, track) = producer(false);
		let expected = payloads(5);
		for payload in &expected {
			producer.append(payload.clone()).unwrap();
		}
		producer.finish().unwrap();

		// One group holds the whole log, unlike snapshot's group-per-value.
		assert_eq!(track.latest(), Some(0));
		assert_eq!(drain(track, false), expected);
	}

	#[test]
	fn compressed_roundtrip_in_order() {
		let (mut producer, track) = producer(true);
		let expected = payloads(20);
		for payload in &expected {
			producer.append(payload.clone()).unwrap();
		}
		producer.finish().unwrap();

		assert_eq!(drain(track, true), expected);
	}

	/// The window spans the whole group, so a repeated payload costs almost nothing after the first.
	#[test]
	fn the_shared_window_shrinks_repetitive_payloads() {
		let (mut producer, track) = producer(true);
		let payload = Bytes::from(b"the quick brown fox".repeat(16));
		for _ in 0..8 {
			producer.append(payload.clone()).unwrap();
		}
		producer.finish().unwrap();

		let waiter = kio::Waiter::noop();
		let Poll::Ready(Ok(Some(mut group))) = track.ordered().poll_next_group(&waiter) else {
			panic!("expected a group");
		};

		let mut sizes = Vec::new();
		while let Poll::Ready(Ok(Some(frame))) = group.poll_read_frame(&waiter) {
			sizes.push(frame.payload.len());
		}

		assert_eq!(sizes.len(), 8);
		assert!(
			*sizes.last().unwrap() < sizes[0],
			"windowed frame {} should be below the first {}",
			sizes.last().unwrap(),
			sizes[0]
		);
	}

	/// Clones share one track and one window, so two owners interleave into a single ordered log.
	#[test]
	fn clones_append_into_one_log() {
		let (mut producer, track) = producer(true);
		let mut clone = producer.clone();

		producer.append(&b"a"[..]).unwrap();
		clone.append(&b"b"[..]).unwrap();
		producer.append(&b"c"[..]).unwrap();
		producer.finish().unwrap();

		assert_eq!(
			drain(track, true),
			vec![
				Bytes::from_static(b"a"),
				Bytes::from_static(b"b"),
				Bytes::from_static(b"c")
			]
		);
	}

	/// A track whose timescale is extreme enough that converting a wall-clock timestamp into it
	/// overflows, so `write_frame` rejects every frame. That stands in for any post-`append_group`
	/// write failure (the real one is a frame over moq-net's 32 MB per-group cache) without
	/// allocating 32 MB to provoke it. Borrowed from moq-json's stream tests.
	fn rejecting_track() -> moq_net::track::Producer {
		let mut info = moq_net::track::Info::default();
		info.timescale = moq_net::Timescale::new((1u64 << 62) - 1).unwrap();

		moq_net::broadcast::Info::new()
			.produce()
			.create_track("test", Some(info))
			.unwrap()
	}

	/// A failed write must reach the consumer, not just the caller. A clean close drains a reader to
	/// `None`, which is exactly what a completed log looks like, so a truncated log would be
	/// indistinguishable from a whole one.
	#[test]
	fn a_failed_write_aborts_the_track() {
		let track = rejecting_track();
		let mut subscriber = track.subscribe(None);
		let mut producer = Producer::new(track, Config::default());

		assert!(producer.append(&b"rejected"[..]).is_err());

		// Aborting the group alone is not enough: the group is dropped from the cache and the reader
		// still sees a clean end, which is exactly what a completed log looks like.
		let waiter = kio::Waiter::noop();
		assert!(
			matches!(subscriber.poll_recv_group(&waiter), Poll::Ready(Err(_))),
			"a truncated log must surface an error rather than read as a completed one"
		);
	}

	/// A record the consumer could never decode is as lost as one the track rejects, so it takes the
	/// track with it rather than leaving a live log missing a record. Guards the guard: an early
	/// return here would bypass the terminal path and let a later append continue the gap.
	#[test]
	fn an_undecodable_record_ends_the_track() {
		let track = moq_net::broadcast::Info::new()
			.produce()
			.create_track("test", None)
			.unwrap();
		let mut subscriber = track.subscribe(None);
		let mut producer = Producer::new(track, cfg(true));

		assert!(producer.is_used());
		let oversized = Bytes::from(vec![0u8; moq_flate::DEFAULT_MAX_FRAME_SIZE as usize + 1]);
		assert!(matches!(
			producer.append(oversized),
			Err(crate::Error::Flate(moq_flate::Error::TooLarge(_)))
		));

		assert!(!producer.is_used());

		// Nothing was published, and the track is terminal rather than merely skipping the record.
		let waiter = kio::Waiter::noop();
		assert!(matches!(subscriber.poll_recv_group(&waiter), Poll::Ready(Err(_))));
		assert!(producer.append(&b"after"[..]).is_err());
	}

	/// A reader already inside the group keeps its own handle, which `track::Producer::abort`
	/// deliberately leaves independent. So the group has to carry the same error, or that reader is
	/// told the producer went away rather than why the log stopped.
	#[test]
	fn a_reader_inside_the_group_sees_the_real_error() {
		let track = moq_net::broadcast::Info::new()
			.produce()
			.create_track("test", None)
			.unwrap();
		let mut subscriber = track.subscribe(None);
		let mut producer = Producer::new(track, Config::default());

		producer.append(&b"first"[..]).unwrap();

		// Pull the group, the way a live reader would, so we hold a handle of our own.
		let waiter = kio::Waiter::noop();
		let Poll::Ready(Ok(Some(mut group))) = subscriber.poll_recv_group(&waiter) else {
			panic!("the group was published, so a subscriber sees it");
		};
		assert!(matches!(group.poll_read_frame(&waiter), Poll::Ready(Ok(Some(_)))));

		// A payload the group cannot hold, so the write fails and ends the log.
		let oversized = Bytes::from(vec![0u8; moq_net::group::MAX_CACHE_BYTES as usize + 1]);
		assert!(producer.append(oversized).is_err());

		match group.poll_read_frame(&waiter) {
			Poll::Ready(Err(err)) => assert!(
				matches!(err, moq_net::Error::FrameTooLarge),
				"the reader should see the write's own error, got {err:?}"
			),
			other => panic!("expected the write error, got {other:?}"),
		}
	}

	/// These tests walk a finished multi-group track, so they need a subscriber that tolerates a
	/// backlog. The default budget is [`Duration::ZERO`], which abandons any group a newer one has
	/// already superseded.
	fn replaying() -> moq_net::track::Subscription {
		moq_net::track::Subscription::default().with_max_age(std::time::Duration::from_secs(30))
	}

	/// The track ends with the group, so nothing opens a second one and splits the log.
	#[test]
	fn a_failed_write_ends_the_track() {
		let track = rejecting_track();
		let mut producer = Producer::new(track, Config::default());

		assert!(producer.append(&b"rejected"[..]).is_err());
		assert!(
			producer.append(&b"again"[..]).is_err(),
			"a second append must fail on the closed track rather than open another group"
		);
	}

	/// A stream is one group. A publisher that rolls to a second one lost whatever would have
	/// completed the first, so the read reports that rather than handing back the remainder as a
	/// continuous log. Written by hand because this producer never rolls.
	#[test]
	fn a_second_group_is_a_rolled_log() {
		let track = moq_net::broadcast::Info::new()
			.produce()
			.create_track("test", None)
			.unwrap();
		let subscriber = track.subscribe(replaying());

		for pair in payloads(4).chunks(2) {
			// Each group is its own DEFLATE stream, which is what a recovery roll would produce.
			let mut flate = moq_flate::Encoder::new();
			let mut group = track.append_group().unwrap();
			for payload in pair {
				group
					.write_frame(moq_net::Timestamp::now(), flate.frame(payload))
					.unwrap();
			}
			group.finish().unwrap();
		}
		track.finish().unwrap();

		let mut consumer = Consumer::new(subscriber, cfg(true));
		let waiter = kio::Waiter::noop();

		// The log's one group reads normally.
		assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(_)))));
		assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(_)))));

		// The second group is a gap, not a continuation.
		assert!(matches!(
			consumer.poll_next(&waiter),
			Poll::Ready(Err(crate::Error::Rolled))
		));
	}

	/// A boundary-only check would never look at the track again while a group is open, so a
	/// publisher that opens a second group and leaves the first one running parks the read forever
	/// on a log that already lost payloads. The payloads already in hand are still delivered.
	#[test]
	fn a_second_group_is_reported_while_the_first_is_open() {
		let track = moq_net::broadcast::Info::new()
			.produce()
			.create_track("test", None)
			.unwrap();
		let subscriber = track.subscribe(replaying());

		// Both groups stay open, the way a publisher writing to two at once leaves them.
		let mut first = track.append_group().unwrap();
		first.write_frame(moq_net::Timestamp::now(), &b"first"[..]).unwrap();
		let mut second = track.append_group().unwrap();
		second.write_frame(moq_net::Timestamp::now(), &b"second"[..]).unwrap();

		let mut consumer = Consumer::new(subscriber, Config::default());
		let waiter = kio::Waiter::noop();

		assert!(matches!(
			consumer.poll_next(&waiter),
			Poll::Ready(Ok(Some(payload))) if payload == Bytes::from_static(b"first")
		));
		assert!(matches!(
			consumer.poll_next(&waiter),
			Poll::Ready(Err(crate::Error::Rolled))
		));

		// Sticky: a later read must not report the rest of the first group as a whole log.
		first.write_frame(moq_net::Timestamp::now(), &b"more"[..]).unwrap();
		assert!(matches!(
			consumer.poll_next(&waiter),
			Poll::Ready(Err(crate::Error::Rolled))
		));
	}

	/// Groups are separate QUIC streams, so a second one can land before the first. Reading in
	/// arrival order is what catches that: the monotonic `next_group` would skip the late lower
	/// sequence and end the log cleanly, reporting a truncated log as a whole one.
	#[test]
	fn a_late_lower_group_is_still_reported() {
		let track = moq_net::broadcast::Info::new()
			.produce()
			.create_track("test", None)
			.unwrap();
		let subscriber = track.subscribe(replaying());

		// Publish sequence 1 before sequence 0, the way reordering delivers them.
		for sequence in [1u64, 0] {
			let mut flate = moq_flate::Encoder::new();
			let mut group = track.create_group(moq_net::group::Info { sequence }).unwrap();
			group
				.write_frame(moq_net::Timestamp::now(), flate.frame(&[sequence as u8; 8]))
				.unwrap();
			group.finish().unwrap();
		}
		track.finish().unwrap();

		let mut consumer = Consumer::new(subscriber, cfg(true));
		let waiter = kio::Waiter::noop();

		assert!(matches!(
			consumer.poll_next(&waiter),
			Poll::Ready(Ok(Some(payload))) if payload == vec![1u8; 8]
		));
		assert!(matches!(
			consumer.poll_next(&waiter),
			Poll::Ready(Err(crate::Error::Rolled))
		));
	}

	/// `finish` closes the underlying track, so a later append fails rather than being silently
	/// accepted, and that holds for every clone since they share one track.
	#[test]
	fn appending_after_finish_fails_on_every_clone() {
		let (mut producer, _track) = producer(false);
		let mut clone = producer.clone();

		producer.append(&b"first"[..]).unwrap();
		producer.finish().unwrap();

		assert!(producer.append(&b"late"[..]).is_err());
		assert!(clone.append(&b"late"[..]).is_err());
	}

	#[test]
	fn a_finished_track_ends_the_consumer() {
		let (mut producer, track) = producer(false);
		producer.append(&b"only"[..]).unwrap();
		producer.finish().unwrap();

		let mut consumer = Consumer::new(track, Config::default());
		let waiter = kio::Waiter::noop();
		assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(Some(_)))));
		assert!(matches!(consumer.poll_next(&waiter), Poll::Ready(Ok(None))));
	}
}