Skip to main content

moq_json/stream/
decoder.rs

1//! The track-free half of stream consuming: frame payloads in, records out.
2
3use std::marker::PhantomData;
4
5use serde::de::DeserializeOwned;
6
7use crate::Result;
8
9/// Configuration for a [`Decoder`], and so for the [`Consumer`](super::Consumer) wrapping one.
10///
11/// Build from [`Default`] and override fields (the struct is `#[non_exhaustive]`, so new options
12/// stay additive).
13#[derive(Debug, Clone, Default)]
14#[non_exhaustive]
15pub struct ConsumerConfig {
16	/// Whether the frames are DEFLATE-compressed. Must match the encoder's
17	/// [`ProducerConfig::compression`](super::ProducerConfig::compression). Defaults to `false`.
18	pub compression: bool,
19}
20
21impl ConsumerConfig {
22	/// Set [`compression`](Self::compression) (a builder, since the struct is `#[non_exhaustive]`).
23	pub fn with_compression(mut self, compression: bool) -> Self {
24		self.compression = compression;
25		self
26	}
27}
28
29/// Decodes JSON records from frame payloads, sharing one DEFLATE window across the log.
30///
31/// The track-free core of [`Consumer`](super::Consumer), and the mirror of
32/// [`Encoder`](super::Encoder). Payloads must be fed in the order they were encoded, since each one
33/// builds on the window the earlier ones left behind. Call [`reset`](Self::reset) at a group
34/// boundary, matching the encoder.
35pub struct Decoder<T> {
36	/// The DEFLATE decoder (one window for the whole log), `Some` while decompressing.
37	flate: Option<moq_flate::Decoder>,
38	compression: bool,
39	_marker: PhantomData<fn() -> T>,
40}
41
42impl<T> Decoder<T> {
43	/// Create a decoder with a cold window.
44	pub fn new(config: ConsumerConfig) -> Self {
45		Self {
46			flate: config.compression.then(moq_flate::Decoder::new),
47			compression: config.compression,
48			_marker: PhantomData,
49		}
50	}
51
52	/// Start a cold DEFLATE window, for a caller that has just moved to a new group.
53	pub fn reset(&mut self) {
54		self.flate = self.compression.then(moq_flate::Decoder::new);
55	}
56}
57
58impl<T: DeserializeOwned> Decoder<T> {
59	/// Decode the next frame payload back into a record.
60	pub fn decode(&mut self, payload: &[u8]) -> Result<T> {
61		Ok(match self.flate.as_mut() {
62			Some(flate) => serde_json::from_slice(&flate.frame(payload)?)?,
63			None => serde_json::from_slice(payload)?,
64		})
65	}
66}
67
68#[cfg(test)]
69mod test {
70	use super::super::{Encoder, ProducerConfig};
71	use super::*;
72	use serde_json::{Value, json};
73
74	/// Round-trip a sequence of records through an encoder and decoder.
75	fn roundtrip(compression: bool, values: &[Value]) -> Vec<Value> {
76		let mut encoder = Encoder::<Value>::new(ProducerConfig::default().with_compression(compression));
77		let mut decoder = Decoder::<Value>::new(ConsumerConfig::default().with_compression(compression));
78
79		values
80			.iter()
81			.map(|value| {
82				let record = encoder.encode(value).unwrap();
83				let decoded = decoder.decode(record.payload()).unwrap();
84				record.commit();
85				decoded
86			})
87			.collect()
88	}
89
90	#[test]
91	fn plaintext_roundtrip_in_order() {
92		let values: Vec<Value> = (0..5).map(|n| json!({ "n": n })).collect();
93		assert_eq!(roundtrip(false, &values), values);
94	}
95
96	#[test]
97	fn compressed_roundtrip_in_order() {
98		let values: Vec<Value> = (0..20).map(|n| json!({ "group": n, "pts": n * 2_000 })).collect();
99		assert_eq!(roundtrip(true, &values), values);
100	}
101
102	#[test]
103	fn the_shared_window_shrinks_repetitive_records() {
104		let mut encoder = Encoder::<Value>::new(ProducerConfig::default().with_compression(true));
105		let sizes: Vec<usize> = (0..8)
106			.map(|n| {
107				let record = encoder.encode(&json!({ "group": n, "pts": n * 2_000 })).unwrap();
108				let len = record.payload().len();
109				record.commit();
110				len
111			})
112			.collect();
113
114		let raw = serde_json::to_vec(&json!({ "group": 7, "pts": 14_000 })).unwrap().len();
115		assert!(
116			*sizes.last().unwrap() < raw / 2,
117			"windowed record {} should be far below its raw size {raw}",
118			sizes.last().unwrap()
119		);
120	}
121
122	/// A caller that rolls a group has to restart both windows, or the new group's frames decode
123	/// against context the decoder on the other side never received.
124	#[test]
125	fn reset_starts_a_cold_window_on_both_sides() {
126		let mut encoder = Encoder::<Value>::new(ProducerConfig::default().with_compression(true));
127		let mut decoder = Decoder::<Value>::new(ConsumerConfig::default().with_compression(true));
128
129		for n in 0..4 {
130			let record = encoder.encode(&json!({ "n": n })).unwrap();
131			assert_eq!(decoder.decode(record.payload()).unwrap(), json!({ "n": n }));
132			record.commit();
133		}
134
135		encoder.reset();
136		decoder.reset();
137
138		let record = encoder.encode(&json!({ "n": 99 })).unwrap();
139		assert_eq!(decoder.decode(record.payload()).unwrap(), json!({ "n": 99 }));
140		record.commit();
141	}
142}
143
144#[cfg(test)]
145mod desync_test {
146	use super::super::{Encoder, ProducerConfig};
147	use super::*;
148	use serde_json::{Value, json};
149
150	/// A compressed record that never reached the wire leaves the window ahead of the consumer, and a
151	/// log has no keyframe to resynchronize on. Continuing would emit frames nothing can decode, so
152	/// the encoder refuses until the caller rolls a new group and resets.
153	#[test]
154	fn an_uncommitted_compressed_record_stops_the_encoder() {
155		let mut encoder = Encoder::<Value>::new(ProducerConfig::default().with_compression(true));
156		encoder.encode(&json!({ "n": 0 })).unwrap().commit();
157
158		// This one fails to write, so the caller never commits it.
159		drop(encoder.encode(&json!({ "n": 1 })).unwrap());
160
161		assert!(matches!(encoder.encode(&json!({ "n": 2 })), Err(crate::Error::Desync)));
162
163		// Rolling a new group gives the consumer a cold window too, so the reset clears it.
164		encoder.reset();
165		let record = encoder.encode(&json!({ "n": 2 })).unwrap();
166		let mut decoder = Decoder::<Value>::new(ConsumerConfig::default().with_compression(true));
167		assert_eq!(decoder.decode(record.payload()).unwrap(), json!({ "n": 2 }));
168		record.commit();
169	}
170
171	/// Without compression a record carries no shared state, so a dropped one leaves a gap in the log
172	/// rather than an undecodable stream, and the encoder keeps going.
173	#[test]
174	fn an_uncommitted_plaintext_record_does_not_stop_the_encoder() {
175		let mut encoder = Encoder::<Value>::new(ProducerConfig::default());
176		drop(encoder.encode(&json!({ "n": 0 })).unwrap());
177
178		let record = encoder
179			.encode(&json!({ "n": 1 }))
180			.expect("plaintext records are independent");
181		let mut decoder = Decoder::<Value>::new(ConsumerConfig::default());
182		assert_eq!(decoder.decode(record.payload()).unwrap(), json!({ "n": 1 }));
183		record.commit();
184	}
185}