Skip to main content

prns_interfaces_embassy/bluetooth_auto/
frame_pool.rs

1use core::cell::Cell;
2use core::marker::PhantomData;
3
4use embassy_sync::blocking_mutex::raw::RawMutex;
5use embassy_sync::mutex::{Mutex, MutexGuard};
6use embassy_sync::semaphore::{FairSemaphore, Semaphore, SemaphoreReleaser};
7use heapless::Vec as FrameBytes;
8use portable_atomic::{AtomicBool, Ordering};
9
10type Availability<M, const WAITERS: usize> = FairSemaphore<M, WAITERS>;
11
12pub struct SharedFramePool<
13    M: RawMutex + 'static,
14    const FRAME: usize,
15    const CAPACITY: usize,
16    const WAITERS: usize,
17> {
18    availability: Availability<M, WAITERS>,
19    slots: [FrameSlot<M, FRAME>; CAPACITY],
20}
21
22struct FrameSlot<M: RawMutex + 'static, const FRAME: usize> {
23    claimed: AtomicBool,
24    frame: Mutex<M, FrameBytes<u8, FRAME>>,
25}
26
27impl<M: RawMutex + 'static, const FRAME: usize> FrameSlot<M, FRAME> {
28    const fn new() -> Self {
29        Self {
30            claimed: AtomicBool::new(false),
31            frame: Mutex::new(FrameBytes::new()),
32        }
33    }
34}
35
36impl<M: RawMutex + 'static, const FRAME: usize, const CAPACITY: usize, const WAITERS: usize>
37    SharedFramePool<M, FRAME, CAPACITY, WAITERS>
38{
39    #[must_use]
40    pub const fn new() -> Self {
41        assert!(CAPACITY > 0);
42        assert!(CAPACITY <= u8::MAX as usize + 1);
43        assert!(WAITERS > 0);
44        Self {
45            availability: FairSemaphore::new(CAPACITY),
46            slots: [const { FrameSlot::new() }; CAPACITY],
47        }
48    }
49
50    pub async fn lease(
51        &'static self,
52    ) -> Result<FrameLease<M, FRAME, CAPACITY, WAITERS>, FramePoolError> {
53        let permit = self
54            .availability
55            .acquire(1)
56            .await
57            .map_err(|_| FramePoolError::WaitQueueFull)?;
58        let lease = self
59            .claim(permit)
60            .ok_or(FramePoolError::PermitWithoutAvailableSlot)?;
61        lease.lock().await.clear();
62        Ok(lease)
63    }
64
65    pub fn try_lease(
66        &'static self,
67    ) -> Result<Option<FrameLease<M, FRAME, CAPACITY, WAITERS>>, FramePoolError> {
68        let Some(permit) = self.availability.try_acquire(1) else {
69            return Ok(None);
70        };
71        let lease = self
72            .claim(permit)
73            .ok_or(FramePoolError::PermitWithoutAvailableSlot)?;
74        lease
75            .slot()
76            .frame
77            .try_lock()
78            .map_err(|_| FramePoolError::SlotBusy)?
79            .clear();
80        Ok(Some(lease))
81    }
82
83    fn claim(
84        &'static self,
85        permit: SemaphoreReleaser<'static, Availability<M, WAITERS>>,
86    ) -> Option<FrameLease<M, FRAME, CAPACITY, WAITERS>> {
87        for index in 0..CAPACITY {
88            if self.slots[index]
89                .claimed
90                .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
91                .is_ok()
92            {
93                permit.disarm();
94                return Some(FrameLease {
95                    pool: self,
96                    index: index as u8,
97                    not_sync: PhantomData,
98                });
99            }
100        }
101        None
102    }
103}
104
105impl<M: RawMutex + 'static, const FRAME: usize, const CAPACITY: usize, const WAITERS: usize> Default
106    for SharedFramePool<M, FRAME, CAPACITY, WAITERS>
107{
108    fn default() -> Self {
109        Self::new()
110    }
111}
112
113#[must_use]
114pub struct FrameLease<
115    M: RawMutex + 'static,
116    const FRAME: usize,
117    const CAPACITY: usize,
118    const WAITERS: usize,
119> {
120    pool: &'static SharedFramePool<M, FRAME, CAPACITY, WAITERS>,
121    index: u8,
122    not_sync: PhantomData<Cell<()>>,
123}
124
125impl<M: RawMutex + 'static, const FRAME: usize, const CAPACITY: usize, const WAITERS: usize>
126    FrameLease<M, FRAME, CAPACITY, WAITERS>
127{
128    fn slot(&self) -> &FrameSlot<M, FRAME> {
129        &self.pool.slots[usize::from(self.index)]
130    }
131
132    pub async fn lock(&self) -> MutexGuard<'_, M, FrameBytes<u8, FRAME>> {
133        self.slot().frame.lock().await
134    }
135
136    pub async fn fill(&self, bytes: &[u8]) -> Result<(), FramePoolError> {
137        if bytes.len() > FRAME {
138            return Err(FramePoolError::FrameTooLarge {
139                len: bytes.len(),
140                capacity: FRAME,
141            });
142        }
143        let mut frame = self.lock().await;
144        frame.clear();
145        frame
146            .extend_from_slice(bytes)
147            .map_err(|_| FramePoolError::FrameTooLarge {
148                len: bytes.len(),
149                capacity: FRAME,
150            })
151    }
152
153    pub fn try_fill(&self, bytes: &[u8]) -> Result<(), FramePoolError> {
154        if bytes.len() > FRAME {
155            return Err(FramePoolError::FrameTooLarge {
156                len: bytes.len(),
157                capacity: FRAME,
158            });
159        }
160        let mut frame = self
161            .slot()
162            .frame
163            .try_lock()
164            .map_err(|_| FramePoolError::SlotBusy)?;
165        frame.clear();
166        frame
167            .extend_from_slice(bytes)
168            .map_err(|_| FramePoolError::FrameTooLarge {
169                len: bytes.len(),
170                capacity: FRAME,
171            })
172    }
173
174    pub async fn append(&self, bytes: &[u8]) -> Result<(), FramePoolError> {
175        let mut frame = self.lock().await;
176        let len = frame.len().saturating_add(bytes.len());
177        if len > FRAME {
178            return Err(FramePoolError::FrameTooLarge {
179                len,
180                capacity: FRAME,
181            });
182        }
183        frame
184            .extend_from_slice(bytes)
185            .map_err(|_| FramePoolError::FrameTooLarge {
186                len,
187                capacity: FRAME,
188            })
189    }
190
191    pub fn try_append(&self, bytes: &[u8]) -> Result<(), FramePoolError> {
192        let mut frame = self
193            .slot()
194            .frame
195            .try_lock()
196            .map_err(|_| FramePoolError::SlotBusy)?;
197        let len = frame.len().saturating_add(bytes.len());
198        if len > FRAME {
199            return Err(FramePoolError::FrameTooLarge {
200                len,
201                capacity: FRAME,
202            });
203        }
204        frame
205            .extend_from_slice(bytes)
206            .map_err(|_| FramePoolError::FrameTooLarge {
207                len,
208                capacity: FRAME,
209            })
210    }
211}
212
213impl<M: RawMutex + 'static, const FRAME: usize, const CAPACITY: usize, const WAITERS: usize> Drop
214    for FrameLease<M, FRAME, CAPACITY, WAITERS>
215{
216    fn drop(&mut self) {
217        self.slot().claimed.store(false, Ordering::Release);
218        self.pool.availability.release(1);
219    }
220}
221
222#[derive(Debug, PartialEq, Eq)]
223pub enum FramePoolError {
224    FrameTooLarge { len: usize, capacity: usize },
225    SlotBusy,
226    WaitQueueFull,
227    PermitWithoutAvailableSlot,
228}
229
230#[cfg(test)]
231mod tests {
232    use core::future::{poll_fn, ready, Future};
233    use core::pin::pin;
234    use core::task::Poll;
235
236    use embassy_futures::block_on;
237    use embassy_futures::select::{select, select3, Either, Either3};
238    use embassy_sync::blocking_mutex::raw::CriticalSectionRawMutex;
239
240    use super::{FramePoolError, SharedFramePool};
241
242    #[test]
243    fn leases_are_exclusive_and_release_on_drop() {
244        static POOL: SharedFramePool<CriticalSectionRawMutex, 8, 2, 2> = SharedFramePool::new();
245
246        let first = POOL.try_lease().ok().flatten();
247        let second = POOL.try_lease().ok().flatten();
248        assert!(first.is_some());
249        assert!(second.is_some());
250        assert_eq!(POOL.try_lease().map(|lease| lease.is_none()), Ok(true));
251
252        drop(first);
253        let replacement = POOL.try_lease().ok().flatten();
254        assert!(replacement.is_some());
255        assert_eq!(POOL.try_lease().map(|lease| lease.is_none()), Ok(true));
256    }
257
258    #[test]
259    fn frames_are_bounded_and_reused() {
260        static POOL: SharedFramePool<CriticalSectionRawMutex, 4, 1, 1> = SharedFramePool::new();
261
262        block_on(async {
263            let lease = POOL.lease().await;
264            assert!(lease.is_ok());
265            if let Ok(lease) = lease {
266                assert_eq!(lease.fill(b"prns").await, Ok(()));
267                assert_eq!(
268                    lease.append(b"!").await,
269                    Err(FramePoolError::FrameTooLarge {
270                        len: 5,
271                        capacity: 4,
272                    })
273                );
274                {
275                    let frame = lease.lock().await;
276                    assert_eq!(frame.as_slice(), b"prns");
277                    assert_eq!(lease.try_fill(b"rns"), Err(FramePoolError::SlotBusy));
278                }
279                assert_eq!(
280                    lease.fill(b"large").await,
281                    Err(FramePoolError::FrameTooLarge {
282                        len: 5,
283                        capacity: 4,
284                    })
285                );
286                drop(lease);
287
288                let reused = POOL.lease().await;
289                assert!(reused.is_ok());
290                if let Ok(reused) = reused {
291                    assert!(reused.lock().await.is_empty());
292                    assert_eq!(reused.fill(b"rns").await, Ok(()));
293                    assert_eq!(reused.try_append(b"!"), Ok(()));
294                    let frame = reused.lock().await;
295                    assert_eq!(frame.as_slice(), b"rns!");
296                }
297            }
298        });
299    }
300
301    #[test]
302    fn cancelled_lease_leaves_capacity_available() {
303        static POOL: SharedFramePool<CriticalSectionRawMutex, 4, 1, 2> = SharedFramePool::new();
304
305        let held = POOL.try_lease().ok().flatten();
306        assert!(held.is_some());
307        block_on(async {
308            assert!(matches!(
309                select(POOL.lease(), ready(())).await,
310                Either::Second(())
311            ));
312        });
313        drop(held);
314        assert_eq!(POOL.try_lease().map(|lease| lease.is_some()), Ok(true));
315    }
316
317    #[test]
318    fn wait_queue_saturation_is_explicit() {
319        static POOL: SharedFramePool<CriticalSectionRawMutex, 4, 1, 2> = SharedFramePool::new();
320
321        let held = POOL.try_lease().ok().flatten();
322        assert!(held.is_some());
323        block_on(async {
324            assert!(matches!(
325                select3(POOL.lease(), POOL.lease(), POOL.lease()).await,
326                Either3::Third(Err(FramePoolError::WaitQueueFull))
327            ));
328        });
329        drop(held);
330        assert_eq!(POOL.try_lease().map(|lease| lease.is_some()), Ok(true));
331    }
332
333    #[test]
334    fn waiters_receive_released_capacity_in_registration_order() {
335        static POOL: SharedFramePool<CriticalSectionRawMutex, 4, 1, 2> = SharedFramePool::new();
336
337        let held = POOL.try_lease().ok().flatten();
338        assert!(held.is_some());
339        block_on(async {
340            let mut first = pin!(POOL.lease());
341            let mut second = pin!(POOL.lease());
342            poll_fn(|cx| {
343                assert!(first.as_mut().poll(cx).is_pending());
344                assert!(second.as_mut().poll(cx).is_pending());
345                Poll::Ready(())
346            })
347            .await;
348            drop(held);
349
350            let winner = select(first.as_mut(), second.as_mut()).await;
351            assert!(matches!(&winner, Either::First(Ok(_))));
352            if let Either::First(Ok(first_lease)) = winner {
353                drop(first_lease);
354                assert!(second.await.is_ok());
355            }
356        });
357    }
358}