Skip to main content

dcp/stream/
ring_buffer.rs

1//! Lock-free ring buffer for streaming support.
2//!
3//! Implements a single-producer single-consumer (SPSC) ring buffer
4//! with atomic operations for thread-safe streaming.
5
6use crate::DCPError;
7use std::sync::atomic::{AtomicU32, AtomicU64, Ordering};
8
9/// Backpressure signal for flow control
10#[derive(Debug)]
11pub struct Backpressure {
12    /// Available capacity in bytes
13    available: AtomicU32,
14}
15
16impl Backpressure {
17    /// Create a new backpressure signal with initial capacity
18    pub fn new(capacity: u32) -> Self {
19        Self {
20            available: AtomicU32::new(capacity),
21        }
22    }
23
24    /// Get available capacity
25    #[inline]
26    pub fn available(&self) -> u32 {
27        self.available.load(Ordering::Acquire)
28    }
29
30    /// Try to reserve capacity, returns true if successful
31    #[inline]
32    pub fn try_reserve(&self, amount: u32) -> bool {
33        let mut current = self.available.load(Ordering::Acquire);
34        loop {
35            if current < amount {
36                return false;
37            }
38            match self.available.compare_exchange_weak(
39                current,
40                current - amount,
41                Ordering::AcqRel,
42                Ordering::Acquire,
43            ) {
44                Ok(_) => return true,
45                Err(new_current) => current = new_current,
46            }
47        }
48    }
49
50    /// Release capacity back
51    #[inline]
52    pub fn release(&self, amount: u32) {
53        self.available.fetch_add(amount, Ordering::Release);
54    }
55
56    /// Check if backpressure is active (no capacity available)
57    #[inline]
58    pub fn is_full(&self) -> bool {
59        self.available() == 0
60    }
61}
62
63/// Lock-free ring buffer for streaming
64///
65/// Uses atomic operations for single-producer single-consumer (SPSC) pattern.
66/// The buffer is power-of-two sized for efficient modulo operations.
67pub struct StreamRingBuffer {
68    /// Ring buffer storage
69    buffer: Box<[u8]>,
70    /// Write position (producer) - monotonically increasing
71    write_pos: AtomicU64,
72    /// Read position (consumer) - monotonically increasing
73    read_pos: AtomicU64,
74    /// Capacity (must be power of 2)
75    capacity: usize,
76    /// Mask for efficient modulo (capacity - 1)
77    mask: usize,
78    /// Backpressure signal
79    backpressure: Backpressure,
80}
81
82impl StreamRingBuffer {
83    /// Create a new ring buffer with the given capacity.
84    /// Capacity will be rounded up to the next power of 2.
85    pub fn new(capacity: usize) -> Self {
86        let capacity = capacity.next_power_of_two().max(16);
87        Self {
88            buffer: vec![0u8; capacity].into_boxed_slice(),
89            write_pos: AtomicU64::new(0),
90            read_pos: AtomicU64::new(0),
91            capacity,
92            mask: capacity - 1,
93            backpressure: Backpressure::new(capacity as u32),
94        }
95    }
96
97    /// Get the capacity of the buffer
98    #[inline]
99    pub fn capacity(&self) -> usize {
100        self.capacity
101    }
102
103    /// Get the number of bytes available to read
104    #[inline]
105    pub fn len(&self) -> usize {
106        let write = self.write_pos.load(Ordering::Acquire);
107        let read = self.read_pos.load(Ordering::Acquire);
108        (write - read) as usize
109    }
110
111    /// Check if the buffer is empty
112    #[inline]
113    pub fn is_empty(&self) -> bool {
114        self.len() == 0
115    }
116
117    /// Get available space for writing
118    #[inline]
119    pub fn available_space(&self) -> usize {
120        self.capacity - self.len()
121    }
122
123    /// Get the backpressure signal
124    pub fn backpressure(&self) -> &Backpressure {
125        &self.backpressure
126    }
127
128    /// Push data into the buffer.
129    /// Returns Err(Backpressure) if there's not enough space.
130    pub fn push(&self, data: &[u8]) -> Result<(), DCPError> {
131        let len = data.len();
132        if len > self.capacity {
133            return Err(DCPError::OutOfBounds);
134        }
135
136        // Check if we have space
137        if !self.backpressure.try_reserve(len as u32) {
138            return Err(DCPError::Backpressure);
139        }
140
141        let write = self.write_pos.load(Ordering::Acquire);
142        let start = (write as usize) & self.mask;
143
144        // Get mutable access to buffer (safe because we're single producer)
145        let buffer = unsafe {
146            std::slice::from_raw_parts_mut(self.buffer.as_ptr() as *mut u8, self.capacity)
147        };
148
149        // Handle wrap-around
150        if start + len <= self.capacity {
151            buffer[start..start + len].copy_from_slice(data);
152        } else {
153            let first_part = self.capacity - start;
154            buffer[start..].copy_from_slice(&data[..first_part]);
155            buffer[..len - first_part].copy_from_slice(&data[first_part..]);
156        }
157
158        // Update write position with release ordering
159        self.write_pos.store(write + len as u64, Ordering::Release);
160        Ok(())
161    }
162
163    /// Pop data from the buffer into the provided slice.
164    /// Returns the number of bytes read.
165    pub fn pop(&self, buf: &mut [u8]) -> usize {
166        let available = self.len();
167        let to_read = buf.len().min(available);
168
169        if to_read == 0 {
170            return 0;
171        }
172
173        let read = self.read_pos.load(Ordering::Acquire);
174        let start = (read as usize) & self.mask;
175
176        // Handle wrap-around
177        if start + to_read <= self.capacity {
178            buf[..to_read].copy_from_slice(&self.buffer[start..start + to_read]);
179        } else {
180            let first_part = self.capacity - start;
181            buf[..first_part].copy_from_slice(&self.buffer[start..]);
182            buf[first_part..to_read].copy_from_slice(&self.buffer[..to_read - first_part]);
183        }
184
185        // Update read position and release backpressure
186        self.read_pos
187            .store(read + to_read as u64, Ordering::Release);
188        self.backpressure.release(to_read as u32);
189
190        to_read
191    }
192
193    /// Peek at data without consuming it.
194    /// Returns the number of bytes peeked.
195    pub fn peek(&self, buf: &mut [u8]) -> usize {
196        let available = self.len();
197        let to_read = buf.len().min(available);
198
199        if to_read == 0 {
200            return 0;
201        }
202
203        let read = self.read_pos.load(Ordering::Acquire);
204        let start = (read as usize) & self.mask;
205
206        // Handle wrap-around
207        if start + to_read <= self.capacity {
208            buf[..to_read].copy_from_slice(&self.buffer[start..start + to_read]);
209        } else {
210            let first_part = self.capacity - start;
211            buf[..first_part].copy_from_slice(&self.buffer[start..]);
212            buf[first_part..to_read].copy_from_slice(&self.buffer[..to_read - first_part]);
213        }
214
215        to_read
216    }
217
218    /// Clear the buffer
219    pub fn clear(&self) {
220        let write = self.write_pos.load(Ordering::Acquire);
221        let read = self.read_pos.load(Ordering::Acquire);
222        let consumed = (write - read) as u32;
223
224        self.read_pos.store(write, Ordering::Release);
225        self.backpressure.release(consumed);
226    }
227}
228
229#[cfg(test)]
230mod tests {
231    use super::*;
232
233    #[test]
234    fn test_backpressure() {
235        let bp = Backpressure::new(100);
236        assert_eq!(bp.available(), 100);
237        assert!(!bp.is_full());
238
239        assert!(bp.try_reserve(50));
240        assert_eq!(bp.available(), 50);
241
242        assert!(!bp.try_reserve(60));
243        assert_eq!(bp.available(), 50);
244
245        bp.release(30);
246        assert_eq!(bp.available(), 80);
247    }
248
249    #[test]
250    fn test_ring_buffer_basic() {
251        let rb = StreamRingBuffer::new(64);
252        assert!(rb.is_empty());
253        assert_eq!(rb.capacity(), 64);
254
255        rb.push(b"hello").unwrap();
256        assert_eq!(rb.len(), 5);
257
258        let mut buf = [0u8; 10];
259        let read = rb.pop(&mut buf);
260        assert_eq!(read, 5);
261        assert_eq!(&buf[..5], b"hello");
262        assert!(rb.is_empty());
263    }
264
265    #[test]
266    fn test_ring_buffer_wrap_around() {
267        let rb = StreamRingBuffer::new(16);
268
269        // Fill most of the buffer
270        rb.push(&[1u8; 12]).unwrap();
271
272        // Read some
273        let mut buf = [0u8; 8];
274        rb.pop(&mut buf);
275
276        // Write more (will wrap around)
277        rb.push(&[2u8; 10]).unwrap();
278
279        // Read all
280        let mut buf = [0u8; 14];
281        let read = rb.pop(&mut buf);
282        assert_eq!(read, 14);
283        assert_eq!(&buf[..4], &[1u8; 4]);
284        assert_eq!(&buf[4..14], &[2u8; 10]);
285    }
286
287    #[test]
288    fn test_ring_buffer_backpressure() {
289        let rb = StreamRingBuffer::new(16);
290
291        // Fill the buffer
292        rb.push(&[0u8; 16]).unwrap();
293
294        // Should fail with backpressure
295        assert_eq!(rb.push(&[0u8; 1]), Err(DCPError::Backpressure));
296
297        // Read some to make space
298        let mut buf = [0u8; 8];
299        rb.pop(&mut buf);
300
301        // Now should succeed
302        rb.push(&[0u8; 8]).unwrap();
303    }
304
305    #[test]
306    fn test_ring_buffer_peek() {
307        let rb = StreamRingBuffer::new(32);
308        rb.push(b"test data").unwrap();
309
310        let mut buf = [0u8; 4];
311        let peeked = rb.peek(&mut buf);
312        assert_eq!(peeked, 4);
313        assert_eq!(&buf, b"test");
314
315        // Data should still be there
316        assert_eq!(rb.len(), 9);
317    }
318}