Skip to main content

vmcircbuffer/
generic.rs

1//! Circular Buffer with generic [Notifier] to implement custom wait/block behavior.
2
3use slab::Slab;
4use spin::Mutex;
5use std::sync::Arc;
6use thiserror::Error;
7
8use crate::double_mapped_buffer::{DoubleMappedBuffer, DoubleMappedBufferError};
9
10/// Error setting up the underlying buffer.
11#[derive(Error, Debug)]
12pub enum CircularError {
13    /// Failed to allocate double mapped buffer.
14    #[error("Failed to allocate double mapped buffer.")]
15    Allocation(DoubleMappedBufferError),
16}
17
18pub use crate::{Metadata, NoMetadata, Notifier};
19
20/// Gerneric Circular Buffer Constructor
21pub struct Circular;
22
23impl Circular {
24    /// Create a buffer that can hold at least `min_items` items of type `T`.
25    ///
26    /// The size is the least common multiple of the page size and the size of `T`.
27    pub fn with_capacity<T, N, M>(min_items: usize) -> Result<Writer<T, N, M>, CircularError>
28    where
29        T: Copy + Default,
30        N: Notifier,
31        M: Metadata,
32    {
33        let buffer = match DoubleMappedBuffer::new(min_items) {
34            Ok(buffer) => Arc::new(buffer),
35            Err(e) => return Err(CircularError::Allocation(e)),
36        };
37
38        let state = Arc::new(Mutex::new(State {
39            writer_offset: 0,
40            writer_ab: false,
41            writer_done: false,
42            readers: Slab::new(),
43        }));
44
45        let writer = Writer {
46            buffer,
47            state,
48            last_space: 0,
49        };
50
51        Ok(writer)
52    }
53}
54
55struct State<N, M>
56where
57    N: Notifier,
58    M: Metadata,
59{
60    writer_offset: usize,
61    writer_ab: bool,
62    writer_done: bool,
63    readers: Slab<ReaderState<N, M>>,
64}
65struct ReaderState<N, M> {
66    ab: bool,
67    offset: usize,
68    reader_notifier: N,
69    writer_notifier: N,
70    meta: M,
71}
72
73/// Writer for a generic circular buffer with items of type `T` and [Notifier] of type `N`.
74pub struct Writer<T, N, M>
75where
76    N: Notifier,
77    M: Metadata,
78{
79    last_space: usize,
80    buffer: Arc<DoubleMappedBuffer<T>>,
81    state: Arc<Mutex<State<N, M>>>,
82}
83
84impl<T, N, M> Writer<T, N, M>
85where
86    N: Notifier,
87    M: Metadata,
88{
89    #[inline(always)]
90    fn writer_space(capacity: usize, w_off: usize, w_ab: bool, r_off: usize, r_ab: bool) -> usize {
91        if w_off > r_off {
92            r_off + capacity - w_off
93        } else if w_off < r_off {
94            r_off - w_off
95        } else if r_ab == w_ab {
96            capacity
97        } else {
98            0
99        }
100    }
101
102    /// Add a [Reader] to the buffer.
103    pub fn add_reader(&self, reader_notifier: N, writer_notifier: N) -> Reader<T, N, M> {
104        let mut state = self.state.lock();
105        let reader_state = ReaderState {
106            ab: state.writer_ab,
107            offset: state.writer_offset,
108            reader_notifier,
109            writer_notifier,
110            meta: M::new(),
111        };
112        let id = state.readers.insert(reader_state);
113
114        Reader {
115            id,
116            last_space: 0,
117            buffer: self.buffer.clone(),
118            state: self.state.clone(),
119        }
120    }
121
122    #[inline(always)]
123    fn space_and_offset_locked(
124        state: &mut State<N, M>,
125        capacity: usize,
126        arm: bool,
127    ) -> (usize, usize) {
128        let w_off = state.writer_offset;
129        let w_ab = state.writer_ab;
130
131        let mut space = capacity;
132
133        for (_, reader) in state.readers.iter_mut() {
134            let s = Self::writer_space(capacity, w_off, w_ab, reader.offset, reader.ab);
135
136            space = std::cmp::min(space, s);
137
138            if s == 0 && arm {
139                reader.writer_notifier.arm();
140                break;
141            }
142            if s == 0 {
143                break;
144            }
145        }
146
147        (space, w_off)
148    }
149
150    /// Get a slice for the output buffer space. Might be empty.
151    pub fn slice(&mut self, arm: bool) -> &mut [T] {
152        let mut state = self.state.lock();
153        let (space, offset) =
154            Self::space_and_offset_locked(&mut state, self.buffer.capacity(), arm);
155        self.last_space = space;
156        unsafe { &mut self.buffer.slice_with_offset_mut(offset)[0..space] }
157    }
158
159    /// Indicates that `n` items were written to the output buffer.
160    ///
161    /// It is ok if `n` is zero.
162    ///
163    /// # Panics
164    ///
165    /// If produced more than space was available in the last provided slice.
166    pub fn produce(&mut self, n: usize, meta: &[M::Item]) {
167        if n == 0 {
168            return;
169        }
170
171        assert!(n <= self.last_space, "vmcircbuffer: produced too much");
172        self.last_space -= n;
173
174        let mut state = self.state.lock();
175        let capacity = self.buffer.capacity();
176
177        debug_assert!(Self::space_and_offset_locked(&mut state, capacity, false).0 >= n);
178
179        let w_off = state.writer_offset;
180        let w_ab = state.writer_ab;
181
182        for (_, r) in state.readers.iter_mut() {
183            let space = Reader::<T, N, M>::reader_space(capacity, w_off, w_ab, r.offset, r.ab);
184
185            if !meta.is_empty() {
186                r.meta.add_from_slice(space, meta);
187            }
188            r.reader_notifier.notify();
189        }
190
191        if state.writer_offset + n >= self.buffer.capacity() {
192            state.writer_ab = !state.writer_ab;
193        }
194        state.writer_offset = (state.writer_offset + n) % self.buffer.capacity();
195    }
196}
197
198impl<T, N, M> Drop for Writer<T, N, M>
199where
200    N: Notifier,
201    M: Metadata,
202{
203    fn drop(&mut self) {
204        let mut state = self.state.lock();
205        state.writer_done = true;
206        for (_, r) in state.readers.iter_mut() {
207            r.reader_notifier.notify();
208        }
209    }
210}
211
212/// Reader for a generic circular buffer with items of type `T` and [Notifier] of type `N`.
213pub struct Reader<T, N, M>
214where
215    N: Notifier,
216    M: Metadata,
217{
218    id: usize,
219    last_space: usize,
220    buffer: Arc<DoubleMappedBuffer<T>>,
221    state: Arc<Mutex<State<N, M>>>,
222}
223
224impl<T, N, M> Reader<T, N, M>
225where
226    N: Notifier,
227    M: Metadata,
228{
229    #[inline(always)]
230    fn reader_space(capacity: usize, w_off: usize, w_ab: bool, r_off: usize, r_ab: bool) -> usize {
231        if r_off > w_off {
232            w_off + capacity - r_off
233        } else if r_off < w_off {
234            w_off - r_off
235        } else if r_ab == w_ab {
236            0
237        } else {
238            capacity
239        }
240    }
241
242    #[inline(always)]
243    fn space_and_offset_locked(
244        state: &mut State<N, M>,
245        id: usize,
246        capacity: usize,
247        arm: bool,
248    ) -> (usize, usize, bool) {
249        let done = state.writer_done;
250        let w_off = state.writer_offset;
251        let w_ab = state.writer_ab;
252
253        let my = unsafe { state.readers.get_unchecked_mut(id) };
254        let space = Self::reader_space(capacity, w_off, w_ab, my.offset, my.ab);
255
256        if space == 0 && arm {
257            my.reader_notifier.arm();
258        }
259
260        (space, my.offset, done)
261    }
262
263    /// Get a slice without fetching metadata.
264    pub fn slice(&mut self, arm: bool) -> Option<&[T]> {
265        let mut state = self.state.lock();
266        let (space, offset, done) =
267            Self::space_and_offset_locked(&mut state, self.id, self.buffer.capacity(), arm);
268        self.last_space = space;
269
270        if space == 0 && done {
271            return None;
272        }
273
274        unsafe { Some(&self.buffer.slice_with_offset(offset)[0..space]) }
275    }
276
277    /// Get a slice and copy metadata into `out` in one call.
278    pub fn slice_with_metadata_into(&mut self, arm: bool, out: &mut Vec<M::Item>) -> Option<&[T]> {
279        let mut state = self.state.lock();
280        let (space, offset, done) =
281            Self::space_and_offset_locked(&mut state, self.id, self.buffer.capacity(), arm);
282        let my = unsafe { state.readers.get_unchecked_mut(self.id) };
283
284        my.meta.get_into(out);
285        self.last_space = space;
286
287        if space == 0 && done {
288            out.clear();
289            return None;
290        }
291
292        unsafe { Some(&self.buffer.slice_with_offset(offset)[0..space]) }
293    }
294
295    /// Indicates that `n` items were read.
296    ///
297    /// # Panics
298    ///
299    /// If consumed more than space was available in the last provided slice.
300    pub fn consume(&mut self, n: usize) {
301        if n == 0 {
302            return;
303        }
304
305        assert!(n <= self.last_space, "vmcircbuffer: consumed too much!");
306        self.last_space -= n;
307
308        let mut state = self.state.lock();
309        debug_assert!(
310            Self::space_and_offset_locked(&mut state, self.id, self.buffer.capacity(), false).0
311                >= n
312        );
313        let my = unsafe { state.readers.get_unchecked_mut(self.id) };
314
315        my.meta.consume(n);
316
317        if my.offset + n >= self.buffer.capacity() {
318            my.ab = !my.ab;
319        }
320        my.offset = (my.offset + n) % self.buffer.capacity();
321
322        my.writer_notifier.notify();
323    }
324}
325
326impl<T, N, M> Drop for Reader<T, N, M>
327where
328    N: Notifier,
329    M: Metadata,
330{
331    fn drop(&mut self) {
332        let mut state = self.state.lock();
333        let mut s = state.readers.remove(self.id);
334        s.writer_notifier.notify();
335    }
336}