Skip to main content

smoltcp/storage/
packet_buffer.rs

1use managed::ManagedSlice;
2
3use crate::storage::{Full, RingBuffer};
4
5use super::Empty;
6
7/// Size and header of a packet.
8#[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    /// Empty packet description.
17    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/// An UDP packet ring buffer.
42#[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    /// Create a new packet buffer with the provided metadata and payload storage.
50    ///
51    /// Metadata storage limits the maximum _number_ of packets in the buffer and payload
52    /// storage limits the maximum _total size_ of packets.
53    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    /// Query whether the buffer is empty.
65    pub fn is_empty(&self) -> bool {
66        self.metadata_ring.is_empty()
67    }
68
69    /// Query whether the buffer is full.
70    pub fn is_full(&self) -> bool {
71        self.metadata_ring.is_full()
72    }
73
74    // There is currently no enqueue_with() because of the complexity of managing padding
75    // in case of failure.
76
77    /// Enqueue a single packet with the given header into the buffer, and
78    /// return a reference to its payload, or return `Err(Full)`
79    /// if the buffer is full.
80    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        // Ring is currently empty.  Clear it (resetting `read_at`) to maximize
86        // for contiguous space.
87        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                // The buffer length is larger than the current contiguous window
99                // and is larger than the contiguous window will be after adding
100                // the padding necessary to circle around to the beginning of the
101                // ring buffer.
102                return Err(Full);
103            } else {
104                // Add padding to the end of the ring buffer so that the
105                // contiguous window is at the beginning of the ring buffer.
106                *self.metadata_ring.enqueue_one()? = PacketMetadata::padding(contig_window);
107                // note(discard): function does not write to the result
108                // enqueued padding buffer location
109                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    /// Call `f` with a packet from the buffer large enough to fit `max_size` bytes. The packet
121    /// is shrunk to the size returned from `f` and enqueued into the buffer.
122    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                // The buffer length is larger than the current contiguous window
143                // and is larger than the contiguous window will be after adding
144                // the padding necessary to circle around to the beginning of the
145                // ring buffer.
146                return Err(Full);
147            } else {
148                // Add padding to the end of the ring buffer so that the
149                // contiguous window is at the beginning of the ring buffer.
150                *self.metadata_ring.enqueue_one()? = PacketMetadata::padding(contig_window);
151                // note(discard): function does not write to the result
152                // enqueued padding buffer location
153                let _buf_enqueued = self.payload_ring.enqueue_many(contig_window);
154            }
155        }
156
157        let metadata_slot = self.metadata_ring.enqueue_one()?;
158
159        // Only call f once we know that we will succeed
160        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                // note(discard): function does not use value of dequeued padding bytes
173                let _buf_dequeued = self.payload_ring.dequeue_many(metadata.size);
174                Ok(()) // dequeue metadata
175            } else {
176                Err(()) // don't dequeue metadata
177            }
178        });
179    }
180
181    /// Call `f` with a single packet from the buffer, and dequeue the packet if `f`
182    /// returns successfully, or return `Err(EmptyError)` if the buffer is empty.
183    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    /// Dequeue a single packet from the buffer, and return a reference to its payload
207    /// as well as its header, or return `Err(Error::Exhausted)` if the buffer is empty.
208    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    /// Peek at a single packet from the buffer without removing it, and return a reference to
219    /// its payload as well as its header, or return `Err(Error:Exhausted)` if the buffer is empty.
220    ///
221    /// This function otherwise behaves identically to [dequeue](#method.dequeue).
222    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    /// Return the maximum number packets that can be stored.
236    pub fn packet_capacity(&self) -> usize {
237        self.metadata_ring.capacity()
238    }
239
240    /// Return the maximum number of bytes in the payload ring buffer.
241    pub fn payload_capacity(&self) -> usize {
242        self.payload_ring.capacity()
243    }
244
245    /// Return the current number of bytes in the payload ring buffer.
246    pub fn payload_bytes_count(&self) -> usize {
247        self.payload_ring.len()
248    }
249
250    /// Reset the packet buffer and clear any staged.
251    #[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        // Fill the metadata buffer except 1 byte and then make room at the start
402        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        // At this point there is room at the start of the payload storage
412        // and 1 byte at the end of the payload storage
413        // but only one metadata slot.
414        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        // Ensure enqueuing data in the buffer fills it somewhat.
428        assert!(buffer.is_empty());
429        assert!(buffer.enqueue(6, ()).is_ok());
430
431        // Ensure that resetting the buffer causes it to be empty.
432        assert!(!buffer.is_empty());
433        buffer.reset();
434        assert!(buffer.is_empty());
435    }
436}