moq_json/stream/
decoder.rs1use std::marker::PhantomData;
4
5use serde::de::DeserializeOwned;
6
7use crate::Result;
8
9#[derive(Debug, Clone, Default)]
14#[non_exhaustive]
15pub struct ConsumerConfig {
16 pub compression: bool,
19}
20
21impl ConsumerConfig {
22 pub fn with_compression(mut self, compression: bool) -> Self {
24 self.compression = compression;
25 self
26 }
27}
28
29pub struct Decoder<T> {
36 flate: Option<moq_flate::Decoder>,
38 compression: bool,
39 _marker: PhantomData<fn() -> T>,
40}
41
42impl<T> Decoder<T> {
43 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 pub fn reset(&mut self) {
54 self.flate = self.compression.then(moq_flate::Decoder::new);
55 }
56}
57
58impl<T: DeserializeOwned> Decoder<T> {
59 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 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 #[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 #[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 drop(encoder.encode(&json!({ "n": 1 })).unwrap());
160
161 assert!(matches!(encoder.encode(&json!({ "n": 2 })), Err(crate::Error::Desync)));
162
163 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 #[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}