1use ::core::cell::Cell;
2
3use embassy_futures::join::join_array;
4use embassy_futures::select::{select, select5, select_array, Either, Either5};
5use embassy_sync::blocking_mutex::raw::{CriticalSectionRawMutex, RawMutex};
6use embassy_sync::blocking_mutex::CriticalSectionMutex;
7use embassy_sync::signal::Signal;
8use embassy_time::{with_deadline, with_timeout, Duration, Instant};
9use portable_atomic::{AtomicBool, AtomicU32, AtomicU64, AtomicU8, Ordering};
10
11use prns_core::engine::FanTarget;
12use prns_core::interfaces::bluetooth_auto::{
13 self as contract, BleAddress, BleIdentity, Control, Endpoint, EstablishedPeer,
14 EstablishedTransport, Handshake, HandshakeOutcome, L2capPlan, LinkCapabilities, LocalPeer,
15 PeerProtocol,
16};
17use prns_core::interfaces::bluetooth_auto::{
18 role_for, ConnectionPolicy, PolicyAction, PolicyInput,
19};
20use prns_core::interfaces::bluetooth_auto::{
21 AdvertisingMode, BleBackend, BleEvent, BleLink, BleSink, BleSource, DialOutcome, Origin,
22 RadioMode, ScanningMode,
23};
24use prns_core::interfaces::{
25 BitrateBps, ConnectionState, InterfaceId, InterfaceKind, InterfaceStatus,
26};
27use prns_runtime::manifold::grant::FrameTarget;
28use prns_runtime::runtime::{EmbassyFleet as Fleet, OutboundFrame};
29
30const DIAL_TRACK: usize = 6;
31
32const ACTION_CAP: usize = 6;
33
34const HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(10);
35const HANDSHAKE_LANES: usize = 2;
36const OUTBOUND_TIMEOUT: Duration = Duration::from_secs(2);
37
38const ACTION_OVERFLOW_REASON: &str = "BLE policy action capacity exceeded";
39const DIAL_INVARIANT_REASON: &str = "BLE dial admission invariant failed";
40const RADIO_CONTROL_REASON: &str = "BLE radio control failed";
41const INGRESS_PRESSURE_REASON: &str = "BLE receive pressure";
42const SETUP_FAILURE_REASON: &str = "BLE setup failed; retrying";
43const TRANSPORT_CLOSURE_REASON: &str = "BLE transport closed; retrying";
44
45#[derive(Debug, Clone, Copy, PartialEq, Eq)]
46pub enum BluetoothRecoveryReason {
47 IngressPressure,
48 SetupFailure,
49 TransportClosure,
50}
51
52#[derive(Debug, Clone, Copy, PartialEq, Eq)]
53pub struct BluetoothRecoveryCounters {
54 pub ingress_pressure: u32,
55 pub setup_failures: u32,
56 pub transport_closures: u32,
57}
58
59impl BluetoothRecoveryCounters {
60 const ZERO: Self = Self {
61 ingress_pressure: 0,
62 setup_failures: 0,
63 transport_closures: 0,
64 };
65}
66
67impl BluetoothRecoveryReason {
68 const fn as_u8(self) -> u8 {
69 match self {
70 Self::IngressPressure => 1,
71 Self::SetupFailure => 2,
72 Self::TransportClosure => 3,
73 }
74 }
75
76 fn from_u8(value: u8) -> Option<Self> {
77 match value {
78 1 => Some(Self::IngressPressure),
79 2 => Some(Self::SetupFailure),
80 3 => Some(Self::TransportClosure),
81 _ => None,
82 }
83 }
84
85 const fn description(self) -> &'static str {
86 match self {
87 Self::IngressPressure => INGRESS_PRESSURE_REASON,
88 Self::SetupFailure => SETUP_FAILURE_REASON,
89 Self::TransportClosure => TRANSPORT_CLOSURE_REASON,
90 }
91 }
92}
93
94pub struct BluetoothMemberStatus {
95 id: CriticalSectionMutex<Cell<InterfaceId>>,
96 connection: AtomicU8,
97 rx: AtomicU64,
98 tx: AtomicU64,
99 active: AtomicBool,
100}
101
102impl BluetoothMemberStatus {
103 const fn new() -> Self {
104 Self {
105 id: CriticalSectionMutex::new(Cell::new(InterfaceId::new([0u8; 8]))),
106 connection: AtomicU8::new(ConnectionState::Disconnected.as_u8()),
107 rx: AtomicU64::new(0),
108 tx: AtomicU64::new(0),
109 active: AtomicBool::new(false),
110 }
111 }
112
113 fn assign(&self, id: InterfaceId) {
114 self.id.lock(|cell| cell.set(id));
115 self.connection
116 .store(ConnectionState::Connected.as_u8(), Ordering::Relaxed);
117 self.rx.store(0, Ordering::Relaxed);
118 self.tx.store(0, Ordering::Relaxed);
119 self.active.store(true, Ordering::Relaxed);
120 }
121
122 fn retire(&self) {
123 self.connection
124 .store(ConnectionState::Disconnected.as_u8(), Ordering::Relaxed);
125 self.active.store(false, Ordering::Relaxed);
126 }
127
128 fn add_rx(&self, bytes: u64) {
129 self.rx.fetch_add(bytes, Ordering::Relaxed);
130 }
131
132 fn add_tx(&self, bytes: u64) {
133 self.tx.fetch_add(bytes, Ordering::Relaxed);
134 }
135}
136
137impl InterfaceStatus for BluetoothMemberStatus {
138 fn id(&self) -> InterfaceId {
139 self.id.lock(|cell| cell.get())
140 }
141
142 fn connection(&self) -> ConnectionState {
143 ConnectionState::from_u8(self.connection.load(Ordering::Relaxed))
144 }
145
146 fn rx_bytes(&self) -> u64 {
147 self.rx.load(Ordering::Relaxed)
148 }
149
150 fn tx_bytes(&self) -> u64 {
151 self.tx.load(Ordering::Relaxed)
152 }
153}
154
155pub struct BluetoothAutoShared<const MEMBERS: usize> {
156 id: InterfaceId,
157 enabled: AtomicBool,
158 enabled_changed: Signal<CriticalSectionRawMutex, bool>,
159 up: AtomicBool,
160 failed: AtomicBool,
161 fatal_failure_reason: CriticalSectionMutex<Cell<Option<&'static str>>>,
162 recovery_reason: AtomicU8,
163 peers: AtomicU32,
164 recovery_counters: CriticalSectionMutex<Cell<BluetoothRecoveryCounters>>,
165 members: [BluetoothMemberStatus; MEMBERS],
166}
167
168impl<const MEMBERS: usize> BluetoothAutoShared<MEMBERS> {
169 #[must_use]
170 pub const fn new(id: InterfaceId) -> Self {
171 Self {
172 id,
173 enabled: AtomicBool::new(true),
174 enabled_changed: Signal::new(),
175 up: AtomicBool::new(false),
176 failed: AtomicBool::new(false),
177 fatal_failure_reason: CriticalSectionMutex::new(Cell::new(None)),
178 recovery_reason: AtomicU8::new(0),
179 peers: AtomicU32::new(0),
180 recovery_counters: CriticalSectionMutex::new(Cell::new(
181 BluetoothRecoveryCounters::ZERO,
182 )),
183 members: [const { BluetoothMemberStatus::new() }; MEMBERS],
184 }
185 }
186}
187
188#[derive(Clone, Copy)]
189pub struct BluetoothAutoStatus<const MEMBERS: usize> {
190 shared: &'static BluetoothAutoShared<MEMBERS>,
191}
192
193impl<const MEMBERS: usize> BluetoothAutoStatus<MEMBERS> {
194 #[must_use]
195 pub const fn new(shared: &'static BluetoothAutoShared<MEMBERS>) -> Self {
196 Self { shared }
197 }
198
199 fn increment_recovery_counter(&self, reason: BluetoothRecoveryReason) {
200 self.shared.recovery_counters.lock(|slot| {
201 let mut counters = slot.get();
202 match reason {
203 BluetoothRecoveryReason::IngressPressure => {
204 counters.ingress_pressure = counters.ingress_pressure.saturating_add(1);
205 }
206 BluetoothRecoveryReason::SetupFailure => {
207 counters.setup_failures = counters.setup_failures.saturating_add(1);
208 }
209 BluetoothRecoveryReason::TransportClosure => {
210 counters.transport_closures = counters.transport_closures.saturating_add(1);
211 }
212 }
213 slot.set(counters);
214 });
215 }
216
217 pub fn note_ingress_pressure(&self) {
218 self.increment_recovery_counter(BluetoothRecoveryReason::IngressPressure);
219 self.shared.recovery_reason.store(
220 BluetoothRecoveryReason::IngressPressure.as_u8(),
221 Ordering::Relaxed,
222 );
223 }
224
225 pub fn note_successful_admission(&self) {
226 let _ = self.shared.recovery_reason.compare_exchange(
227 BluetoothRecoveryReason::IngressPressure.as_u8(),
228 0,
229 Ordering::Relaxed,
230 Ordering::Relaxed,
231 );
232 }
233
234 pub fn note_setup_failure(&self) {
235 self.increment_recovery_counter(BluetoothRecoveryReason::SetupFailure);
236 self.shared.recovery_reason.store(
237 BluetoothRecoveryReason::SetupFailure.as_u8(),
238 Ordering::Relaxed,
239 );
240 }
241
242 pub fn note_transport_closure(&self) {
243 self.increment_recovery_counter(BluetoothRecoveryReason::TransportClosure);
244 self.shared.recovery_reason.store(
245 BluetoothRecoveryReason::TransportClosure.as_u8(),
246 Ordering::Relaxed,
247 );
248 }
249
250 fn note_settled_link(&self) {
251 self.shared.recovery_reason.store(0, Ordering::Relaxed);
252 }
253
254 #[must_use]
255 pub fn recovery_reason(&self) -> Option<BluetoothRecoveryReason> {
256 BluetoothRecoveryReason::from_u8(self.shared.recovery_reason.load(Ordering::Relaxed))
257 }
258
259 #[must_use]
260 pub fn recovery_counters(&self) -> BluetoothRecoveryCounters {
261 self.shared.recovery_counters.lock(Cell::get)
262 }
263
264 #[must_use]
265 pub fn ingress_pressure_events(&self) -> u32 {
266 self.recovery_counters().ingress_pressure
267 }
268
269 #[must_use]
270 pub fn setup_failure_events(&self) -> u32 {
271 self.recovery_counters().setup_failures
272 }
273
274 #[must_use]
275 pub fn transport_closure_events(&self) -> u32 {
276 self.recovery_counters().transport_closures
277 }
278
279 fn mark_up(&self) {
280 self.shared.up.store(true, Ordering::Relaxed);
281 }
282
283 fn mark_failed(&self, reason: &'static str) {
284 self.shared
285 .fatal_failure_reason
286 .lock(|slot| slot.set(Some(reason)));
287 self.shared.failed.store(true, Ordering::Relaxed);
288 }
289
290 fn is_failed(&self) -> bool {
291 self.shared.failed.load(Ordering::Relaxed)
292 }
293
294 pub fn enable(&self) {
295 self.update_enabled(true);
296 }
297
298 pub fn disable(&self) {
299 self.update_enabled(false);
300 }
301
302 pub fn toggle_enabled(&self) {
303 let enabled = !self.shared.enabled.fetch_xor(true, Ordering::Relaxed);
304 self.shared.enabled_changed.signal(enabled);
305 }
306
307 fn update_enabled(&self, enabled: bool) {
308 if self.shared.enabled.swap(enabled, Ordering::Relaxed) != enabled {
309 self.shared.enabled_changed.signal(enabled);
310 }
311 }
312
313 #[must_use]
314 pub fn is_enabled(&self) -> bool {
315 self.shared.enabled.load(Ordering::Relaxed)
316 }
317
318 async fn wait_until_enabled(&self) {
319 self.wait_for_enabled_state(true).await;
320 }
321
322 async fn wait_until_disabled(&self) {
323 self.wait_for_enabled_state(false).await;
324 }
325
326 async fn wait_for_enabled_state(&self, enabled: bool) {
327 loop {
328 if self.is_enabled() == enabled {
329 return;
330 }
331 if self.shared.enabled_changed.wait().await == enabled {
332 return;
333 }
334 }
335 }
336
337 fn member(&self, slot: usize) -> &'static BluetoothMemberStatus {
338 &self.shared.members[slot]
339 }
340
341 fn republish_peer_count(&self) {
342 let count = self
343 .shared
344 .members
345 .iter()
346 .filter(|member| member.active.load(Ordering::Relaxed))
347 .count();
348 self.shared.peers.store(count as u32, Ordering::Relaxed);
349 }
350
351 pub fn members(&self) -> impl Iterator<Item = &'static BluetoothMemberStatus> {
352 self.shared
353 .members
354 .iter()
355 .filter(|member| member.active.load(Ordering::Relaxed))
356 }
357}
358
359impl<const MEMBERS: usize> InterfaceStatus for BluetoothAutoStatus<MEMBERS> {
360 fn id(&self) -> InterfaceId {
361 self.shared.id
362 }
363
364 fn connection(&self) -> ConnectionState {
365 if !self.is_enabled() {
366 ConnectionState::Disabled
367 } else if self.is_failed() {
368 ConnectionState::Failed
369 } else if !self.shared.up.load(Ordering::Relaxed) {
370 ConnectionState::Initializing
371 } else if matches!(
372 self.recovery_reason(),
373 Some(BluetoothRecoveryReason::IngressPressure)
374 ) && self.shared.peers.load(Ordering::Relaxed) > 0
375 {
376 ConnectionState::Degraded
377 } else if self.recovery_reason().is_some() && self.shared.peers.load(Ordering::Relaxed) == 0
378 {
379 ConnectionState::Reconnecting
380 } else if self.shared.peers.load(Ordering::Relaxed) > 0 {
381 ConnectionState::Connected
382 } else {
383 ConnectionState::Disconnected
384 }
385 }
386
387 fn rx_bytes(&self) -> u64 {
388 self.shared
389 .members
390 .iter()
391 .map(|member| member.rx.load(Ordering::Relaxed))
392 .sum()
393 }
394
395 fn tx_bytes(&self) -> u64 {
396 self.shared
397 .members
398 .iter()
399 .map(|member| member.tx.load(Ordering::Relaxed))
400 .sum()
401 }
402
403 fn failure_reason(&self) -> Option<&'static str> {
404 self.shared
405 .fatal_failure_reason
406 .lock(Cell::get)
407 .or_else(|| {
408 self.recovery_reason()
409 .map(BluetoothRecoveryReason::description)
410 })
411 }
412}
413
414struct Active<L: BleLink> {
415 identity: BleIdentity,
416 id: InterfaceId,
417 slot: usize,
418 address: BleAddress,
419 source: L::Source,
420 sink: L::Sink,
421}
422
423struct PendingActions<const CAP: usize> {
424 actions: heapless::Vec<PolicyAction, CAP>,
425 overflowed: bool,
426}
427
428impl<const CAP: usize> PendingActions<CAP> {
429 const fn new() -> Self {
430 Self {
431 actions: heapless::Vec::new(),
432 overflowed: false,
433 }
434 }
435
436 fn push(&mut self, action: PolicyAction) {
437 if self.actions.push(action).is_err() {
438 self.overflowed = true;
439 }
440 }
441
442 fn take(&mut self) -> heapless::Vec<PolicyAction, CAP> {
443 ::core::mem::take(&mut self.actions)
444 }
445
446 fn clear(&mut self) {
447 self.actions.clear();
448 self.overflowed = false;
449 }
450}
451
452enum HandshakeStage {
453 ColumbaReceive,
454 ColumbaSend {
455 identity: BleIdentity,
456 },
457 NativeSend {
458 handshake: Option<Handshake>,
459 control: Control,
460 },
461 NativeReceive {
462 handshake: Option<Handshake>,
463 },
464 NativeReply {
465 handshake: Option<Handshake>,
466 control: Control,
467 outcome: HandshakeOutcome,
468 },
469}
470
471struct PendingHandshake<L: BleLink> {
472 link: L,
473 address: BleAddress,
474 origin: Origin,
475 deadline: Instant,
476 stage: HandshakeStage,
477}
478
479impl<L: BleLink> PendingHandshake<L> {
480 fn new(link: L, origin: Origin, local: LocalPeer) -> Self {
481 let stage = if link.peer_protocol() == PeerProtocol::Columba {
482 HandshakeStage::ColumbaReceive
483 } else {
484 let (handshake, opening) = Handshake::begin(role_for(origin), local, None);
485 match opening {
486 Some(control) => HandshakeStage::NativeSend {
487 handshake: Some(handshake),
488 control,
489 },
490 None => HandshakeStage::NativeReceive {
491 handshake: Some(handshake),
492 },
493 }
494 };
495 Self {
496 address: link.address(),
497 link,
498 origin,
499 deadline: Instant::now() + HANDSHAKE_TIMEOUT,
500 stage,
501 }
502 }
503}
504
505#[derive(Debug, Clone, Copy, PartialEq, Eq)]
506enum HandshakeFailure {
507 Timeout,
508 Link,
509 Aborted,
510 InvariantViolation,
511}
512
513struct HandshakeDone<L: BleLink> {
514 address: BleAddress,
515 origin: Origin,
516 outcome: Result<(EstablishedPeer, L), HandshakeFailure>,
517}
518
519enum HandshakeStep<L: BleLink> {
520 Advanced,
521 Done(HandshakeDone<L>),
522}
523
524#[derive(Clone, Copy, PartialEq, Eq)]
525enum SendState {
526 NotSelected,
527 Pending,
528 Sent,
529 Failed,
530}
531
532enum SupervisorStep<L: BleLink> {
533 Disabled,
534 Handshake(HandshakeStep<L>),
535 Backend(BleEvent<L>),
536 Inbound(usize, Result<usize, <L::Source as BleSource>::Error>),
537 Outbound,
538}
539
540pub struct BluetoothAuto<B, const MEMBERS: usize> {
541 backend: B,
542 local: LocalPeer,
543 status: BluetoothAutoStatus<MEMBERS>,
544 bitrate: BitrateBps,
545}
546
547impl<B, const MEMBERS: usize> BluetoothAuto<B, MEMBERS>
548where
549 B: BleBackend<MEMBERS>,
550{
551 #[must_use]
552 pub fn new(
553 backend: B,
554 identity: BleIdentity,
555 endpoint: Endpoint,
556 capabilities: LinkCapabilities,
557 shared: &'static BluetoothAutoShared<MEMBERS>,
558 ) -> Self {
559 Self {
560 backend,
561 local: LocalPeer {
562 identity,
563 endpoint,
564 capabilities,
565 },
566 status: BluetoothAutoStatus::new(shared),
567 bitrate: contract::BLE_BITRATE_GUESS_BPS,
568 }
569 }
570
571 #[must_use]
572 pub fn status(&self) -> BluetoothAutoStatus<MEMBERS> {
573 self.status
574 }
575
576 pub async fn run<M, const FRAME: usize, const NOTIFY: usize, const LIFECYCLE: usize>(
577 self,
578 mut fleet: Fleet<M, FRAME, NOTIFY, LIFECYCLE>,
579 ) where
580 M: RawMutex + 'static,
581 {
582 let Self {
583 mut backend,
584 local: configured_local,
585 status,
586 bitrate,
587 } = self;
588 if let Some(reason) = backend.blocked() {
589 status.mark_failed(reason);
590 ::core::future::pending::<()>().await;
591 return;
592 }
593 let configured_capabilities = configured_local.capabilities;
594 let mut local = configured_local;
595 prepare_radio(&mut backend, &mut local, configured_capabilities, &status).await;
596 if status.is_failed() {
597 let _ = backend.set_radio_mode(RadioMode::Off).await;
598 ::core::future::pending::<()>().await;
599 return;
600 }
601 let mut manager = ConnectionPolicy::<MEMBERS, DIAL_TRACK>::new(local);
602 let mut members: [Option<Active<B::Link>>; MEMBERS] = [const { None }; MEMBERS];
603 let mut inbufs: [[u8; contract::BLE_HW_MTU]; MEMBERS] =
604 [[0u8; contract::BLE_HW_MTU]; MEMBERS];
605 let mut handshakes: [Option<PendingHandshake<B::Link>>; HANDSHAKE_LANES] =
606 [const { None }; HANDSHAKE_LANES];
607 let mut pending = PendingActions::<ACTION_CAP>::new();
608 let mut outbound_first = false;
609 status.mark_up();
610 manager.start(&mut |action| pending.push(action));
611 apply_radio(
612 &mut pending,
613 &mut manager,
614 &status,
615 &mut fleet,
616 &mut backend,
617 &mut members,
618 )
619 .await;
620
621 loop {
622 if status.is_failed() {
623 handshakes.fill_with(|| None);
624 pending.clear();
625 disable_members(&status, &mut fleet, &mut backend, &mut members).await;
626 ::core::future::pending::<()>().await;
627 return;
628 }
629 if !status.is_enabled() {
630 handshakes.fill_with(|| None);
631 disable_members(&status, &mut fleet, &mut backend, &mut members).await;
632 pending.clear();
633 status.wait_until_enabled().await;
634 local = configured_local;
635 prepare_radio(&mut backend, &mut local, configured_capabilities, &status).await;
636 if status.is_failed() {
637 continue;
638 }
639 manager = ConnectionPolicy::<MEMBERS, DIAL_TRACK>::new(local);
640 manager.start(&mut |action| pending.push(action));
641 apply_radio(
642 &mut pending,
643 &mut manager,
644 &status,
645 &mut fleet,
646 &mut backend,
647 &mut members,
648 )
649 .await;
650 continue;
651 }
652 let step = next_step(
653 &status,
654 &mut backend,
655 &mut handshakes,
656 local,
657 &mut members,
658 &mut inbufs,
659 &fleet,
660 outbound_first,
661 )
662 .await;
663 outbound_first = !matches!(&step, SupervisorStep::Outbound);
664 let now_ms = Instant::now().as_millis();
665 match step {
666 SupervisorStep::Disabled => {}
667 SupervisorStep::Handshake(HandshakeStep::Advanced) => {}
668 SupervisorStep::Handshake(HandshakeStep::Done(HandshakeDone {
669 address,
670 origin,
671 outcome,
672 })) => match outcome {
673 Ok((established, link)) => {
674 manager.handle(
675 PolicyInput::Settled {
676 address,
677 origin,
678 established,
679 now_ms,
680 },
681 &mut |action| pending.push(action),
682 );
683 apply_settled(
684 link,
685 bitrate,
686 &mut manager,
687 &mut pending,
688 &status,
689 &mut fleet,
690 &mut backend,
691 &mut members,
692 )
693 .await;
694 }
695 Err(_) => {
696 status.note_setup_failure();
697 manager.handle(
698 PolicyInput::HandshakeFailed { address, origin },
699 &mut |action| pending.push(action),
700 );
701 apply_radio(
702 &mut pending,
703 &mut manager,
704 &status,
705 &mut fleet,
706 &mut backend,
707 &mut members,
708 )
709 .await;
710 }
711 },
712 SupervisorStep::Backend(BleEvent::Sighting { address, .. }) => {
713 manager.handle(PolicyInput::Sighting { address, now_ms }, &mut |action| {
714 pending.push(action)
715 });
716 apply_radio(
717 &mut pending,
718 &mut manager,
719 &status,
720 &mut fleet,
721 &mut backend,
722 &mut members,
723 )
724 .await;
725 }
726 SupervisorStep::Backend(BleEvent::Inbound(link)) => {
727 queue_handshake(
728 link,
729 Origin::Accepted,
730 local,
731 &mut manager,
732 &mut handshakes,
733 &mut backend,
734 )
735 .await;
736 }
737 SupervisorStep::Backend(BleEvent::LinkReady { link, origin, .. }) => {
738 queue_handshake(
739 link,
740 origin,
741 local,
742 &mut manager,
743 &mut handshakes,
744 &mut backend,
745 )
746 .await;
747 }
748 SupervisorStep::Backend(BleEvent::DialFailed { address }) => {
749 status.note_setup_failure();
750 manager.handle(PolicyInput::DialFailed { address, now_ms }, &mut |action| {
751 pending.push(action)
752 });
753 apply_radio(
754 &mut pending,
755 &mut manager,
756 &status,
757 &mut fleet,
758 &mut backend,
759 &mut members,
760 )
761 .await;
762 }
763 SupervisorStep::Inbound(index, received) => {
764 deliver_inbound(
765 index,
766 received.map_err(|_| ()),
767 &mut manager,
768 &mut pending,
769 &status,
770 &mut fleet,
771 &mut backend,
772 &mut members,
773 &mut inbufs,
774 )
775 .await;
776 }
777 SupervisorStep::Outbound => {
778 while let Some(frame) = fleet.try_next_outbound() {
783 send_outbound(
784 &frame,
785 &mut manager,
786 &mut pending,
787 &status,
788 &mut fleet,
789 &mut backend,
790 &mut members,
791 )
792 .await;
793 }
794 }
795 }
796 }
797 }
798}
799
800#[expect(
801 clippy::too_many_arguments,
802 reason = "the scheduler owns one borrow per independently wakeable supervisor branch"
803)]
804async fn next_step<
805 B,
806 M: RawMutex + 'static,
807 const FRAME: usize,
808 const NOTIFY: usize,
809 const LIFECYCLE: usize,
810 const MEMBERS: usize,
811>(
812 status: &BluetoothAutoStatus<MEMBERS>,
813 backend: &mut B,
814 handshakes: &mut [Option<PendingHandshake<B::Link>>; HANDSHAKE_LANES],
815 local: LocalPeer,
816 members: &mut [Option<Active<B::Link>>; MEMBERS],
817 inbufs: &mut [[u8; contract::BLE_HW_MTU]; MEMBERS],
818 fleet: &Fleet<M, FRAME, NOTIFY, LIFECYCLE>,
819 outbound_first: bool,
820) -> SupervisorStep<B::Link>
821where
822 B: BleBackend<MEMBERS>,
823{
824 if outbound_first {
825 return match select5(
826 status.wait_until_disabled(),
827 fleet.outbound_ready(),
828 advance_handshakes(handshakes, local),
829 backend.next_event(),
830 recv_any(members, inbufs),
831 )
832 .await
833 {
834 Either5::First(()) => SupervisorStep::Disabled,
835 Either5::Second(()) => SupervisorStep::Outbound,
836 Either5::Third(step) => SupervisorStep::Handshake(step),
837 Either5::Fourth(event) => SupervisorStep::Backend(event),
838 Either5::Fifth((index, received)) => SupervisorStep::Inbound(index, received),
839 };
840 }
841 match select5(
842 status.wait_until_disabled(),
843 advance_handshakes(handshakes, local),
844 backend.next_event(),
845 recv_any(members, inbufs),
846 fleet.outbound_ready(),
847 )
848 .await
849 {
850 Either5::First(()) => SupervisorStep::Disabled,
851 Either5::Second(step) => SupervisorStep::Handshake(step),
852 Either5::Third(event) => SupervisorStep::Backend(event),
853 Either5::Fourth((index, received)) => SupervisorStep::Inbound(index, received),
854 Either5::Fifth(()) => SupervisorStep::Outbound,
855 }
856}
857
858async fn prepare_radio<B, const MEMBERS: usize>(
859 backend: &mut B,
860 local: &mut LocalPeer,
861 configured_capabilities: LinkCapabilities,
862 status: &BluetoothAutoStatus<MEMBERS>,
863) where
864 B: BleBackend<MEMBERS>,
865{
866 if backend.set_radio_mode(RadioMode::On).await.is_err() {
867 status.mark_failed(RADIO_CONTROL_REASON);
868 return;
869 }
870 match backend.local_capabilities(configured_capabilities).await {
871 Ok(capabilities) => local.capabilities = capabilities,
872 Err(_) => status.mark_failed(RADIO_CONTROL_REASON),
873 }
874}
875
876async fn queue_handshake<B, const MEMBERS: usize>(
877 link: B::Link,
878 origin: Origin,
879 local: LocalPeer,
880 manager: &mut ConnectionPolicy<MEMBERS, DIAL_TRACK>,
881 handshakes: &mut [Option<PendingHandshake<B::Link>>; HANDSHAKE_LANES],
882 backend: &mut B,
883) where
884 B: BleBackend<MEMBERS>,
885{
886 let address = link.address();
887 match handshakes.iter_mut().find(|entry| entry.is_none()) {
888 Some(entry) if manager.begin_handshake(origin) => {
889 *entry = Some(PendingHandshake::new(link, origin, local));
890 }
891 _ => {
892 drop(link);
893 backend.on_link_closed(address).await;
894 }
895 }
896}
897
898async fn advance_handshakes<L: BleLink>(
899 handshakes: &mut [Option<PendingHandshake<L>>; HANDSHAKE_LANES],
900 local: LocalPeer,
901) -> HandshakeStep<L> {
902 let [first, second] = handshakes;
903 match select(
904 advance_handshake(first, local),
905 advance_handshake(second, local),
906 )
907 .await
908 {
909 Either::First(step) | Either::Second(step) => step,
910 }
911}
912
913async fn advance_handshake<L: BleLink>(
914 pending: &mut Option<PendingHandshake<L>>,
915 local: LocalPeer,
916) -> HandshakeStep<L> {
917 let completion = match pending.as_mut() {
918 Some(pending) => {
919 let deadline = pending.deadline;
920 match with_deadline(deadline, async {
921 let mut next_stage = None;
922 let completion = match &mut pending.stage {
923 HandshakeStage::ColumbaReceive => {
924 match pending.link.receive_columba_peer_identity().await {
925 Ok(identity) if pending.origin == Origin::Dialed => {
926 next_stage = Some(HandshakeStage::ColumbaSend { identity });
927 None
928 }
929 Ok(identity) => Some(Ok(EstablishedPeer {
930 identity,
931 transport: EstablishedTransport::ColumbaGatt,
932 peer_rssi: None,
933 })),
934 Err(_) => Some(Err(HandshakeFailure::Link)),
935 }
936 }
937 HandshakeStage::ColumbaSend { identity } => {
938 let identity = *identity;
939 match pending.link.send_columba_identity(local.identity).await {
940 Ok(()) => Some(Ok(EstablishedPeer {
941 identity,
942 transport: EstablishedTransport::ColumbaGatt,
943 peer_rssi: None,
944 })),
945 Err(_) => Some(Err(HandshakeFailure::Link)),
946 }
947 }
948 HandshakeStage::NativeSend { handshake, control } => {
949 let control = *control;
950 match pending.link.control_send(&control).await {
951 Ok(()) => match handshake.take() {
952 Some(handshake) => {
953 next_stage = Some(HandshakeStage::NativeReceive {
954 handshake: Some(handshake),
955 });
956 None
957 }
958 None => Some(Err(HandshakeFailure::InvariantViolation)),
959 },
960 Err(_) => Some(Err(HandshakeFailure::Link)),
961 }
962 }
963 HandshakeStage::NativeReceive { handshake } => {
964 match pending.link.control_recv().await {
965 Ok(control) => match handshake.as_mut() {
966 Some(active) => {
967 let reaction = active.absorb(control);
968 match reaction.reply {
969 Some(control) => match handshake.take() {
970 Some(handshake) => {
971 next_stage = Some(HandshakeStage::NativeReply {
972 handshake: Some(handshake),
973 control,
974 outcome: reaction.outcome,
975 });
976 None
977 }
978 None => Some(Err(HandshakeFailure::InvariantViolation)),
979 },
980 None => match reaction.outcome {
981 HandshakeOutcome::Pending => None,
982 HandshakeOutcome::Settled(established) => {
983 Some(Ok(established))
984 }
985 HandshakeOutcome::Aborted(_) => {
986 Some(Err(HandshakeFailure::Aborted))
987 }
988 },
989 }
990 }
991 None => Some(Err(HandshakeFailure::InvariantViolation)),
992 },
993 Err(_) => Some(Err(HandshakeFailure::Link)),
994 }
995 }
996 HandshakeStage::NativeReply {
997 handshake,
998 control,
999 outcome,
1000 } => {
1001 let control = *control;
1002 let outcome = *outcome;
1003 match pending.link.control_send(&control).await {
1004 Ok(()) => match outcome {
1005 HandshakeOutcome::Pending => match handshake.take() {
1006 Some(handshake) => {
1007 next_stage = Some(HandshakeStage::NativeReceive {
1008 handshake: Some(handshake),
1009 });
1010 None
1011 }
1012 None => Some(Err(HandshakeFailure::InvariantViolation)),
1013 },
1014 HandshakeOutcome::Settled(established) => Some(Ok(established)),
1015 HandshakeOutcome::Aborted(_) => {
1016 Some(Err(HandshakeFailure::Aborted))
1017 }
1018 },
1019 Err(_) => Some(Err(HandshakeFailure::Link)),
1020 }
1021 }
1022 };
1023 if let Some(stage) = next_stage {
1024 pending.stage = stage;
1025 }
1026 completion
1027 })
1028 .await
1029 {
1030 Ok(completion) => completion,
1031 Err(_) => Some(Err(HandshakeFailure::Timeout)),
1032 }
1033 }
1034 None => ::core::future::pending().await,
1035 };
1036 let Some(outcome) = completion else {
1037 return HandshakeStep::Advanced;
1038 };
1039 let Some(pending) = pending.take() else {
1040 return HandshakeStep::Advanced;
1041 };
1042 HandshakeStep::Done(HandshakeDone {
1043 address: pending.address,
1044 origin: pending.origin,
1045 outcome: outcome.map(|established| (established, pending.link)),
1046 })
1047}
1048
1049async fn recv_or_pending<L: BleLink>(
1050 member: &mut Option<Active<L>>,
1051 buf: &mut [u8; contract::BLE_HW_MTU],
1052) -> Result<usize, <L::Source as BleSource>::Error> {
1053 match member {
1054 Some(active) => active.source.recv_frame(buf).await,
1055 None => ::core::future::pending().await,
1056 }
1057}
1058
1059#[expect(
1060 clippy::expect_used,
1061 reason = "from_fn runs exactly MEMBERS times over a zip of two [_; MEMBERS] arrays, so the iterator cannot run dry; the disjoint-borrow trick has no panic-free spelling without unsafe"
1062)]
1063async fn recv_any<L: BleLink, const MEMBERS: usize>(
1064 members: &mut [Option<Active<L>>; MEMBERS],
1065 bufs: &mut [[u8; contract::BLE_HW_MTU]; MEMBERS],
1066) -> (usize, Result<usize, <L::Source as BleSource>::Error>) {
1067 let mut pairs = members.iter_mut().zip(bufs.iter_mut());
1068 let futures: [_; MEMBERS] = ::core::array::from_fn(|_| {
1069 let (member, buf) = pairs.next().expect("one pair per member slot");
1070 recv_or_pending(member, buf)
1071 });
1072 let (result, index) = select_array(futures).await;
1073 (index, result)
1074}
1075
1076async fn apply_one<
1077 B,
1078 M: RawMutex + 'static,
1079 const FRAME: usize,
1080 const NOTIFY: usize,
1081 const LIFECYCLE: usize,
1082 const MEMBERS: usize,
1083>(
1084 action: PolicyAction,
1085 pending: &mut PendingActions<ACTION_CAP>,
1086 manager: &mut ConnectionPolicy<MEMBERS, DIAL_TRACK>,
1087 status: &BluetoothAutoStatus<MEMBERS>,
1088 fleet: &mut Fleet<M, FRAME, NOTIFY, LIFECYCLE>,
1089 backend: &mut B,
1090 members: &mut [Option<Active<B::Link>>; MEMBERS],
1091) where
1092 B: BleBackend<MEMBERS>,
1093{
1094 match action {
1095 PolicyAction::Dial(address) => match backend.dial(address).await {
1096 DialOutcome::Started => {}
1097 DialOutcome::Busy | DialOutcome::UnknownPeer | DialOutcome::RadioOff => {
1098 manager.handle(
1099 PolicyInput::DialFailed {
1100 address,
1101 now_ms: Instant::now().as_millis(),
1102 },
1103 &mut |action| pending.push(action),
1104 );
1105 }
1106 DialOutcome::InvariantViolation => status.mark_failed(DIAL_INVARIANT_REASON),
1107 },
1108 PolicyAction::Evict { slot, .. } => {
1109 if let Some(member) = members[slot].take() {
1110 fleet.deregister_member(member.id).await;
1111 status.member(slot).retire();
1112 status.republish_peer_count();
1113 backend.on_link_closed(member.address).await;
1114 }
1115 }
1116 PolicyAction::NotifyClosed(address) => backend.on_link_closed(address).await,
1117 PolicyAction::SetAdvertising(mode) => {
1118 if backend.set_advertising(mode).await.is_err() {
1119 status.mark_failed(RADIO_CONTROL_REASON);
1120 }
1121 }
1122 PolicyAction::SetScanning(mode) => {
1123 if backend.set_scanning(mode).await.is_err() {
1124 status.mark_failed(RADIO_CONTROL_REASON);
1125 }
1126 }
1127 PolicyAction::Admit { .. } | PolicyAction::Reject { .. } => {}
1128 }
1129}
1130
1131async fn apply_radio<
1132 B,
1133 M: RawMutex + 'static,
1134 const FRAME: usize,
1135 const NOTIFY: usize,
1136 const LIFECYCLE: usize,
1137 const MEMBERS: usize,
1138>(
1139 pending: &mut PendingActions<ACTION_CAP>,
1140 manager: &mut ConnectionPolicy<MEMBERS, DIAL_TRACK>,
1141 status: &BluetoothAutoStatus<MEMBERS>,
1142 fleet: &mut Fleet<M, FRAME, NOTIFY, LIFECYCLE>,
1143 backend: &mut B,
1144 members: &mut [Option<Active<B::Link>>; MEMBERS],
1145) where
1146 B: BleBackend<MEMBERS>,
1147{
1148 loop {
1149 if pending.overflowed {
1150 status.mark_failed(ACTION_OVERFLOW_REASON);
1151 pending.actions.clear();
1152 return;
1153 }
1154 let actions = pending.take();
1155 if actions.is_empty() {
1156 return;
1157 }
1158 for action in actions {
1159 apply_one(action, pending, manager, status, fleet, backend, members).await;
1160 if status.is_failed() {
1161 return;
1162 }
1163 }
1164 }
1165}
1166
1167async fn disable_members<
1168 B,
1169 M: RawMutex + 'static,
1170 const FRAME: usize,
1171 const NOTIFY: usize,
1172 const LIFECYCLE: usize,
1173 const MEMBERS: usize,
1174>(
1175 status: &BluetoothAutoStatus<MEMBERS>,
1176 fleet: &mut Fleet<M, FRAME, NOTIFY, LIFECYCLE>,
1177 backend: &mut B,
1178 members: &mut [Option<Active<B::Link>>; MEMBERS],
1179) where
1180 B: BleBackend<MEMBERS>,
1181{
1182 let advertising = backend.set_advertising(AdvertisingMode::Off).await;
1183 let scanning = backend.set_scanning(ScanningMode::Off).await;
1184 let radio = backend.set_radio_mode(RadioMode::Off).await;
1185 if advertising.is_err() || scanning.is_err() || radio.is_err() {
1186 status.mark_failed(RADIO_CONTROL_REASON);
1187 }
1188 let mut changed = false;
1189 for (slot, entry) in members.iter_mut().enumerate() {
1190 if let Some(id) = entry.as_ref().map(|member| member.id) {
1191 fleet.deregister_member(id).await;
1192 let Some(member) = entry.take() else {
1193 continue;
1194 };
1195 status.member(slot).retire();
1196 backend.on_link_closed(member.address).await;
1197 changed = true;
1198 }
1199 }
1200 if changed {
1201 status.republish_peer_count();
1202 }
1203}
1204
1205async fn close_member<
1206 B,
1207 M: RawMutex + 'static,
1208 const FRAME: usize,
1209 const NOTIFY: usize,
1210 const LIFECYCLE: usize,
1211 const MEMBERS: usize,
1212>(
1213 slot: usize,
1214 manager: &mut ConnectionPolicy<MEMBERS, DIAL_TRACK>,
1215 pending: &mut PendingActions<ACTION_CAP>,
1216 status: &BluetoothAutoStatus<MEMBERS>,
1217 fleet: &mut Fleet<M, FRAME, NOTIFY, LIFECYCLE>,
1218 backend: &mut B,
1219 members: &mut [Option<Active<B::Link>>; MEMBERS],
1220) where
1221 B: BleBackend<MEMBERS>,
1222{
1223 let Some(id) = members[slot].as_ref().map(|member| member.id) else {
1224 return;
1225 };
1226 fleet.deregister_member(id).await;
1227 let Some(member) = members[slot].take() else {
1228 return;
1229 };
1230 status.member(slot).retire();
1231 status.republish_peer_count();
1232 manager.handle(
1233 PolicyInput::Closed {
1234 identity: member.identity,
1235 address: member.address,
1236 },
1237 &mut |action| pending.push(action),
1238 );
1239 apply_radio(pending, manager, status, fleet, backend, members).await;
1240}
1241
1242fn selected<L: BleLink>(member: &Active<L>, target: FrameTarget) -> bool {
1243 match target {
1244 FrameTarget::Direct(id) => member.id == id,
1245 FrameTarget::Fan(FanTarget::Only(id)) => member.id == id,
1246 FrameTarget::Fan(FanTarget::All) => true,
1247 FrameTarget::Fan(FanTarget::AllExcept(id)) => member.id != id,
1248 }
1249}
1250
1251async fn send_member<L: BleLink>(
1252 member: &mut Option<Active<L>>,
1253 state: &mut SendState,
1254 frame: &[u8],
1255) {
1256 if *state != SendState::Pending {
1257 return;
1258 }
1259 let Some(member) = member.as_mut() else {
1260 *state = SendState::Failed;
1261 return;
1262 };
1263 *state = if member.sink.send_frame(frame).await.is_ok() {
1264 SendState::Sent
1265 } else {
1266 SendState::Failed
1267 };
1268}
1269
1270#[expect(
1271 clippy::expect_used,
1272 reason = "from_fn runs exactly MEMBERS times over a zip of two [_; MEMBERS] arrays, so the iterator cannot run dry; the disjoint-borrow trick has no panic-free spelling without unsafe"
1273)]
1274async fn send_members<L: BleLink, const MEMBERS: usize>(
1275 members: &mut [Option<Active<L>>; MEMBERS],
1276 states: &mut [SendState; MEMBERS],
1277 frame: &[u8],
1278) {
1279 let mut pairs = members.iter_mut().zip(states.iter_mut());
1280 let futures: [_; MEMBERS] = ::core::array::from_fn(|_| {
1281 let (member, state) = pairs.next().expect("one pair per member slot");
1282 send_member(member, state, frame)
1283 });
1284 join_array(futures).await;
1285}
1286
1287async fn send_outbound<
1288 B,
1289 M: RawMutex + 'static,
1290 const FRAME: usize,
1291 const NOTIFY: usize,
1292 const LIFECYCLE: usize,
1293 const MEMBERS: usize,
1294>(
1295 frame: &OutboundFrame<FRAME>,
1296 manager: &mut ConnectionPolicy<MEMBERS, DIAL_TRACK>,
1297 pending: &mut PendingActions<ACTION_CAP>,
1298 status: &BluetoothAutoStatus<MEMBERS>,
1299 fleet: &mut Fleet<M, FRAME, NOTIFY, LIFECYCLE>,
1300 backend: &mut B,
1301 members: &mut [Option<Active<B::Link>>; MEMBERS],
1302) where
1303 B: BleBackend<MEMBERS>,
1304{
1305 if frame.is_empty() {
1306 return;
1307 }
1308 let mut states = ::core::array::from_fn(|slot| match members[slot].as_ref() {
1309 Some(member) if selected(member, frame.target()) => SendState::Pending,
1310 _ => SendState::NotSelected,
1311 });
1312 let sends = send_members(members, &mut states, frame.bytes());
1313 match select(
1314 status.wait_until_disabled(),
1315 with_timeout(OUTBOUND_TIMEOUT, sends),
1316 )
1317 .await
1318 {
1319 Either::First(()) => return,
1320 Either::Second(_) => {}
1321 }
1322 for (slot, state) in states.into_iter().enumerate() {
1323 match state {
1324 SendState::NotSelected => {}
1325 SendState::Sent => status.member(slot).add_tx(frame.len() as u64),
1326 SendState::Pending | SendState::Failed => {
1327 status.note_transport_closure();
1328 close_member(slot, manager, pending, status, fleet, backend, members).await;
1329 }
1330 }
1331 }
1332}
1333
1334#[expect(
1335 clippy::too_many_arguments,
1336 reason = "embedded serve-loop internals pass the loop's split-borrowed locals; bundling awaits an on-hardware validation pass"
1337)]
1338async fn deliver_inbound<
1339 B,
1340 M: RawMutex + 'static,
1341 const FRAME: usize,
1342 const NOTIFY: usize,
1343 const LIFECYCLE: usize,
1344 const MEMBERS: usize,
1345>(
1346 index: usize,
1347 received: Result<usize, ()>,
1348 manager: &mut ConnectionPolicy<MEMBERS, DIAL_TRACK>,
1349 pending: &mut PendingActions<ACTION_CAP>,
1350 status: &BluetoothAutoStatus<MEMBERS>,
1351 fleet: &mut Fleet<M, FRAME, NOTIFY, LIFECYCLE>,
1352 backend: &mut B,
1353 members: &mut [Option<Active<B::Link>>; MEMBERS],
1354 inbufs: &mut [[u8; contract::BLE_HW_MTU]; MEMBERS],
1355) where
1356 B: BleBackend<MEMBERS>,
1357{
1358 match received {
1359 Ok(0) => {}
1360 Ok(len) => {
1361 if let Some(member) = members[index].as_ref() {
1362 if fleet
1363 .deliver_inbound(member.id, &inbufs[index][..len])
1364 .await
1365 .is_ok()
1366 {
1367 status.member(member.slot).add_rx(len as u64);
1368 }
1369 }
1370 }
1371 Err(()) => {
1372 status.note_transport_closure();
1373 close_member(index, manager, pending, status, fleet, backend, members).await;
1374 }
1375 }
1376}
1377
1378#[expect(
1379 clippy::too_many_arguments,
1380 reason = "embedded serve-loop internals pass the loop's split-borrowed locals; bundling awaits an on-hardware validation pass"
1381)]
1382async fn apply_settled<
1383 B,
1384 M: RawMutex + 'static,
1385 const FRAME: usize,
1386 const NOTIFY: usize,
1387 const LIFECYCLE: usize,
1388 const MEMBERS: usize,
1389>(
1390 link: B::Link,
1391 bitrate: BitrateBps,
1392 manager: &mut ConnectionPolicy<MEMBERS, DIAL_TRACK>,
1393 pending: &mut PendingActions<ACTION_CAP>,
1394 status: &BluetoothAutoStatus<MEMBERS>,
1395 fleet: &mut Fleet<M, FRAME, NOTIFY, LIFECYCLE>,
1396 backend: &mut B,
1397 members: &mut [Option<Active<B::Link>>; MEMBERS],
1398) where
1399 B: BleBackend<MEMBERS>,
1400{
1401 let address = link.address();
1402 if pending.overflowed {
1403 status.mark_failed(ACTION_OVERFLOW_REASON);
1404 drop(link);
1405 backend.on_link_closed(address).await;
1406 return;
1407 }
1408 let actions = pending.take();
1409 let mut held = Some(link);
1410 for action in actions {
1411 match action {
1412 PolicyAction::Admit {
1413 identity,
1414 slot,
1415 address,
1416 lane,
1417 } => {
1418 if let Some(mut link) = held.take() {
1419 if !matches!(lane, L2capPlan::None) {
1420 let _ = link.upgrade(&lane).await;
1421 }
1422 let (source, sink) = link.into_data();
1423 let id = InterfaceId::from_channel_tag(
1424 InterfaceKind::BluetoothPeer,
1425 identity.as_bytes(),
1426 );
1427 fleet
1428 .register_member(contract::descriptor(id, bitrate))
1429 .await;
1430 status.member(slot).assign(id);
1431 status.republish_peer_count();
1432 status.note_settled_link();
1433 members[slot] = Some(Active {
1434 identity,
1435 id,
1436 slot,
1437 address,
1438 source,
1439 sink,
1440 });
1441 }
1442 }
1443 PolicyAction::Reject { address, .. } => {
1444 held = None;
1445 backend.on_link_closed(address).await;
1446 }
1447 other => {
1448 apply_one(other, pending, manager, status, fleet, backend, members).await;
1449 }
1450 }
1451 }
1452 if held.is_some() {
1453 drop(held.take());
1454 backend.on_link_closed(address).await;
1455 }
1456 apply_radio(pending, manager, status, fleet, backend, members).await;
1457}
1458
1459#[cfg(test)]
1460mod tests {
1461 use core::sync::atomic::{AtomicUsize, Ordering as AtomicOrdering};
1462
1463 use embassy_futures::block_on;
1464 use embassy_futures::select::{select, Either};
1465 use prns_core::interfaces::bluetooth_auto::{Endpoint, Nrf52Host};
1466
1467 use super::*;
1468
1469 const CAPS: LinkCapabilities = LinkCapabilities {
1470 l2cap: None,
1471 link_mtu: contract::BLE_HW_MTU as u16,
1472 };
1473
1474 #[derive(Debug)]
1475 struct MockError;
1476
1477 #[derive(Clone, Copy)]
1478 enum MockSinkMode {
1479 Ready,
1480 Blocked,
1481 }
1482
1483 struct MockSource;
1484
1485 struct MockSink {
1486 mode: MockSinkMode,
1487 }
1488
1489 struct MockLink {
1490 address: BleAddress,
1491 protocol: PeerProtocol,
1492 incoming: Option<Control>,
1493 identity: BleIdentity,
1494 track_drop: bool,
1495 }
1496
1497 static DROPS: AtomicUsize = AtomicUsize::new(0);
1498
1499 impl Drop for MockLink {
1500 fn drop(&mut self) {
1501 if self.track_drop {
1502 DROPS.fetch_add(1, AtomicOrdering::Relaxed);
1503 }
1504 }
1505 }
1506
1507 impl BleLink for MockLink {
1508 type Error = MockError;
1509 type Source = MockSource;
1510 type Sink = MockSink;
1511
1512 fn peer_protocol(&self) -> PeerProtocol {
1513 self.protocol
1514 }
1515
1516 fn address(&self) -> BleAddress {
1517 self.address
1518 }
1519
1520 async fn receive_columba_peer_identity(&mut self) -> Result<BleIdentity, MockError> {
1521 Ok(self.identity)
1522 }
1523
1524 async fn send_columba_identity(&mut self, _identity: BleIdentity) -> Result<(), MockError> {
1525 Ok(())
1526 }
1527
1528 async fn control_send(&mut self, _msg: &Control) -> Result<(), MockError> {
1529 Ok(())
1530 }
1531
1532 async fn control_recv(&mut self) -> Result<Control, MockError> {
1533 match self.incoming.take() {
1534 Some(control) => Ok(control),
1535 None => ::core::future::pending().await,
1536 }
1537 }
1538
1539 async fn upgrade(&mut self, _plan: &L2capPlan) -> Result<(), MockError> {
1540 Ok(())
1541 }
1542
1543 fn into_data(self) -> (MockSource, MockSink) {
1544 (
1545 MockSource,
1546 MockSink {
1547 mode: MockSinkMode::Ready,
1548 },
1549 )
1550 }
1551 }
1552
1553 impl BleSource for MockSource {
1554 type Error = MockError;
1555
1556 async fn recv_frame(&mut self, _out: &mut [u8]) -> Result<usize, MockError> {
1557 ::core::future::pending().await
1558 }
1559 }
1560
1561 impl BleSink for MockSink {
1562 type Error = MockError;
1563
1564 async fn send_frame(&mut self, _frame: &[u8]) -> Result<(), MockError> {
1565 match self.mode {
1566 MockSinkMode::Ready => Ok(()),
1567 MockSinkMode::Blocked => ::core::future::pending().await,
1568 }
1569 }
1570 }
1571
1572 fn local(identity: u8) -> LocalPeer {
1573 LocalPeer {
1574 identity: BleIdentity::new([identity; 16]),
1575 endpoint: Endpoint::Nrf52(Nrf52Host::Nrf52),
1576 capabilities: CAPS,
1577 }
1578 }
1579
1580 fn link(address: u8, incoming: Option<Control>, track_drop: bool) -> MockLink {
1581 MockLink {
1582 address: BleAddress::new([address; 6]),
1583 protocol: PeerProtocol::Native,
1584 incoming,
1585 identity: BleIdentity::new([address; 16]),
1586 track_drop,
1587 }
1588 }
1589
1590 fn active(id: u8, mode: MockSinkMode) -> Active<MockLink> {
1591 Active {
1592 identity: BleIdentity::new([id; 16]),
1593 id: InterfaceId::new([id; 8]),
1594 slot: usize::from(id),
1595 address: BleAddress::new([id; 6]),
1596 source: MockSource,
1597 sink: MockSink { mode },
1598 }
1599 }
1600
1601 #[test]
1602 fn simultaneous_handshakes_let_the_ready_peer_settle() {
1603 let local = local(1);
1604 let hello = Control::Hello {
1605 identity: BleIdentity::new([3; 16]),
1606 endpoint: Endpoint::Nrf52(Nrf52Host::Nrf52),
1607 capabilities: CAPS,
1608 peer_rssi: None,
1609 };
1610 let mut handshakes = [
1611 Some(PendingHandshake::new(
1612 link(2, None, false),
1613 Origin::Dialed,
1614 local,
1615 )),
1616 Some(PendingHandshake::new(
1617 link(3, Some(hello), false),
1618 Origin::Accepted,
1619 local,
1620 )),
1621 ];
1622
1623 block_on(async {
1624 assert!(matches!(
1625 advance_handshakes(&mut handshakes, local).await,
1626 HandshakeStep::Advanced
1627 ));
1628 assert!(matches!(
1629 advance_handshakes(&mut handshakes, local).await,
1630 HandshakeStep::Advanced
1631 ));
1632 let step = advance_handshakes(&mut handshakes, local).await;
1633 assert!(matches!(&step, HandshakeStep::Done(_)));
1634 if let HandshakeStep::Done(done) = step {
1635 assert_eq!(done.address, BleAddress::new([3; 6]));
1636 assert_eq!(done.origin, Origin::Accepted);
1637 assert!(done.outcome.is_ok());
1638 }
1639 });
1640 assert!(handshakes[0].is_some());
1641 assert!(handshakes[1].is_none());
1642 }
1643
1644 #[test]
1645 fn cancelling_a_handshake_drops_its_link() {
1646 DROPS.store(0, AtomicOrdering::Relaxed);
1647 let mut handshakes = [
1648 Some(PendingHandshake::new(
1649 link(2, None, true),
1650 Origin::Dialed,
1651 local(1),
1652 )),
1653 None,
1654 ];
1655
1656 handshakes[0] = None;
1657
1658 assert_eq!(DROPS.load(AtomicOrdering::Relaxed), 1);
1659 }
1660
1661 #[test]
1662 fn blocked_peer_does_not_hide_completed_fanout() {
1663 let mut members = [
1664 Some(active(0, MockSinkMode::Ready)),
1665 Some(active(1, MockSinkMode::Blocked)),
1666 ];
1667 let mut states = [SendState::Pending, SendState::Pending];
1668
1669 block_on(async {
1670 assert!(matches!(
1671 select(send_members(&mut members, &mut states, b"frame"), async {}).await,
1672 Either::Second(())
1673 ));
1674 });
1675
1676 assert!(matches!(states, [SendState::Sent, SendState::Pending]));
1677 }
1678
1679 #[test]
1680 fn policy_action_overflow_is_explicit() {
1681 let mut pending = PendingActions::<1>::new();
1682 pending.push(PolicyAction::SetAdvertising(AdvertisingMode::On));
1683 pending.push(PolicyAction::SetScanning(ScanningMode::On));
1684
1685 assert!(pending.overflowed);
1686 assert_eq!(pending.actions.len(), 1);
1687 }
1688
1689 #[test]
1690 fn failed_status_survives_disable_and_reenable() {
1691 static SHARED: BluetoothAutoShared<1> = BluetoothAutoShared::new(InterfaceId::new([9; 8]));
1692 let status = BluetoothAutoStatus::new(&SHARED);
1693 status.mark_up();
1694 status.mark_failed(ACTION_OVERFLOW_REASON);
1695
1696 assert_eq!(status.connection(), ConnectionState::Failed);
1697 assert_eq!(status.failure_reason(), Some(ACTION_OVERFLOW_REASON));
1698 status.disable();
1699 assert_eq!(status.connection(), ConnectionState::Disabled);
1700 status.enable();
1701 assert_eq!(status.connection(), ConnectionState::Failed);
1702 }
1703
1704 fn recovery_view<const MEMBERS: usize>(
1705 status: &BluetoothAutoStatus<MEMBERS>,
1706 ) -> (
1707 ConnectionState,
1708 Option<BluetoothRecoveryReason>,
1709 BluetoothRecoveryCounters,
1710 Option<&'static str>,
1711 ) {
1712 (
1713 status.connection(),
1714 status.recovery_reason(),
1715 status.recovery_counters(),
1716 status.failure_reason(),
1717 )
1718 }
1719
1720 #[test]
1721 fn recovery_status_distinguishes_pressure_setup_and_closure() {
1722 static SHARED: BluetoothAutoShared<2> = BluetoothAutoShared::new(InterfaceId::new([10; 8]));
1723 let status = BluetoothAutoStatus::new(&SHARED);
1724
1725 assert_eq!(
1726 recovery_view(&status),
1727 (
1728 ConnectionState::Initializing,
1729 None,
1730 BluetoothRecoveryCounters::ZERO,
1731 None,
1732 )
1733 );
1734 status.mark_up();
1735 assert_eq!(
1736 recovery_view(&status),
1737 (
1738 ConnectionState::Disconnected,
1739 None,
1740 BluetoothRecoveryCounters::ZERO,
1741 None,
1742 )
1743 );
1744
1745 status.note_setup_failure();
1746 assert_eq!(
1747 recovery_view(&status),
1748 (
1749 ConnectionState::Reconnecting,
1750 Some(BluetoothRecoveryReason::SetupFailure),
1751 BluetoothRecoveryCounters {
1752 ingress_pressure: 0,
1753 setup_failures: 1,
1754 transport_closures: 0,
1755 },
1756 Some(SETUP_FAILURE_REASON),
1757 )
1758 );
1759
1760 status.member(0).assign(InterfaceId::new([11; 8]));
1761 status.republish_peer_count();
1762 assert_eq!(status.connection(), ConnectionState::Connected);
1763 assert_eq!(status.failure_reason(), Some(SETUP_FAILURE_REASON));
1764 status.note_settled_link();
1765 assert_eq!(status.connection(), ConnectionState::Connected);
1766
1767 status.note_ingress_pressure();
1768 assert_eq!(
1769 recovery_view(&status),
1770 (
1771 ConnectionState::Degraded,
1772 Some(BluetoothRecoveryReason::IngressPressure),
1773 BluetoothRecoveryCounters {
1774 ingress_pressure: 1,
1775 setup_failures: 1,
1776 transport_closures: 0,
1777 },
1778 Some(INGRESS_PRESSURE_REASON),
1779 )
1780 );
1781 status.note_successful_admission();
1782 assert_eq!(status.connection(), ConnectionState::Connected);
1783 assert_eq!(status.failure_reason(), None);
1784
1785 status.note_transport_closure();
1786 assert_eq!(status.connection(), ConnectionState::Connected);
1787 assert_eq!(status.failure_reason(), Some(TRANSPORT_CLOSURE_REASON));
1788 status.member(0).retire();
1789 status.republish_peer_count();
1790 assert_eq!(status.connection(), ConnectionState::Reconnecting);
1791
1792 status.member(1).assign(InterfaceId::new([12; 8]));
1793 status.republish_peer_count();
1794 status.note_settled_link();
1795 assert_eq!(
1796 recovery_view(&status),
1797 (
1798 ConnectionState::Connected,
1799 None,
1800 BluetoothRecoveryCounters {
1801 ingress_pressure: 1,
1802 setup_failures: 1,
1803 transport_closures: 1,
1804 },
1805 None,
1806 )
1807 );
1808 }
1809
1810 #[test]
1811 fn recovery_counters_saturate_as_one_snapshot() {
1812 static SHARED: BluetoothAutoShared<1> = BluetoothAutoShared::new(InterfaceId::new([13; 8]));
1813 let status = BluetoothAutoStatus::new(&SHARED);
1814 SHARED.recovery_counters.lock(|slot| {
1815 slot.set(BluetoothRecoveryCounters {
1816 ingress_pressure: u32::MAX,
1817 setup_failures: u32::MAX,
1818 transport_closures: u32::MAX,
1819 });
1820 });
1821
1822 status.note_ingress_pressure();
1823 status.note_setup_failure();
1824 status.note_transport_closure();
1825
1826 assert_eq!(
1827 status.recovery_counters(),
1828 BluetoothRecoveryCounters {
1829 ingress_pressure: u32::MAX,
1830 setup_failures: u32::MAX,
1831 transport_closures: u32::MAX,
1832 }
1833 );
1834 }
1835}