1use crate::error::{Result, StreamingError};
4use bytes::Bytes;
5use std::collections::VecDeque;
6use std::sync::Arc;
7use tokio::sync::RwLock;
8use tracing::debug;
9
10#[derive(Debug, Clone)]
12pub struct ChunkDescriptor {
13 pub offset: u64,
15
16 pub length: usize,
18
19 pub index: usize,
21
22 pub total_chunks: usize,
24
25 pub is_last: bool,
27}
28
29impl ChunkDescriptor {
30 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 pub fn end_offset(&self) -> u64 {
43 self.offset + self.length as u64
44 }
45}
46
47pub struct ChunkedBuffer {
49 inner: Arc<RwLock<ChunkedBufferInner>>,
51
52 chunk_size: usize,
54
55 max_size: usize,
57}
58
59struct ChunkedBufferInner {
60 chunks: VecDeque<BufferedChunk>,
62
63 current_size: usize,
65
66 next_read_index: usize,
68
69 next_write_index: usize,
71
72 total_chunks: Option<usize>,
74
75 write_complete: bool,
77}
78
79struct BufferedChunk {
80 descriptor: ChunkDescriptor,
81 data: Bytes,
82}
83
84impl ChunkedBuffer {
85 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 pub fn with_defaults() -> Self {
103 Self::new(1024 * 1024, 100 * 1024 * 1024) }
105
106 pub fn calculate_chunks(&self, total_size: u64) -> usize {
108 total_size.div_ceil(self.chunk_size as u64) as usize
109 }
110
111 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 pub async fn push(&self, descriptor: ChunkDescriptor, data: Bytes) -> Result<()> {
123 let mut inner = self.inner.write().await;
124
125 if inner.current_size + data.len() > self.max_size {
127 return Err(StreamingError::BufferFull);
128 }
129
130 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 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 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 pub async fn len(&self) -> usize {
199 let inner = self.inner.read().await;
200 inner.chunks.len()
201 }
202
203 pub async fn is_empty(&self) -> bool {
205 let inner = self.inner.read().await;
206 inner.chunks.is_empty()
207 }
208
209 pub async fn size_bytes(&self) -> usize {
211 let inner = self.inner.read().await;
212 inner.current_size
213 }
214
215 pub async fn is_complete(&self) -> bool {
217 let inner = self.inner.read().await;
218 inner.write_complete
219 }
220
221 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 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#[derive(Debug, Clone)]
247pub struct BufferStats {
248 pub chunks_buffered: usize,
250
251 pub bytes_buffered: usize,
253
254 pub max_bytes: usize,
256
257 pub utilization: f64,
259
260 pub chunks_read: usize,
262
263 pub chunks_written: usize,
265
266 pub total_chunks: Option<usize>,
268
269 pub complete: bool,
271}
272
273impl BufferStats {
274 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
286pub struct CircularChunkBuffer {
288 buffer: Vec<u8>,
290
291 read_pos: usize,
293
294 write_pos: usize,
296
297 available: usize,
299
300 capacity: usize,
302}
303
304impl CircularChunkBuffer {
305 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 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 self.buffer[self.write_pos..end_pos].copy_from_slice(&data[..to_write]);
329 self.write_pos = end_pos % self.capacity;
330 } else {
331 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 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 buf[..to_read].copy_from_slice(&self.buffer[self.read_pos..end_pos]);
354 self.read_pos = end_pos % self.capacity;
355 } else {
356 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 pub fn available(&self) -> usize {
369 self.available
370 }
371
372 pub fn space_available(&self) -> usize {
374 self.capacity - self.available
375 }
376
377 pub fn is_empty(&self) -> bool {
379 self.available == 0
380 }
381
382 pub fn is_full(&self) -> bool {
384 self.available == self.capacity
385 }
386
387 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 let written = buffer.write(&[1, 2, 3, 4, 5]).ok();
423 assert_eq!(written, Some(5));
424 assert_eq!(buffer.available(), 5);
425
426 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 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}