1use slab::Slab;
4use spin::Mutex;
5use std::sync::Arc;
6use thiserror::Error;
7
8use crate::double_mapped_buffer::{DoubleMappedBuffer, DoubleMappedBufferError};
9
10#[derive(Error, Debug)]
12pub enum CircularError {
13 #[error("Failed to allocate double mapped buffer.")]
15 Allocation(DoubleMappedBufferError),
16}
17
18pub use crate::{Metadata, NoMetadata, Notifier};
19
20pub struct Circular;
22
23impl Circular {
24 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
73pub 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 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 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 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
212pub 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 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 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 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}