Skip to main content

prns_runtime/runtime/node_introspection/
core.rs

1use prns_core::interfaces::IfacSize;
2use prns_core::interfaces::{
3    ConnectionState, InterfaceGravity, InterfaceId, InterfaceMode, InterfaceOriginKind,
4    InterfaceSnapshot, Membership, TransferRates,
5};
6
7#[derive(Debug, Clone, PartialEq, Eq)]
8pub struct InterfaceIfacSnapshot<Label> {
9    pub signature: [u8; 64],
10    pub size: IfacSize,
11    pub network_name: Option<Label>,
12}
13
14#[derive(Debug, Clone, PartialEq, Eq)]
15pub struct InterfaceInventoryEntry<Label> {
16    pub name: Option<Label>,
17    pub origin: InterfaceOriginKind,
18    pub snapshot: InterfaceSnapshot,
19    pub ifac: Option<InterfaceIfacSnapshot<Label>>,
20}
21
22struct FoldedInterface<Label> {
23    id: InterfaceId,
24    name: Option<Label>,
25    origin: InterfaceOriginKind,
26    root: Option<InterfaceSnapshot>,
27    ifac: Option<InterfaceIfacSnapshot<Label>>,
28    member_connection: ConnectionState,
29    member_mode: Option<InterfaceMode>,
30    member_gravity: Option<InterfaceGravity>,
31    member_failure_reason: Option<&'static str>,
32    member_rx_bytes: u64,
33    member_tx_bytes: u64,
34    member_rates: Option<TransferRates>,
35    has_members: bool,
36    destinations: u32,
37    links: u32,
38    transported_links: u32,
39}
40
41impl<Label> FoldedInterface<Label> {
42    fn new(id: InterfaceId, origin: InterfaceOriginKind) -> Self {
43        Self {
44            id,
45            name: None,
46            origin,
47            root: None,
48            ifac: None,
49            member_connection: ConnectionState::Unknown,
50            member_mode: None,
51            member_gravity: None,
52            member_failure_reason: None,
53            member_rx_bytes: 0,
54            member_tx_bytes: 0,
55            member_rates: None,
56            has_members: false,
57            destinations: 0,
58            links: 0,
59            transported_links: 0,
60        }
61    }
62
63    fn add(&mut self, entry: &mut InterfaceInventoryEntry<Label>) {
64        let snapshot = entry.snapshot;
65        self.destinations = self.destinations.saturating_add(snapshot.destinations);
66        self.links = self.links.saturating_add(snapshot.links);
67        self.transported_links = self
68            .transported_links
69            .saturating_add(snapshot.transported_links);
70        match snapshot.membership {
71            Membership::Independent => {
72                self.root = Some(snapshot);
73                self.origin = entry.origin;
74                if entry.name.is_some() {
75                    self.name = entry.name.take();
76                }
77                if entry.ifac.is_some() {
78                    self.ifac = entry.ifac.take();
79                }
80            }
81            Membership::FleetMember { .. } => {
82                if self.root.is_none() && entry.origin == InterfaceOriginKind::Discovered {
83                    self.origin = InterfaceOriginKind::Discovered;
84                }
85                self.has_members = true;
86                self.member_mode.get_or_insert(snapshot.mode);
87                self.member_gravity.get_or_insert(snapshot.gravity);
88                self.member_connection =
89                    preferred_connection(self.member_connection, snapshot.connection);
90                self.member_failure_reason = self.member_failure_reason.or(snapshot.failure_reason);
91                self.member_rx_bytes = self.member_rx_bytes.saturating_add(snapshot.rx_bytes);
92                self.member_tx_bytes = self.member_tx_bytes.saturating_add(snapshot.tx_bytes);
93                if let Some(rates) = snapshot.transfer_rates {
94                    let aggregate = self.member_rates.get_or_insert(TransferRates {
95                        rx_bps: 0,
96                        tx_bps: 0,
97                    });
98                    aggregate.rx_bps = aggregate.rx_bps.saturating_add(rates.rx_bps);
99                    aggregate.tx_bps = aggregate.tx_bps.saturating_add(rates.tx_bps);
100                }
101                if self.name.is_none() {
102                    self.name = entry.name.take();
103                }
104                if self.ifac.is_none() {
105                    self.ifac = entry.ifac.take();
106                }
107            }
108        }
109    }
110
111    fn finish(self) -> InterfaceInventoryEntry<Label> {
112        let connection = self
113            .root
114            .map_or(self.member_connection, |snapshot| snapshot.connection);
115        let failure_reason = self
116            .root
117            .and_then(|snapshot| snapshot.failure_reason)
118            .or(self.member_failure_reason);
119        let mode = self
120            .root
121            .map(|snapshot| snapshot.mode)
122            .or(self.member_mode)
123            .unwrap_or(InterfaceMode::Full);
124        let gravity = self
125            .root
126            .map(|snapshot| snapshot.gravity)
127            .or(self.member_gravity)
128            .unwrap_or(InterfaceGravity::ZERO);
129        let root_rx_bytes = self.root.map_or(0, |snapshot| snapshot.rx_bytes);
130        let root_tx_bytes = self.root.map_or(0, |snapshot| snapshot.tx_bytes);
131        let rx_bytes = root_rx_bytes.max(self.member_rx_bytes);
132        let tx_bytes = root_tx_bytes.max(self.member_tx_bytes);
133        let transfer_rates = if self.has_members {
134            self.member_rates
135        } else {
136            self.root.and_then(|snapshot| snapshot.transfer_rates)
137        };
138        InterfaceInventoryEntry {
139            name: self.name,
140            origin: self.origin,
141            snapshot: InterfaceSnapshot {
142                id: self.id,
143                mode,
144                gravity,
145                connection,
146                failure_reason,
147                rx_bytes,
148                tx_bytes,
149                transfer_rates,
150                destinations: self.destinations,
151                links: self.links,
152                transported_links: self.transported_links,
153                membership: Membership::Independent,
154            },
155            ifac: self.ifac,
156        }
157    }
158}
159
160#[must_use]
161pub fn fold_logical_interface_inventory<Label: Ord>(
162    inventory: &mut [InterfaceInventoryEntry<Label>],
163) -> &mut [InterfaceInventoryEntry<Label>] {
164    inventory.sort_unstable_by(|left, right| {
165        logical_interface_id(left)
166            .cmp(&logical_interface_id(right))
167            .then_with(|| membership_rank(left).cmp(&membership_rank(right)))
168            .then_with(|| left.snapshot.id.cmp(&right.snapshot.id))
169    });
170
171    let mut read = 0;
172    let mut write = 0;
173    while read < inventory.len() {
174        let logical_id = logical_interface_id(&inventory[read]);
175        let mut end = read + 1;
176        while end < inventory.len() && logical_interface_id(&inventory[end]) == logical_id {
177            end += 1;
178        }
179        let mut folded = FoldedInterface::new(logical_id, inventory[read].origin);
180        for entry in &mut inventory[read..end] {
181            folded.add(entry);
182        }
183        inventory[write] = folded.finish();
184        write += 1;
185        read = end;
186    }
187
188    let logical = &mut inventory[..write];
189    logical.sort_unstable_by(|left, right| {
190        left.name
191            .cmp(&right.name)
192            .then_with(|| left.snapshot.id.cmp(&right.snapshot.id))
193    });
194    logical
195}
196
197fn logical_interface_id<Label>(entry: &InterfaceInventoryEntry<Label>) -> InterfaceId {
198    match entry.snapshot.membership {
199        Membership::Independent => entry.snapshot.id,
200        Membership::FleetMember { supervisor_id } => supervisor_id,
201    }
202}
203
204fn membership_rank<Label>(entry: &InterfaceInventoryEntry<Label>) -> u8 {
205    match entry.snapshot.membership {
206        Membership::Independent => 0,
207        Membership::FleetMember { .. } => 1,
208    }
209}
210
211fn preferred_connection(left: ConnectionState, right: ConnectionState) -> ConnectionState {
212    if connection_rank(left) <= connection_rank(right) {
213        left
214    } else {
215        right
216    }
217}
218
219fn connection_rank(state: ConnectionState) -> u8 {
220    match state {
221        ConnectionState::Connected => 0,
222        ConnectionState::Degraded => 1,
223        ConnectionState::Initializing => 2,
224        ConnectionState::Reconnecting => 3,
225        ConnectionState::Failed => 4,
226        ConnectionState::Disconnected => 5,
227        ConnectionState::Disabled => 6,
228        ConnectionState::Unknown => 7,
229    }
230}
231
232#[cfg(test)]
233mod tests {
234    use super::*;
235    use prns_core::interfaces::InterfaceKind;
236
237    fn snapshot(
238        id: InterfaceId,
239        membership: Membership,
240        rx_bytes: u64,
241        destinations: u32,
242        links: u32,
243        name: Option<&'static str>,
244    ) -> InterfaceInventoryEntry<&'static str> {
245        InterfaceInventoryEntry {
246            name,
247            origin: InterfaceOriginKind::Configured,
248            snapshot: InterfaceSnapshot {
249                id,
250                mode: InterfaceMode::Full,
251                gravity: InterfaceGravity::ZERO,
252                connection: ConnectionState::Connected,
253                failure_reason: None,
254                rx_bytes,
255                tx_bytes: rx_bytes / 2,
256                transfer_rates: Some(TransferRates {
257                    rx_bps: rx_bytes as u32,
258                    tx_bps: (rx_bytes / 2) as u32,
259                }),
260                destinations,
261                links,
262                transported_links: 0,
263                membership,
264            },
265            ifac: None,
266        }
267    }
268
269    #[test]
270    fn fleet_members_fold_into_the_named_supervisor() {
271        let supervisor = InterfaceId::from_channel_tag(InterfaceKind::TcpServer, b"server");
272        let first = InterfaceId::from_channel_tag(InterfaceKind::TcpServerPeer, b"first");
273        let second = InterfaceId::from_channel_tag(InterfaceKind::TcpServerPeer, b"second");
274        let membership = Membership::FleetMember {
275            supervisor_id: supervisor,
276        };
277        let mut snapshots = [
278            snapshot(second, membership, 60, 3, 2, None),
279            snapshot(
280                supervisor,
281                Membership::Independent,
282                10,
283                0,
284                0,
285                Some("Public server"),
286            ),
287            snapshot(first, membership, 40, 2, 1, None),
288        ];
289
290        let logical = fold_logical_interface_inventory(&mut snapshots);
291
292        assert_eq!(logical.len(), 1);
293        assert_eq!(logical[0].name, Some("Public server"));
294        assert_eq!(logical[0].origin, InterfaceOriginKind::Configured);
295        assert_eq!(logical[0].snapshot.id, supervisor);
296        assert_eq!(logical[0].snapshot.rx_bytes, 100);
297        assert_eq!(logical[0].snapshot.destinations, 5);
298        assert_eq!(logical[0].snapshot.links, 3);
299        assert_eq!(logical[0].snapshot.membership, Membership::Independent);
300    }
301
302    #[test]
303    fn retained_supervisor_traffic_is_not_readded_to_live_members() {
304        let supervisor = InterfaceId::from_channel_tag(InterfaceKind::AutoWifi, b"default");
305        let peer = InterfaceId::from_channel_tag(InterfaceKind::WifiPeer, b"peer");
306        let membership = Membership::FleetMember {
307            supervisor_id: supervisor,
308        };
309        let mut snapshots = [
310            snapshot(
311                supervisor,
312                Membership::Independent,
313                120,
314                0,
315                0,
316                Some("Default Interface"),
317            ),
318            snapshot(peer, membership, 30, 0, 0, None),
319        ];
320
321        let logical = fold_logical_interface_inventory(&mut snapshots);
322
323        assert_eq!(logical.len(), 1);
324        assert_eq!(logical[0].snapshot.rx_bytes, 120);
325        assert_eq!(logical[0].snapshot.tx_bytes, 60);
326    }
327
328    #[test]
329    fn discovered_origin_survives_logical_inventory_folding() {
330        let id = InterfaceId::from_channel_tag(InterfaceKind::BackboneClient, b"discovered");
331        let mut snapshots = [snapshot(
332            id,
333            Membership::Independent,
334            100,
335            2,
336            1,
337            Some("Discovered backbone"),
338        )];
339        snapshots[0].origin = InterfaceOriginKind::Discovered;
340        let expected = snapshots[0].clone();
341
342        let logical = fold_logical_interface_inventory(&mut snapshots);
343
344        assert_eq!(logical, core::slice::from_ref(&expected));
345    }
346}