dcp/stream/
ring_buffer.rs1use crate::DCPError;
7use std::sync::atomic::{AtomicU32, AtomicU64, Ordering};
8
9#[derive(Debug)]
11pub struct Backpressure {
12 available: AtomicU32,
14}
15
16impl Backpressure {
17 pub fn new(capacity: u32) -> Self {
19 Self {
20 available: AtomicU32::new(capacity),
21 }
22 }
23
24 #[inline]
26 pub fn available(&self) -> u32 {
27 self.available.load(Ordering::Acquire)
28 }
29
30 #[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 #[inline]
52 pub fn release(&self, amount: u32) {
53 self.available.fetch_add(amount, Ordering::Release);
54 }
55
56 #[inline]
58 pub fn is_full(&self) -> bool {
59 self.available() == 0
60 }
61}
62
63pub struct StreamRingBuffer {
68 buffer: Box<[u8]>,
70 write_pos: AtomicU64,
72 read_pos: AtomicU64,
74 capacity: usize,
76 mask: usize,
78 backpressure: Backpressure,
80}
81
82impl StreamRingBuffer {
83 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 #[inline]
99 pub fn capacity(&self) -> usize {
100 self.capacity
101 }
102
103 #[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 #[inline]
113 pub fn is_empty(&self) -> bool {
114 self.len() == 0
115 }
116
117 #[inline]
119 pub fn available_space(&self) -> usize {
120 self.capacity - self.len()
121 }
122
123 pub fn backpressure(&self) -> &Backpressure {
125 &self.backpressure
126 }
127
128 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 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 let buffer = unsafe {
146 std::slice::from_raw_parts_mut(self.buffer.as_ptr() as *mut u8, self.capacity)
147 };
148
149 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 self.write_pos.store(write + len as u64, Ordering::Release);
160 Ok(())
161 }
162
163 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 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 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 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 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 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 rb.push(&[1u8; 12]).unwrap();
271
272 let mut buf = [0u8; 8];
274 rb.pop(&mut buf);
275
276 rb.push(&[2u8; 10]).unwrap();
278
279 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 rb.push(&[0u8; 16]).unwrap();
293
294 assert_eq!(rb.push(&[0u8; 1]), Err(DCPError::Backpressure));
296
297 let mut buf = [0u8; 8];
299 rb.pop(&mut buf);
300
301 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 assert_eq!(rb.len(), 9);
317 }
318}