subetha_cxc/
interleave.rs1#[derive(Debug)]
27pub struct Interleaver {
28 depth: usize,
29 pending: Vec<Vec<Vec<u8>>>,
31}
32
33impl Interleaver {
34 pub fn new(depth: usize) -> Self {
36 Self {
37 depth: depth.max(1),
38 pending: Vec::new(),
39 }
40 }
41
42 pub fn depth(&self) -> usize {
44 self.depth
45 }
46
47 pub fn set_depth(&mut self, depth: usize) -> Vec<Vec<u8>> {
51 let flushed = self.flush();
52 self.depth = depth.max(1);
53 flushed
54 }
55
56 pub fn buffered_blocks(&self) -> usize {
58 self.pending.len()
59 }
60
61 pub fn add_block(&mut self, datagrams: Vec<Vec<u8>>) -> Vec<Vec<u8>> {
66 if datagrams.is_empty() {
67 return Vec::new();
68 }
69 self.pending.push(datagrams);
70 if self.pending.len() >= self.depth {
71 self.emit()
72 } else {
73 Vec::new()
74 }
75 }
76
77 pub fn flush(&mut self) -> Vec<Vec<u8>> {
80 if self.pending.is_empty() {
81 Vec::new()
82 } else {
83 self.emit()
84 }
85 }
86
87 fn emit(&mut self) -> Vec<Vec<u8>> {
89 let mut blocks = std::mem::take(&mut self.pending);
90 let max_len = blocks.iter().map(|b| b.len()).max().unwrap_or(0);
91 let total: usize = blocks.iter().map(|b| b.len()).sum();
92 let mut out = Vec::with_capacity(total);
93 for col in 0..max_len {
100 for block in blocks.iter_mut() {
101 if col < block.len() {
102 out.push(std::mem::take(&mut block[col]));
103 }
104 }
105 }
106 out
107 }
108}
109
110#[cfg(test)]
111mod tests {
112 use super::*;
113
114 fn make_blocks(n: usize, shards: usize) -> Vec<Vec<Vec<u8>>> {
118 (0..n)
119 .map(|b| (0..shards).map(|s| vec![b as u8, s as u8]).collect())
120 .collect()
121 }
122
123 #[test]
124 fn depth_one_is_passthrough() {
125 let mut il = Interleaver::new(1);
126 let blocks = make_blocks(1, 5);
127 let out = il.add_block(blocks[0].clone());
128 assert_eq!(out, blocks[0], "depth 1 emits the block immediately");
129 assert_eq!(il.buffered_blocks(), 0);
130 }
131
132 #[test]
133 fn emits_when_depth_reached_and_is_a_permutation() {
134 let depth = 4;
135 let shards = 6;
136 let blocks = make_blocks(depth, shards);
137 let mut il = Interleaver::new(depth);
138 let mut out = Vec::new();
139 for (i, b) in blocks.iter().enumerate() {
140 let emitted = il.add_block(b.clone());
141 if i < depth - 1 {
142 assert!(emitted.is_empty(), "no emit before depth reached");
143 } else {
144 out = emitted;
145 }
146 }
147 let mut got = out.clone();
149 let mut want: Vec<Vec<u8>> = blocks.into_iter().flatten().collect();
150 got.sort();
151 want.sort();
152 assert_eq!(got, want, "interleave is a permutation of the input");
153 }
154
155 fn burst_property(depth: usize, shards: usize) {
158 let blocks = make_blocks(depth, shards);
159 let mut il = Interleaver::new(depth);
160 let mut out = Vec::new();
161 for b in &blocks {
162 out.extend(il.add_block(b.clone()));
163 }
164 out.extend(il.flush());
165 assert_eq!(out.len(), depth * shards);
166 for start in 0..=out.len() - depth {
168 let mut per_block = vec![0u32; depth];
169 for pkt in &out[start..start + depth] {
170 per_block[pkt[0] as usize] += 1;
171 }
172 assert!(
173 per_block.iter().all(|&c| c <= 1),
174 "depth={depth} shards={shards} window@{start}: a block lost >1 shard to a burst of {depth}"
175 );
176 }
177 }
178
179 #[test]
180 fn burst_of_depth_hits_at_most_one_shard_per_block() {
181 burst_property(4, 6);
182 burst_property(8, 10);
183 burst_property(3, 3);
184 burst_property(16, 8);
185 }
186
187 #[test]
188 fn set_depth_flushes_pending() {
189 let mut il = Interleaver::new(4);
190 let blocks = make_blocks(2, 5); il.add_block(blocks[0].clone());
192 let flushed = il.add_block(blocks[1].clone());
193 assert!(flushed.is_empty(), "2 of 4 staged, nothing emitted yet");
194 let out = il.set_depth(2);
195 assert_eq!(out.len(), 10, "changing depth flushes the 2 staged blocks");
196 assert_eq!(il.depth(), 2);
197 assert_eq!(il.buffered_blocks(), 0);
198 }
199}
200
201#[cfg(test)]
206mod gilbert_elliott {
207 use super::Interleaver;
208 use crate::reliable_udp::{Decoder, Encoder};
209
210 struct Ge {
212 bad: bool,
213 rng: u64,
214 p_gb: u32, p_bg: u32, p_b: u32, p_g: u32, }
219
220 impl Ge {
221 fn new(seed: u64) -> Self {
222 Self { bad: false, rng: seed | 1, p_gb: 30, p_bg: 200, p_b: 900, p_g: 0 }
223 }
224 fn rand(&mut self) -> u32 {
225 self.rng = self
226 .rng
227 .wrapping_mul(6364136223846793005)
228 .wrapping_add(1442695040888963407);
229 (self.rng >> 33) as u32
230 }
231 fn drop(&mut self) -> bool {
233 if self.bad {
234 if self.rand() % 1000 < self.p_bg {
235 self.bad = false;
236 }
237 } else if self.rand() % 1000 < self.p_gb {
238 self.bad = true;
239 }
240 let p = if self.bad { self.p_b } else { self.p_g };
241 self.rand() % 1000 < p
242 }
243 }
244
245 fn run(depth: usize, n: u64, seed: u64) -> (bool, usize) {
249 let (k, r) = (8usize, 2usize);
250 let mut enc = Encoder::new(k, r, 8);
251 let mut il = Interleaver::new(depth);
252 let mut dec = Decoder::new();
253
254 let mut wire: Vec<Vec<u8>> = Vec::new();
257 for i in 0..n {
258 let block = enc.push(&i.to_le_bytes());
259 if !block.is_empty() {
260 wire.extend(il.add_block(block));
261 }
262 }
263 let tail = enc.flush();
264 if !tail.is_empty() {
265 wire.extend(il.add_block(tail));
266 }
267 wire.extend(il.flush());
268
269 let mut ge = Ge::new(seed);
270 let mut delivered: Vec<u64> = Vec::new();
271 for pkt in &wire {
272 if ge.drop() {
273 continue;
274 }
275 for it in dec.on_packet(pkt) {
276 delivered.push(u64::from_le_bytes(it.try_into().unwrap()));
277 }
278 }
279
280 let mut retransmits = 0usize;
282 let mut rounds = 0u32;
283 while (delivered.len() as u64) < n {
284 rounds += 1;
285 assert!(rounds < 20_000, "no convergence at depth {depth}");
286 let fb = dec.feedback(true);
287 for pkt in enc.on_feedback(&fb) {
288 retransmits += 1;
289 if ge.drop() {
290 continue;
291 }
292 for it in dec.on_packet(&pkt) {
293 delivered.push(u64::from_le_bytes(it.try_into().unwrap()));
294 }
295 }
296 }
297 let ok = delivered == (0..n).collect::<Vec<_>>();
298 (ok, retransmits)
299 }
300
301 #[test]
302 fn interleaving_cuts_arq_under_bursty_loss() {
303 let n = 240;
304 let seed = 0x00C0_FFEE_1234_5678;
305 let (ok1, rtx1) = run(1, n, seed);
306 let (ok8, rtx8) = run(8, n, seed);
307 assert!(ok1 && ok8, "both deliver exactly via the ARQ floor");
308 assert!(
309 rtx8 < rtx1,
310 "interleaving must cut ARQ under bursts: depth8={rtx8} retransmits vs depth1={rtx1}"
311 );
312 }
313}