Skip to main content

prns_runtime_embassy/manifold/driver/
interface_status.rs

1use 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}