Skip to main content

oxigdal_streaming/io/
buffer.rs

1//! Buffer management for chunked I/O operations.
2
3use crate::error::{Result, StreamingError};
4use bytes::Bytes;
5use std::collections::VecDeque;
6use std::sync::Arc;
7use tokio::sync::RwLock;
8use tracing::debug;
9
10/// Descriptor for a data chunk.
11#[derive(Debug, Clone)]
12pub struct ChunkDescriptor {
13    /// Offset in bytes from the start of the data
14    pub offset: u64,
15
16    /// Length of the chunk in bytes
17    pub length: usize,
18
19    /// Chunk index
20    pub index: usize,
21
22    /// Total number of chunks
23    pub total_chunks: usize,
24
25    /// Whether this is the last chunk
26    pub is_last: bool,
27}
28
29impl ChunkDescriptor {
30    /// Create a new chunk descriptor.
31    pub fn new(offset: u64, length: usize, index: usize, total_chunks: usize) -> Self {
32        Self {
33            offset,
34            length,
35            index,
36            total_chunks,
37            is_last: index + 1 == total_chunks,
38        }
39    }
40
41    /// Get the end offset of this chunk.
42    pub fn end_offset(&self) -> u64 {
43        self.offset + self.length as u64
44    }
45}
46
47/// A buffer that manages chunked data.
48pub struct ChunkedBuffer {
49    /// The underlying buffer
50    inner: Arc<RwLock<ChunkedBufferInner>>,
51
52    /// Chunk size in bytes
53    chunk_size: usize,
54
55    /// Maximum buffer size in bytes
56    max_size: usize,
57}
58
59struct ChunkedBufferInner {
60    /// Queue of buffered chunks
61    chunks: VecDeque<BufferedChunk>,
62
63    /// Current size in bytes
64    current_size: usize,
65
66    /// Next chunk index to read
67    next_read_index: usize,
68
69    /// Next chunk index to write
70    next_write_index: usize,
71
72    /// Total number of chunks
73    total_chunks: Option<usize>,
74
75    /// Whether writing is complete
76    write_complete: bool,
77}
78
79struct BufferedChunk {
80    descriptor: ChunkDescriptor,
81    data: Bytes,
82}
83
84impl ChunkedBuffer {
85    /// Create a new chunked buffer.
86    pub fn new(chunk_size: usize, max_size: usize) -> Self {
87        Self {
88            inner: Arc::new(RwLock::new(ChunkedBufferInner {
89                chunks: VecDeque::new(),
90                current_size: 0,
91                next_read_index: 0,
92                next_write_index: 0,
93                total_chunks: None,
94                write_complete: false,
95            })),
96            chunk_size,
97            max_size,
98        }
99    }
100
101    /// Create a new chunked buffer with default settings.
102    pub fn with_defaults() -> Self {
103        Self::new(1024 * 1024, 100 * 1024 * 1024) // 1MB chunks, 100MB max
104    }
105
106    /// Calculate the number of chunks needed for a given size.
107    pub fn calculate_chunks(&self, total_size: u64) -> usize {
108        total_size.div_ceil(self.chunk_size as u64) as usize
109    }
110
111    /// Get a chunk descriptor for a given index.
112    pub fn descriptor_for_index(&self, index: usize, total_size: u64) -> ChunkDescriptor {
113        let total_chunks = self.calculate_chunks(total_size);
114        let offset = (index as u64) * (self.chunk_size as u64);
115        let remaining = total_size.saturating_sub(offset);
116        let length = remaining.min(self.chunk_size as u64) as usize;
117
118        ChunkDescriptor::new(offset, length, index, total_chunks)
119    }
120
121    /// Push a chunk into the buffer.
122    pub async fn push(&self, descriptor: ChunkDescriptor, data: Bytes) -> Result<()> {
123        let mut inner = self.inner.write().await;
124
125        // Check if buffer is full
126        if inner.current_size + data.len() > self.max_size {
127            return Err(StreamingError::BufferFull);
128        }
129
130        // Verify chunk index
131        if descriptor.index != inner.next_write_index {
132            return Err(StreamingError::InvalidOperation(format!(
133                "Expected chunk {}, got {}",
134                inner.next_write_index, descriptor.index
135            )));
136        }
137
138        inner.chunks.push_back(BufferedChunk {
139            descriptor: descriptor.clone(),
140            data,
141        });
142
143        inner.current_size += descriptor.length;
144        inner.next_write_index += 1;
145
146        if let Some(total) = inner.total_chunks {
147            if descriptor.index + 1 == total {
148                inner.write_complete = true;
149            }
150        } else if descriptor.is_last {
151            inner.total_chunks = Some(descriptor.total_chunks);
152            inner.write_complete = true;
153        }
154
155        debug!(
156            "Pushed chunk {} ({} bytes), buffer size: {}",
157            descriptor.index, descriptor.length, inner.current_size
158        );
159
160        Ok(())
161    }
162
163    /// Pop a chunk from the buffer.
164    pub async fn pop(&self) -> Result<Option<(ChunkDescriptor, Bytes)>> {
165        let mut inner = self.inner.write().await;
166
167        if inner.chunks.is_empty() {
168            if inner.write_complete {
169                return Ok(None);
170            } else {
171                return Err(StreamingError::Other("No chunks available".to_string()));
172            }
173        }
174
175        let chunk = inner
176            .chunks
177            .pop_front()
178            .ok_or_else(|| StreamingError::Other("Failed to pop chunk".to_string()))?;
179
180        inner.current_size = inner.current_size.saturating_sub(chunk.descriptor.length);
181        inner.next_read_index += 1;
182
183        debug!(
184            "Popped chunk {} ({} bytes), buffer size: {}",
185            chunk.descriptor.index, chunk.descriptor.length, inner.current_size
186        );
187
188        Ok(Some((chunk.descriptor, chunk.data)))
189    }
190
191    /// Peek at the next chunk without removing it.
192    pub async fn peek(&self) -> Result<Option<ChunkDescriptor>> {
193        let inner = self.inner.read().await;
194        Ok(inner.chunks.front().map(|c| c.descriptor.clone()))
195    }
196
197    /// Get the number of chunks currently in the buffer.
198    pub async fn len(&self) -> usize {
199        let inner = self.inner.read().await;
200        inner.chunks.len()
201    }
202
203    /// Check if the buffer is empty.
204    pub async fn is_empty(&self) -> bool {
205        let inner = self.inner.read().await;
206        inner.chunks.is_empty()
207    }
208
209    /// Get the current buffer size in bytes.
210    pub async fn size_bytes(&self) -> usize {
211        let inner = self.inner.read().await;
212        inner.current_size
213    }
214
215    /// Check if writing is complete.
216    pub async fn is_complete(&self) -> bool {
217        let inner = self.inner.read().await;
218        inner.write_complete
219    }
220
221    /// Clear all chunks from the buffer.
222    pub async fn clear(&self) {
223        let mut inner = self.inner.write().await;
224        inner.chunks.clear();
225        inner.current_size = 0;
226        debug!("Buffer cleared");
227    }
228
229    /// Get buffer statistics.
230    pub async fn stats(&self) -> BufferStats {
231        let inner = self.inner.read().await;
232        BufferStats {
233            chunks_buffered: inner.chunks.len(),
234            bytes_buffered: inner.current_size,
235            max_bytes: self.max_size,
236            utilization: (inner.current_size as f64) / (self.max_size as f64),
237            chunks_read: inner.next_read_index,
238            chunks_written: inner.next_write_index,
239            total_chunks: inner.total_chunks,
240            complete: inner.write_complete,
241        }
242    }
243}
244
245/// Buffer statistics.
246#[derive(Debug, Clone)]
247pub struct BufferStats {
248    /// Number of chunks currently buffered
249    pub chunks_buffered: usize,
250
251    /// Number of bytes currently buffered
252    pub bytes_buffered: usize,
253
254    /// Maximum buffer size in bytes
255    pub max_bytes: usize,
256
257    /// Buffer utilization (0.0 to 1.0)
258    pub utilization: f64,
259
260    /// Number of chunks read
261    pub chunks_read: usize,
262
263    /// Number of chunks written
264    pub chunks_written: usize,
265
266    /// Total number of chunks (if known)
267    pub total_chunks: Option<usize>,
268
269    /// Whether writing is complete
270    pub complete: bool,
271}
272
273impl BufferStats {
274    /// Calculate progress percentage.
275    pub fn progress(&self) -> Option<f64> {
276        self.total_chunks.map(|total| {
277            if total == 0 {
278                100.0
279            } else {
280                (self.chunks_read as f64 / total as f64) * 100.0
281            }
282        })
283    }
284}
285
286/// A circular buffer for efficient chunk management.
287pub struct CircularChunkBuffer {
288    /// The underlying buffer
289    buffer: Vec<u8>,
290
291    /// Read position
292    read_pos: usize,
293
294    /// Write position
295    write_pos: usize,
296
297    /// Number of bytes available
298    available: usize,
299
300    /// Buffer capacity
301    capacity: usize,
302}
303
304impl CircularChunkBuffer {
305    /// Create a new circular buffer.
306    pub fn new(capacity: usize) -> Self {
307        Self {
308            buffer: vec![0; capacity],
309            read_pos: 0,
310            write_pos: 0,
311            available: 0,
312            capacity,
313        }
314    }
315
316    /// Write data to the buffer.
317    pub fn write(&mut self, data: &[u8]) -> Result<usize> {
318        let space_available = self.capacity - self.available;
319        let to_write = data.len().min(space_available);
320
321        if to_write == 0 {
322            return Ok(0);
323        }
324
325        let end_pos = self.write_pos + to_write;
326        if end_pos <= self.capacity {
327            // Simple case: write doesn't wrap
328            self.buffer[self.write_pos..end_pos].copy_from_slice(&data[..to_write]);
329            self.write_pos = end_pos % self.capacity;
330        } else {
331            // Write wraps around
332            let first_part = self.capacity - self.write_pos;
333            self.buffer[self.write_pos..].copy_from_slice(&data[..first_part]);
334            self.buffer[..to_write - first_part].copy_from_slice(&data[first_part..to_write]);
335            self.write_pos = to_write - first_part;
336        }
337
338        self.available += to_write;
339        Ok(to_write)
340    }
341
342    /// Read data from the buffer.
343    pub fn read(&mut self, buf: &mut [u8]) -> Result<usize> {
344        let to_read = buf.len().min(self.available);
345
346        if to_read == 0 {
347            return Ok(0);
348        }
349
350        let end_pos = self.read_pos + to_read;
351        if end_pos <= self.capacity {
352            // Simple case: read doesn't wrap
353            buf[..to_read].copy_from_slice(&self.buffer[self.read_pos..end_pos]);
354            self.read_pos = end_pos % self.capacity;
355        } else {
356            // Read wraps around
357            let first_part = self.capacity - self.read_pos;
358            buf[..first_part].copy_from_slice(&self.buffer[self.read_pos..]);
359            buf[first_part..to_read].copy_from_slice(&self.buffer[..to_read - first_part]);
360            self.read_pos = to_read - first_part;
361        }
362
363        self.available -= to_read;
364        Ok(to_read)
365    }
366
367    /// Get the number of bytes available to read.
368    pub fn available(&self) -> usize {
369        self.available
370    }
371
372    /// Get the amount of free space in the buffer.
373    pub fn space_available(&self) -> usize {
374        self.capacity - self.available
375    }
376
377    /// Check if the buffer is empty.
378    pub fn is_empty(&self) -> bool {
379        self.available == 0
380    }
381
382    /// Check if the buffer is full.
383    pub fn is_full(&self) -> bool {
384        self.available == self.capacity
385    }
386
387    /// Clear the buffer.
388    pub fn clear(&mut self) {
389        self.read_pos = 0;
390        self.write_pos = 0;
391        self.available = 0;
392    }
393}
394
395#[cfg(test)]
396mod tests {
397    use super::*;
398
399    #[tokio::test]
400    async fn test_chunked_buffer() {
401        let buffer = ChunkedBuffer::new(1024, 10240);
402
403        let desc = ChunkDescriptor::new(0, 1024, 0, 10);
404        let data = Bytes::from(vec![0u8; 1024]);
405
406        buffer.push(desc.clone(), data.clone()).await.ok();
407
408        assert_eq!(buffer.len().await, 1);
409        assert_eq!(buffer.size_bytes().await, 1024);
410
411        let popped = buffer.pop().await.ok().flatten();
412        assert!(popped.is_some());
413
414        assert_eq!(buffer.len().await, 0);
415    }
416
417    #[test]
418    fn test_circular_buffer() {
419        let mut buffer = CircularChunkBuffer::new(10);
420
421        // Write some data
422        let written = buffer.write(&[1, 2, 3, 4, 5]).ok();
423        assert_eq!(written, Some(5));
424        assert_eq!(buffer.available(), 5);
425
426        // Read some data
427        let mut read_buf = [0u8; 3];
428        let read = buffer.read(&mut read_buf).ok();
429        assert_eq!(read, Some(3));
430        assert_eq!(read_buf, [1, 2, 3]);
431        assert_eq!(buffer.available(), 2);
432
433        // Write more data (should wrap)
434        let written = buffer.write(&[6, 7, 8, 9, 10]).ok();
435        assert_eq!(written, Some(5));
436        assert_eq!(buffer.available(), 7);
437    }
438
439    #[test]
440    fn test_chunk_descriptor() {
441        let desc = ChunkDescriptor::new(0, 1024, 0, 10);
442        assert_eq!(desc.offset, 0);
443        assert_eq!(desc.length, 1024);
444        assert_eq!(desc.end_offset(), 1024);
445        assert!(!desc.is_last);
446
447        let last_desc = ChunkDescriptor::new(9216, 1024, 9, 10);
448        assert!(last_desc.is_last);
449    }
450}