Skip to main content

driven/streaming/
chunk_streamer.rs

1//! Chunk Streamer
2//!
3//! Streaming delivery of rule updates in chunks.
4
5use crate::Result;
6
7use super::MAX_CHUNK_SIZE;
8
9/// Stream chunk
10#[derive(Debug, Clone)]
11pub struct StreamChunk {
12    /// Chunk sequence number
13    pub sequence: u32,
14    /// Chunk flags
15    pub flags: ChunkFlags,
16    /// Chunk data
17    pub data: Vec<u8>,
18}
19
20/// Chunk flags
21#[derive(Debug, Clone, Copy, Default)]
22pub struct ChunkFlags(u8);
23
24impl ChunkFlags {
25    /// First chunk in stream
26    pub const FIRST: u8 = 1 << 0;
27    /// Last chunk in stream
28    pub const LAST: u8 = 1 << 1;
29    /// Chunk is compressed
30    pub const COMPRESSED: u8 = 1 << 2;
31    /// Chunk is encrypted
32    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    /// Create a new chunk
57    pub fn new(sequence: u32, flags: ChunkFlags, data: Vec<u8>) -> Self {
58        Self {
59            sequence,
60            flags,
61            data,
62        }
63    }
64
65    /// Serialize to bytes
66    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]); // Reserved/padding
71        output.extend_from_slice(&self.data);
72        output
73    }
74
75    /// Deserialize from bytes
76    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/// Chunk streamer for breaking large payloads into chunks
94#[derive(Debug)]
95pub struct ChunkStreamer {
96    /// Maximum chunk size
97    chunk_size: usize,
98    /// Current sequence number
99    sequence: u32,
100}
101
102impl ChunkStreamer {
103    /// Create a new streamer
104    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    /// Split data into chunks
112    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    /// Reassemble chunks into data
133    pub fn reassemble(&self, chunks: &[StreamChunk]) -> Result<Vec<u8>> {
134        let mut result = Vec::new();
135
136        // Sort by sequence
137        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    /// Reset sequence number
148    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/// Streaming receiver for accumulating chunks
160#[derive(Debug)]
161pub struct StreamReceiver {
162    /// Accumulated chunks
163    chunks: Vec<StreamChunk>,
164    /// Expected next sequence
165    next_sequence: u32,
166    /// Complete flag
167    complete: bool,
168}
169
170impl StreamReceiver {
171    /// Create a new receiver
172    pub fn new() -> Self {
173        Self {
174            chunks: Vec::new(),
175            next_sequence: 0,
176            complete: false,
177        }
178    }
179
180    /// Receive a chunk
181    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    /// Check if complete
197    pub fn is_complete(&self) -> bool {
198        self.complete
199    }
200
201    /// Get accumulated data
202    pub fn data(&self) -> Result<Vec<u8>> {
203        ChunkStreamer::default().reassemble(&self.chunks)
204    }
205
206    /// Reset receiver
207    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}