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}