1use managed::ManagedSlice;
2
3use crate::storage::{Full, RingBuffer};
4
5use super::Empty;
6
7#[derive(Debug, Clone, Copy)]
9#[cfg_attr(feature = "defmt", derive(defmt::Format))]
10pub struct PacketMetadata<H> {
11 size: usize,
12 header: Option<H>,
13}
14
15impl<H> PacketMetadata<H> {
16 pub const EMPTY: PacketMetadata<H> = PacketMetadata {
18 size: 0,
19 header: None,
20 };
21
22 fn padding(size: usize) -> PacketMetadata<H> {
23 PacketMetadata {
24 size: size,
25 header: None,
26 }
27 }
28
29 fn packet(size: usize, header: H) -> PacketMetadata<H> {
30 PacketMetadata {
31 size: size,
32 header: Some(header),
33 }
34 }
35
36 fn is_padding(&self) -> bool {
37 self.header.is_none()
38 }
39}
40
41#[derive(Debug)]
43pub struct PacketBuffer<'a, H: 'a> {
44 metadata_ring: RingBuffer<'a, PacketMetadata<H>>,
45 payload_ring: RingBuffer<'a, u8>,
46}
47
48impl<'a, H> PacketBuffer<'a, H> {
49 pub fn new<MS, PS>(metadata_storage: MS, payload_storage: PS) -> PacketBuffer<'a, H>
54 where
55 MS: Into<ManagedSlice<'a, PacketMetadata<H>>>,
56 PS: Into<ManagedSlice<'a, u8>>,
57 {
58 PacketBuffer {
59 metadata_ring: RingBuffer::new(metadata_storage),
60 payload_ring: RingBuffer::new(payload_storage),
61 }
62 }
63
64 pub fn is_empty(&self) -> bool {
66 self.metadata_ring.is_empty()
67 }
68
69 pub fn is_full(&self) -> bool {
71 self.metadata_ring.is_full()
72 }
73
74 pub fn enqueue(&mut self, size: usize, header: H) -> Result<&mut [u8], Full> {
81 if self.payload_ring.capacity() < size || self.metadata_ring.is_full() {
82 return Err(Full);
83 }
84
85 if self.payload_ring.is_empty() {
88 self.payload_ring.clear();
89 }
90
91 let window = self.payload_ring.window();
92 let contig_window = self.payload_ring.contiguous_window();
93
94 if window < size {
95 return Err(Full);
96 } else if contig_window < size {
97 if window - contig_window < size {
98 return Err(Full);
103 } else {
104 *self.metadata_ring.enqueue_one()? = PacketMetadata::padding(contig_window);
107 let _buf_enqueued = self.payload_ring.enqueue_many(contig_window);
110 }
111 }
112
113 *self.metadata_ring.enqueue_one()? = PacketMetadata::packet(size, header);
114
115 let payload_buf = self.payload_ring.enqueue_many(size);
116 debug_assert!(payload_buf.len() == size);
117 Ok(payload_buf)
118 }
119
120 pub fn enqueue_with_infallible<'b, F>(
123 &'b mut self,
124 max_size: usize,
125 header: H,
126 f: F,
127 ) -> Result<usize, Full>
128 where
129 F: FnOnce(&'b mut [u8]) -> usize,
130 {
131 if self.payload_ring.capacity() < max_size || self.metadata_ring.is_full() {
132 return Err(Full);
133 }
134
135 let window = self.payload_ring.window();
136 let contig_window = self.payload_ring.contiguous_window();
137
138 if window < max_size {
139 return Err(Full);
140 } else if contig_window < max_size {
141 if window - contig_window < max_size {
142 return Err(Full);
147 } else {
148 *self.metadata_ring.enqueue_one()? = PacketMetadata::padding(contig_window);
151 let _buf_enqueued = self.payload_ring.enqueue_many(contig_window);
154 }
155 }
156
157 let metadata_slot = self.metadata_ring.enqueue_one()?;
158
159 let (size, _) = self
161 .payload_ring
162 .enqueue_many_with(|data| (f(&mut data[..max_size]), ()));
163
164 *metadata_slot = PacketMetadata::packet(size, header);
165
166 Ok(size)
167 }
168
169 fn dequeue_padding(&mut self) {
170 let _ = self.metadata_ring.dequeue_one_with(|metadata| {
171 if metadata.is_padding() {
172 let _buf_dequeued = self.payload_ring.dequeue_many(metadata.size);
174 Ok(()) } else {
176 Err(()) }
178 });
179 }
180
181 pub fn dequeue_with<'c, R, E, F>(&'c mut self, f: F) -> Result<Result<R, E>, Empty>
184 where
185 F: FnOnce(&mut H, &'c mut [u8]) -> Result<R, E>,
186 {
187 self.dequeue_padding();
188
189 self.metadata_ring.dequeue_one_with(|metadata| {
190 self.payload_ring
191 .dequeue_many_with(|payload_buf| {
192 debug_assert!(payload_buf.len() >= metadata.size);
193
194 match f(
195 metadata.header.as_mut().unwrap(),
196 &mut payload_buf[..metadata.size],
197 ) {
198 Ok(val) => (metadata.size, Ok(val)),
199 Err(err) => (0, Err(err)),
200 }
201 })
202 .1
203 })
204 }
205
206 pub fn dequeue(&mut self) -> Result<(H, &mut [u8]), Empty> {
209 self.dequeue_padding();
210
211 let meta = self.metadata_ring.dequeue_one()?;
212
213 let payload_buf = self.payload_ring.dequeue_many(meta.size);
214 debug_assert!(payload_buf.len() == meta.size);
215 Ok((meta.header.take().unwrap(), payload_buf))
216 }
217
218 pub fn peek(&mut self) -> Result<(&H, &[u8]), Empty> {
223 self.dequeue_padding();
224
225 if let Some(metadata) = self.metadata_ring.get_allocated(0, 1).first() {
226 Ok((
227 metadata.header.as_ref().unwrap(),
228 self.payload_ring.get_allocated(0, metadata.size),
229 ))
230 } else {
231 Err(Empty)
232 }
233 }
234
235 pub fn packet_capacity(&self) -> usize {
237 self.metadata_ring.capacity()
238 }
239
240 pub fn payload_capacity(&self) -> usize {
242 self.payload_ring.capacity()
243 }
244
245 pub fn payload_bytes_count(&self) -> usize {
247 self.payload_ring.len()
248 }
249
250 #[allow(unused)]
252 pub(crate) fn reset(&mut self) {
253 self.payload_ring.clear();
254 self.metadata_ring.clear();
255 }
256}
257
258#[cfg(test)]
259mod test {
260 use super::*;
261
262 fn buffer() -> PacketBuffer<'static, ()> {
263 PacketBuffer::new(vec![PacketMetadata::EMPTY; 4], vec![0u8; 16])
264 }
265
266 #[test]
267 fn test_simple() {
268 let mut buffer = buffer();
269 buffer.enqueue(6, ()).unwrap().copy_from_slice(b"abcdef");
270 assert_eq!(buffer.enqueue(16, ()), Err(Full));
271 assert_eq!(buffer.metadata_ring.len(), 1);
272 assert_eq!(buffer.dequeue().unwrap().1, &b"abcdef"[..]);
273 assert_eq!(buffer.dequeue(), Err(Empty));
274 }
275
276 #[test]
277 fn test_peek() {
278 let mut buffer = buffer();
279 assert_eq!(buffer.peek(), Err(Empty));
280 buffer.enqueue(6, ()).unwrap().copy_from_slice(b"abcdef");
281 assert_eq!(buffer.metadata_ring.len(), 1);
282 assert_eq!(buffer.peek().unwrap().1, &b"abcdef"[..]);
283 assert_eq!(buffer.dequeue().unwrap().1, &b"abcdef"[..]);
284 assert_eq!(buffer.peek(), Err(Empty));
285 }
286
287 #[test]
288 fn test_padding() {
289 let mut buffer = buffer();
290 assert!(buffer.enqueue(6, ()).is_ok());
291 assert!(buffer.enqueue(8, ()).is_ok());
292 assert!(buffer.dequeue().is_ok());
293 buffer.enqueue(4, ()).unwrap().copy_from_slice(b"abcd");
294 assert_eq!(buffer.metadata_ring.len(), 3);
295 assert!(buffer.dequeue().is_ok());
296
297 assert_eq!(buffer.dequeue().unwrap().1, &b"abcd"[..]);
298 assert_eq!(buffer.metadata_ring.len(), 0);
299 }
300
301 #[test]
302 fn test_padding_with_large_payload() {
303 let mut buffer = buffer();
304 assert!(buffer.enqueue(12, ()).is_ok());
305 assert!(buffer.dequeue().is_ok());
306 buffer
307 .enqueue(12, ())
308 .unwrap()
309 .copy_from_slice(b"abcdefghijkl");
310 }
311
312 #[test]
313 fn test_dequeue_with() {
314 let mut buffer = buffer();
315 assert!(buffer.enqueue(6, ()).is_ok());
316 assert!(buffer.enqueue(8, ()).is_ok());
317 assert!(buffer.dequeue().is_ok());
318 buffer.enqueue(4, ()).unwrap().copy_from_slice(b"abcd");
319 assert_eq!(buffer.metadata_ring.len(), 3);
320 assert!(buffer.dequeue().is_ok());
321
322 assert!(matches!(
323 buffer.dequeue_with(|_, _| Result::<(), u32>::Err(123)),
324 Ok(Err(_))
325 ));
326 assert_eq!(buffer.metadata_ring.len(), 1);
327
328 assert!(
329 buffer
330 .dequeue_with(|&mut (), payload| {
331 assert_eq!(payload, &b"abcd"[..]);
332 Result::<(), ()>::Ok(())
333 })
334 .is_ok()
335 );
336 assert_eq!(buffer.metadata_ring.len(), 0);
337 }
338
339 #[test]
340 fn test_metadata_full_empty() {
341 let mut buffer = buffer();
342 assert!(buffer.is_empty());
343 assert!(!buffer.is_full());
344 assert!(buffer.enqueue(1, ()).is_ok());
345 assert!(!buffer.is_empty());
346 assert!(buffer.enqueue(1, ()).is_ok());
347 assert!(buffer.enqueue(1, ()).is_ok());
348 assert!(!buffer.is_full());
349 assert!(!buffer.is_empty());
350 assert!(buffer.enqueue(1, ()).is_ok());
351 assert!(buffer.is_full());
352 assert!(!buffer.is_empty());
353 assert_eq!(buffer.metadata_ring.len(), 4);
354 assert_eq!(buffer.enqueue(1, ()), Err(Full));
355 }
356
357 #[test]
358 fn test_window_too_small() {
359 let mut buffer = buffer();
360 assert!(buffer.enqueue(4, ()).is_ok());
361 assert!(buffer.enqueue(8, ()).is_ok());
362 assert!(buffer.dequeue().is_ok());
363 assert_eq!(buffer.enqueue(16, ()), Err(Full));
364 assert_eq!(buffer.metadata_ring.len(), 1);
365 }
366
367 #[test]
368 fn test_contiguous_window_too_small() {
369 let mut buffer = buffer();
370 assert!(buffer.enqueue(4, ()).is_ok());
371 assert!(buffer.enqueue(8, ()).is_ok());
372 assert!(buffer.dequeue().is_ok());
373 assert_eq!(buffer.enqueue(8, ()), Err(Full));
374 assert_eq!(buffer.metadata_ring.len(), 1);
375 }
376
377 #[test]
378 fn test_contiguous_window_wrap() {
379 let mut buffer = buffer();
380 assert!(buffer.enqueue(15, ()).is_ok());
381 assert!(buffer.dequeue().is_ok());
382 assert!(buffer.enqueue(16, ()).is_ok());
383 }
384
385 #[test]
386 fn test_capacity_too_small() {
387 let mut buffer = buffer();
388 assert_eq!(buffer.enqueue(32, ()), Err(Full));
389 }
390
391 #[test]
392 fn test_contig_window_prioritized() {
393 let mut buffer = buffer();
394 assert!(buffer.enqueue(4, ()).is_ok());
395 assert!(buffer.dequeue().is_ok());
396 assert!(buffer.enqueue(5, ()).is_ok());
397 }
398
399 #[test]
400 fn test_enqueue_fallible_full_metadata_fn_not_called() {
401 let mut buffer = PacketBuffer::new(vec![PacketMetadata::EMPTY; 4], vec![0u8; 16]);
403
404 assert!(buffer.enqueue(12, ()).is_ok());
405 assert!(buffer.enqueue(1, ()).is_ok());
406 assert!(buffer.enqueue(1, ()).is_ok());
407 assert!(buffer.enqueue(1, ()).is_ok());
408 let dequeued_buf = buffer.dequeue().unwrap();
409 assert_eq!(dequeued_buf.1.len(), 12);
410
411 assert!(
415 buffer
416 .enqueue_with_infallible(5, (), |_| {
417 panic!("This enqueue should fail and this closure should not be called")
418 })
419 .is_err()
420 );
421 }
422
423 #[test]
424 fn clear() {
425 let mut buffer = buffer();
426
427 assert!(buffer.is_empty());
429 assert!(buffer.enqueue(6, ()).is_ok());
430
431 assert!(!buffer.is_empty());
433 buffer.reset();
434 assert!(buffer.is_empty());
435 }
436}