prns_runtime_embassy/manifold/driver/
interface_status.rs1use embassy_sync::blocking_mutex::raw::CriticalSectionRawMutex;
2use embassy_sync::signal::Signal;
3use portable_atomic::{AtomicBool, AtomicU32, AtomicU64, AtomicU8, Ordering};
4
5use crate::interfaces::{
6 AirtimeUtilization, ConnectionState, InterfaceId, InterfaceStatus, TransferRates,
7};
8
9pub struct EmbassyInterfaceStatus {
10 id: AtomicU64,
11 connection: AtomicU8,
12 rx: AtomicU64,
13 tx: AtomicU64,
14 airtime: AtomicU32,
15 transfer_rates: AtomicU64,
16 enabled: AtomicBool,
17 enabled_changed: Signal<CriticalSectionRawMutex, bool>,
18}
19
20const AIRTIME_UNPUBLISHED: u32 = u32::MAX;
21const RATES_UNPUBLISHED: u64 = u64::MAX;
22
23impl EmbassyInterfaceStatus {
24 #[must_use]
25 pub const fn new(id: InterfaceId, connection: ConnectionState) -> Self {
26 Self {
27 id: AtomicU64::new(u64::from_be_bytes(*id.as_bytes())),
28 connection: AtomicU8::new(connection.as_u8()),
29 rx: AtomicU64::new(0),
30 tx: AtomicU64::new(0),
31 airtime: AtomicU32::new(AIRTIME_UNPUBLISHED),
32 transfer_rates: AtomicU64::new(RATES_UNPUBLISHED),
33 enabled: AtomicBool::new(true),
34 enabled_changed: Signal::new(),
35 }
36 }
37
38 pub fn set_connection(&self, connection: ConnectionState) {
39 self.connection.store(connection.as_u8(), Ordering::Relaxed);
40 }
41
42 pub fn set_id(&self, id: InterfaceId) {
43 self.id
44 .store(u64::from_be_bytes(*id.as_bytes()), Ordering::Relaxed);
45 }
46
47 pub fn enable(&self) {
48 self.update_enabled(true);
49 }
50
51 pub fn disable(&self) {
52 self.update_enabled(false);
53 }
54
55 pub fn toggle_enabled(&self) {
56 let enabled = !self.enabled.fetch_xor(true, Ordering::Relaxed);
57 self.enabled_changed.signal(enabled);
58 }
59
60 fn update_enabled(&self, enabled: bool) {
61 if self.enabled.swap(enabled, Ordering::Relaxed) != enabled {
62 self.enabled_changed.signal(enabled);
63 }
64 }
65
66 #[must_use]
67 pub fn is_enabled(&self) -> bool {
68 self.enabled.load(Ordering::Relaxed)
69 }
70
71 pub async fn wait_until_enabled(&self) {
72 self.wait_for_enabled_state(true).await;
73 }
74
75 pub async fn wait_until_disabled(&self) {
76 self.wait_for_enabled_state(false).await;
77 }
78
79 async fn wait_for_enabled_state(&self, enabled: bool) {
80 loop {
81 if self.is_enabled() == enabled {
82 return;
83 }
84 if self.enabled_changed.wait().await == enabled {
85 return;
86 }
87 }
88 }
89
90 pub fn add_rx(&self, bytes: u64) {
91 self.rx.fetch_add(bytes, Ordering::Relaxed);
92 }
93
94 pub fn add_tx(&self, bytes: u64) {
95 self.tx.fetch_add(bytes, Ordering::Relaxed);
96 }
97
98 pub fn set_airtime(&self, utilization: AirtimeUtilization) {
99 let packed =
100 (u32::from(utilization.short_per_mille) << 16) | u32::from(utilization.long_per_mille);
101 self.airtime.store(packed, Ordering::Relaxed);
102 }
103
104 pub fn set_transfer_rates(&self, rates: TransferRates) {
105 let packed = (u64::from(rates.rx_bps) << 32) | u64::from(rates.tx_bps);
106 self.transfer_rates.store(packed, Ordering::Relaxed);
107 }
108}
109
110impl InterfaceStatus for EmbassyInterfaceStatus {
111 fn id(&self) -> InterfaceId {
112 InterfaceId::new(self.id.load(Ordering::Relaxed).to_be_bytes())
113 }
114
115 fn connection(&self) -> ConnectionState {
116 if !self.is_enabled() {
117 return ConnectionState::Disabled;
118 }
119 ConnectionState::from_u8(self.connection.load(Ordering::Relaxed))
120 }
121
122 fn rx_bytes(&self) -> u64 {
123 self.rx.load(Ordering::Relaxed)
124 }
125
126 fn tx_bytes(&self) -> u64 {
127 self.tx.load(Ordering::Relaxed)
128 }
129
130 fn airtime(&self) -> Option<AirtimeUtilization> {
131 let packed = self.airtime.load(Ordering::Relaxed);
132 if packed == AIRTIME_UNPUBLISHED {
133 return None;
134 }
135 Some(AirtimeUtilization {
136 short_per_mille: (packed >> 16) as u16,
137 long_per_mille: packed as u16,
138 })
139 }
140
141 fn transfer_rates(&self) -> Option<TransferRates> {
142 let packed = self.transfer_rates.load(Ordering::Relaxed);
143 if packed == RATES_UNPUBLISHED {
144 return None;
145 }
146 Some(TransferRates {
147 rx_bps: (packed >> 32) as u32,
148 tx_bps: packed as u32,
149 })
150 }
151}
152
153#[cfg(test)]
154mod tests {
155 use super::*;
156 use embassy_futures::{block_on, join::join};
157
158 #[test]
159 fn enabled_state_changes_wake_waiters() {
160 let status =
161 EmbassyInterfaceStatus::new(InterfaceId::new([0x5A; 8]), ConnectionState::Initializing);
162
163 block_on(async {
164 join(status.wait_until_disabled(), async {
165 status.disable();
166 })
167 .await;
168 join(status.wait_until_enabled(), async {
169 status.toggle_enabled();
170 })
171 .await;
172 });
173 assert!(status.is_enabled());
174 status.toggle_enabled();
175 assert!(!status.is_enabled());
176 status.enable();
177 assert!(status.is_enabled());
178 }
179}