prns_interfaces_embassy/bluetooth_auto/
frame_pool.rs1use 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}