driven/streaming/
chunk_streamer.rs1use crate::Result;
6
7use super::MAX_CHUNK_SIZE;
8
9#[derive(Debug, Clone)]
11pub struct StreamChunk {
12 pub sequence: u32,
14 pub flags: ChunkFlags,
16 pub data: Vec<u8>,
18}
19
20#[derive(Debug, Clone, Copy, Default)]
22pub struct ChunkFlags(u8);
23
24impl ChunkFlags {
25 pub const FIRST: u8 = 1 << 0;
27 pub const LAST: u8 = 1 << 1;
29 pub const COMPRESSED: u8 = 1 << 2;
31 pub const ENCRYPTED: u8 = 1 << 3;
33
34 pub fn new() -> Self {
35 Self(0)
36 }
37
38 pub fn is_first(self) -> bool {
39 self.0 & Self::FIRST != 0
40 }
41
42 pub fn is_last(self) -> bool {
43 self.0 & Self::LAST != 0
44 }
45
46 pub fn is_compressed(self) -> bool {
47 self.0 & Self::COMPRESSED != 0
48 }
49
50 pub fn set(&mut self, flag: u8) {
51 self.0 |= flag;
52 }
53}
54
55impl StreamChunk {
56 pub fn new(sequence: u32, flags: ChunkFlags, data: Vec<u8>) -> Self {
58 Self {
59 sequence,
60 flags,
61 data,
62 }
63 }
64
65 pub fn serialize(&self) -> Vec<u8> {
67 let mut output = Vec::with_capacity(8 + self.data.len());
68 output.extend_from_slice(&self.sequence.to_le_bytes());
69 output.push(self.flags.0);
70 output.extend_from_slice(&[0, 0, 0]); output.extend_from_slice(&self.data);
72 output
73 }
74
75 pub fn deserialize(data: &[u8]) -> Result<Self> {
77 if data.len() < 8 {
78 return Err(crate::DrivenError::InvalidBinary("Chunk too small".into()));
79 }
80
81 let sequence = u32::from_le_bytes([data[0], data[1], data[2], data[3]]);
82 let flags = ChunkFlags(data[4]);
83 let chunk_data = data[8..].to_vec();
84
85 Ok(Self {
86 sequence,
87 flags,
88 data: chunk_data,
89 })
90 }
91}
92
93#[derive(Debug)]
95pub struct ChunkStreamer {
96 chunk_size: usize,
98 sequence: u32,
100}
101
102impl ChunkStreamer {
103 pub fn new(chunk_size: usize) -> Self {
105 Self {
106 chunk_size: chunk_size.min(MAX_CHUNK_SIZE),
107 sequence: 0,
108 }
109 }
110
111 pub fn chunk(&mut self, data: &[u8]) -> Vec<StreamChunk> {
113 let mut chunks = Vec::new();
114 let total_chunks = data.len().div_ceil(self.chunk_size);
115
116 for (i, chunk_data) in data.chunks(self.chunk_size).enumerate() {
117 let mut flags = ChunkFlags::new();
118 if i == 0 {
119 flags.set(ChunkFlags::FIRST);
120 }
121 if i == total_chunks - 1 {
122 flags.set(ChunkFlags::LAST);
123 }
124
125 chunks.push(StreamChunk::new(self.sequence, flags, chunk_data.to_vec()));
126 self.sequence += 1;
127 }
128
129 chunks
130 }
131
132 pub fn reassemble(&self, chunks: &[StreamChunk]) -> Result<Vec<u8>> {
134 let mut result = Vec::new();
135
136 let mut sorted: Vec<_> = chunks.iter().collect();
138 sorted.sort_by_key(|c| c.sequence);
139
140 for chunk in sorted {
141 result.extend_from_slice(&chunk.data);
142 }
143
144 Ok(result)
145 }
146
147 pub fn reset(&mut self) {
149 self.sequence = 0;
150 }
151}
152
153impl Default for ChunkStreamer {
154 fn default() -> Self {
155 Self::new(MAX_CHUNK_SIZE)
156 }
157}
158
159#[derive(Debug)]
161pub struct StreamReceiver {
162 chunks: Vec<StreamChunk>,
164 next_sequence: u32,
166 complete: bool,
168}
169
170impl StreamReceiver {
171 pub fn new() -> Self {
173 Self {
174 chunks: Vec::new(),
175 next_sequence: 0,
176 complete: false,
177 }
178 }
179
180 pub fn receive(&mut self, chunk: StreamChunk) -> bool {
182 if chunk.flags.is_first() && chunk.sequence == 0 {
183 self.chunks.clear();
184 self.next_sequence = 0;
185 }
186
187 self.chunks.push(chunk.clone());
188
189 if chunk.flags.is_last() {
190 self.complete = true;
191 }
192
193 self.complete
194 }
195
196 pub fn is_complete(&self) -> bool {
198 self.complete
199 }
200
201 pub fn data(&self) -> Result<Vec<u8>> {
203 ChunkStreamer::default().reassemble(&self.chunks)
204 }
205
206 pub fn reset(&mut self) {
208 self.chunks.clear();
209 self.next_sequence = 0;
210 self.complete = false;
211 }
212}
213
214impl Default for StreamReceiver {
215 fn default() -> Self {
216 Self::new()
217 }
218}
219
220#[cfg(test)]
221mod tests {
222 use super::*;
223
224 #[test]
225 fn test_chunk_roundtrip() {
226 let mut flags = ChunkFlags::new();
227 flags.set(ChunkFlags::FIRST);
228 flags.set(ChunkFlags::LAST);
229
230 let chunk = StreamChunk::new(42, flags, vec![1, 2, 3, 4]);
231 let bytes = chunk.serialize();
232 let parsed = StreamChunk::deserialize(&bytes).unwrap();
233
234 assert_eq!(parsed.sequence, 42);
235 assert!(parsed.flags.is_first());
236 assert!(parsed.flags.is_last());
237 assert_eq!(parsed.data, vec![1, 2, 3, 4]);
238 }
239
240 #[test]
241 fn test_streamer() {
242 let mut streamer = ChunkStreamer::new(10);
243 let data = b"Hello, World! This is a test of chunking.";
244
245 let chunks = streamer.chunk(data);
246 assert!(chunks.len() > 1);
247 assert!(chunks[0].flags.is_first());
248 assert!(chunks.last().unwrap().flags.is_last());
249
250 let reassembled = streamer.reassemble(&chunks).unwrap();
251 assert_eq!(reassembled, data);
252 }
253
254 #[test]
255 fn test_receiver() {
256 let mut streamer = ChunkStreamer::new(10);
257 let data = b"Test data for streaming";
258
259 let chunks = streamer.chunk(data);
260
261 let mut receiver = StreamReceiver::new();
262 for chunk in chunks {
263 receiver.receive(chunk);
264 }
265
266 assert!(receiver.is_complete());
267 assert_eq!(receiver.data().unwrap(), data);
268 }
269}