Skip to main content

wlan_telemetry/
lib.rs

1// Copyright 2024 The Fuchsia Authors. All rights reserved.
2// Use of this source code is governed by a BSD-style license that can be
3// found in the LICENSE file.
4use anyhow::{Context as _, Error, format_err};
5use fidl_fuchsia_power_battery as fidl_battery;
6use fidl_fuchsia_wlan_ieee80211 as fidl_ieee80211;
7use fidl_fuchsia_wlan_internal as fidl_internal;
8use fuchsia_async as fasync;
9use fuchsia_inspect::Node as InspectNode;
10use futures::channel::mpsc;
11use futures::{Future, StreamExt, select};
12use log::error;
13use std::boxed::Box;
14use std::sync::Arc;
15use windowed_stats::experimental::inspect::TimeMatrixClient;
16use wlan_common::bss::BssDescription;
17use wlan_legacy_metrics_registry as metrics;
18
19mod config;
20mod convert;
21mod processors;
22pub(crate) mod util;
23pub use crate::config::{CobaltAllowlist, DeviceMobility, TelemetryConfig};
24pub use crate::processors::connect_disconnect::DisconnectInfo;
25pub use crate::processors::pno_scan::PnoScanDisabledReason;
26pub use crate::processors::power::{IfacePowerLevel, UnclearPowerDemand};
27pub use crate::processors::scan::ScanResult;
28pub use crate::processors::sme_timeout::TimeoutSource;
29pub use crate::processors::toggle_events::ClientConnectionsToggleEvent;
30pub use util::sender::TelemetrySender;
31#[cfg(test)]
32mod testing;
33
34#[derive(Debug)]
35pub enum TelemetryEvent {
36    ConnectResult {
37        result: fidl_ieee80211::StatusCode,
38        bss: Box<BssDescription>,
39        is_credential_rejected: bool,
40        is_owe_transition: bool,
41    },
42    Disconnect {
43        info: DisconnectInfo,
44    },
45    ChannelSwitched {
46        channel: wlan_common::channel::Channel,
47    },
48    // We should maintain docstrings if we can see any possibility of ambiguity for an enum
49    /// Client connections enabled or disabled
50    ClientConnectionsToggle {
51        event: ClientConnectionsToggleEvent,
52    },
53    ClientIfaceCreated {
54        iface_id: u16,
55    },
56    ClientIfaceDestroyed {
57        iface_id: u16,
58    },
59    IfaceCreationFailure,
60    IfaceDestructionFailure,
61    ScanStart,
62    ScanResult {
63        result: ScanResult,
64    },
65    IfacePowerLevelChanged {
66        iface_power_level: IfacePowerLevel,
67        iface_id: u16,
68    },
69    /// System suspension imminent
70    SuspendImminent,
71    /// Unclear power level requested by policy layer
72    UnclearPowerDemand(UnclearPowerDemand),
73    BatteryChargeStatus(fidl_battery::ChargeStatus),
74    RecoveryEvent,
75    RecoveryResult {
76        result: Result<(), ()>,
77    },
78    SmeTimeout {
79        source: TimeoutSource,
80    },
81    ChipPowerUpFailure,
82    ChipPowerDownFailure,
83    ResetTxPowerScenario,
84    SetTxPowerScenario {
85        scenario: fidl_internal::TxPowerScenario,
86    },
87    PnoScanFailure,
88    PnoScanEnabled,
89    PnoScanResultsReceived,
90    PnoScanDisabled {
91        reason: PnoScanDisabledReason,
92    },
93}
94
95/// If metrics cannot be reported for extended periods of time, logging new metrics will fail and
96/// the error messages tend to clutter up the logs.  This container limits the rate at which such
97/// potentially noisy logs are reported.  Duplicate error messages are aggregated periodically
98/// reported.
99pub struct ThrottledErrorLogger {
100    time_of_last_log: fasync::MonotonicInstant,
101    pub suppressed_errors: std::collections::HashMap<String, usize>,
102    minutes_between_reports: i64,
103}
104
105impl ThrottledErrorLogger {
106    pub fn new(minutes_between_reports: i64) -> Self {
107        Self {
108            time_of_last_log: fasync::MonotonicInstant::from_nanos(0),
109            suppressed_errors: std::collections::HashMap::new(),
110            minutes_between_reports,
111        }
112    }
113
114    pub fn throttle_log(&mut self, message: String, level: log::Level) {
115        let curr_time = fasync::MonotonicInstant::now();
116        let time_since_last_log = curr_time - self.time_of_last_log;
117
118        if time_since_last_log.into_minutes() > self.minutes_between_reports {
119            log::log!(level, "{}", message);
120            if !self.suppressed_errors.is_empty() {
121                for (log, count) in self.suppressed_errors.iter() {
122                    log::warn!("Suppressed {} instances: {}", count, log);
123                }
124                self.suppressed_errors.clear();
125            }
126            self.time_of_last_log = curr_time;
127        } else {
128            let count = self.suppressed_errors.entry(message).or_default();
129            *count += 1;
130        }
131    }
132
133    pub fn throttle_error(&mut self, result: Result<(), Error>) {
134        if let Err(e) = result {
135            self.throttle_log(e.to_string(), log::Level::Error);
136        }
137    }
138}
139
140/// Attempts to connect to the Cobalt service.
141pub async fn setup_cobalt_proxy()
142-> Result<fidl_fuchsia_metrics::MetricEventLoggerProxy, anyhow::Error> {
143    let cobalt_svc = fuchsia_component::client::connect_to_protocol::<
144        fidl_fuchsia_metrics::MetricEventLoggerFactoryMarker,
145    >()
146    .context("failed to connect to metrics service")?;
147
148    let (cobalt_proxy, cobalt_server) =
149        fidl::endpoints::create_proxy::<fidl_fuchsia_metrics::MetricEventLoggerMarker>();
150
151    let project_spec = fidl_fuchsia_metrics::ProjectSpec {
152        customer_id: Some(metrics::CUSTOMER_ID),
153        project_id: Some(metrics::PROJECT_ID),
154        ..Default::default()
155    };
156
157    match cobalt_svc.create_metric_event_logger(&project_spec, cobalt_server).await {
158        Ok(_) => Ok(cobalt_proxy),
159        Err(err) => Err(format_err!("failed to create metrics event logger: {:?}", err)),
160    }
161}
162
163/// Attempts to create a disconnected FIDL channel with types matching the Cobalt service. This
164/// allows for a fallback with a uniform code path in case of a failure to connect to Cobalt.
165pub fn setup_disconnected_cobalt_proxy()
166-> Result<fidl_fuchsia_metrics::MetricEventLoggerProxy, anyhow::Error> {
167    // Create a disconnected proxy
168    Ok(fidl::endpoints::create_proxy::<fidl_fuchsia_metrics::MetricEventLoggerMarker>().0)
169}
170
171/// How often to refresh time series stats. Also how often to request packet counters.
172const TELEMETRY_QUERY_INTERVAL: zx::MonotonicDuration = zx::MonotonicDuration::from_seconds(10);
173
174pub fn serve_telemetry(
175    cobalt_proxy: fidl_fuchsia_metrics::MetricEventLoggerProxy,
176    monitor_svc_proxy: fidl_fuchsia_wlan_device_service::DeviceMonitorProxy,
177    inspect_node: InspectNode,
178    inspect_path: &str,
179    config: TelemetryConfig,
180    allowlist: CobaltAllowlist,
181) -> (TelemetrySender, impl Future<Output = Result<(), Error>> + use<>) {
182    let (sender, mut receiver) =
183        mpsc::channel::<TelemetryEvent>(util::sender::TELEMETRY_EVENT_BUFFER_SIZE);
184    let sender = TelemetrySender::new(sender);
185
186    let cobalt_logger =
187        Arc::new(util::cobalt_logger::FilteredCobaltLogger::new(cobalt_proxy, allowlist));
188
189    // Inspect nodes to hold time series and metadata for other nodes
190    const METADATA_NODE_NAME: &str = "metadata";
191    let inspect_metadata_node = inspect_node.create_child(METADATA_NODE_NAME);
192    let inspect_metadata_path = format!("{inspect_path}/{METADATA_NODE_NAME}");
193    let inspect_time_series_node = inspect_node.create_child("time_series");
194
195    let time_matrix_client = TimeMatrixClient::new(inspect_time_series_node.clone_weak());
196
197    // Create and initialize modules
198    let connect_disconnect = config.enable_connect_disconnect.then(|| {
199        processors::connect_disconnect::ConnectDisconnectLogger::new(
200            cobalt_logger.clone(),
201            &inspect_node,
202            &inspect_metadata_node,
203            &inspect_metadata_path,
204            &time_matrix_client,
205            config.device_mobility,
206        )
207    });
208    let iface_logger = config
209        .enable_iface_logger
210        .then(|| processors::iface::IfaceLogger::new(cobalt_logger.clone()));
211    let power_logger = config
212        .enable_power_logger
213        .then(|| processors::power::PowerLogger::new(cobalt_logger.clone(), &inspect_node));
214    let recovery_logger = config
215        .enable_recovery_logger
216        .then(|| processors::recovery::RecoveryLogger::new(cobalt_logger.clone()));
217    let mut scan_logger = config
218        .enable_scan_logger
219        .then(|| processors::scan::ScanLogger::new(cobalt_logger.clone(), &time_matrix_client));
220    let mut pno_scan_logger = config
221        .enable_pno_scan_logger
222        .then(|| processors::pno_scan::PnoScanLogger::new(cobalt_logger.clone()));
223    let sme_timeout_logger = config
224        .enable_sme_timeout_logger
225        .then(|| processors::sme_timeout::SmeTimeoutLogger::new(cobalt_logger.clone()));
226    let mut toggle_logger = config.enable_toggle_logger.then(|| {
227        processors::toggle_events::ToggleLogger::new(cobalt_logger.clone(), &inspect_node)
228    });
229    let tx_power_scenario_logger = config
230        .enable_tx_power_scenario_logger
231        .then(|| processors::tx_power_scenario::TxPowerScenarioLogger::new(cobalt_logger.clone()));
232
233    let client_iface_counters_logger = config.enable_client_iface_counters_logger.then(|| {
234        let driver_specific_time_series_node =
235            inspect_time_series_node.create_child("driver_specific");
236        let driver_counters_time_series_node =
237            driver_specific_time_series_node.create_child("counters");
238        let driver_gauges_time_series_node =
239            driver_specific_time_series_node.create_child("gauges");
240
241        let driver_counters_time_series_client =
242            TimeMatrixClient::new(driver_counters_time_series_node.clone_weak());
243        let driver_gauges_time_series_client =
244            TimeMatrixClient::new(driver_gauges_time_series_node.clone_weak());
245
246        inspect_time_series_node.record(driver_specific_time_series_node);
247        inspect_time_series_node.record(driver_counters_time_series_node);
248        inspect_time_series_node.record(driver_gauges_time_series_node);
249
250        processors::client_iface_counters::ClientIfaceCountersLogger::new(
251            cobalt_logger.clone(),
252            monitor_svc_proxy,
253            &inspect_metadata_node,
254            &inspect_metadata_path,
255            &time_matrix_client,
256            driver_counters_time_series_client,
257            driver_gauges_time_series_client,
258        )
259    });
260
261    let fut = async move {
262        // Prevent the inspect nodes from being dropped while the loop is running.
263        let _inspect_node = inspect_node;
264        let _inspect_metadata_node = inspect_metadata_node;
265        let _inspect_time_series_node = inspect_time_series_node;
266
267        let mut telemetry_interval = fasync::Interval::new(TELEMETRY_QUERY_INTERVAL);
268        loop {
269            select! {
270                event = receiver.next() => {
271                    let Some(event) = event else {
272                        error!("Telemetry event stream unexpectedly terminated.");
273                        return Err(format_err!("Telemetry event stream unexpectedly terminated."));
274                    };
275                    use TelemetryEvent::*;
276                    match event {
277                        ConnectResult { result, bss, is_credential_rejected, is_owe_transition } => {
278                            if let Some(ref connect_disconnect) = connect_disconnect {
279                                connect_disconnect.handle_connect_attempt(result, &bss, is_credential_rejected, is_owe_transition).await;
280                            }
281                        }
282                        Disconnect { info } => {
283                            if let Some(ref connect_disconnect) = connect_disconnect {
284                                connect_disconnect.log_disconnect(&info).await;
285                            }
286                            if let Some(ref power_logger) = power_logger {
287                                power_logger.handle_iface_disconnect(info.iface_id).await;
288                            }
289                        }
290                        ChannelSwitched { channel } => {
291                            if let Some(ref connect_disconnect) = connect_disconnect {
292                                connect_disconnect.handle_channel_switched(channel).await;
293                            }
294                        }
295                        ClientConnectionsToggle { event } => {
296                            if let Some(ref connect_disconnect) = connect_disconnect {
297                                connect_disconnect.handle_client_connections_toggle(&event).await;
298                            }
299                            if let Some(ref mut toggle_logger) = toggle_logger {
300                                toggle_logger.handle_toggle_event(event).await;
301                            }
302                        }
303                        ClientIfaceCreated { iface_id } => {
304                            if let Some(ref client_iface_counters_logger) = client_iface_counters_logger {
305                                client_iface_counters_logger.handle_iface_created(iface_id).await;
306                            }
307                        }
308                        ClientIfaceDestroyed { iface_id } => {
309                            if let Some(ref connect_disconnect) = connect_disconnect {
310                                connect_disconnect.handle_iface_destroyed().await;
311                            }
312                            if let Some(ref client_iface_counters_logger) = client_iface_counters_logger {
313                                client_iface_counters_logger.handle_iface_destroyed(iface_id).await;
314                            }
315                            if let Some(ref power_logger) = power_logger {
316                                power_logger.handle_iface_destroyed(iface_id).await;
317                            }
318                        }
319                        IfaceCreationFailure => {
320                            if let Some(ref iface_logger) = iface_logger {
321                                iface_logger.handle_iface_creation_failure().await;
322                            }
323                        }
324                        IfaceDestructionFailure => {
325                            if let Some(ref iface_logger) = iface_logger {
326                                iface_logger.handle_iface_destruction_failure().await;
327                            }
328                        }
329                        ScanStart => {
330                            if let Some(ref mut scan_logger) = scan_logger {
331                                scan_logger.handle_scan_start().await;
332                            }
333                        }
334                        ScanResult { result } => {
335                            if let Some(ref mut scan_logger) = scan_logger {
336                                scan_logger.handle_scan_result(result).await;
337                            }
338                        }
339                        IfacePowerLevelChanged { iface_power_level, iface_id } => {
340                            if let Some(ref power_logger) = power_logger {
341                                power_logger.log_iface_power_event(iface_power_level, iface_id).await;
342                            }
343                        }
344                        // TODO(b/340921554): either watch for suspension directly in the library,
345                        // or plumb this from callers once suspend mechanisms are integrated
346                        SuspendImminent => {
347                            if let Some(ref power_logger) = power_logger {
348                                power_logger.handle_suspend_imminent().await;
349                            }
350                            if let Some(ref connect_disconnect) = connect_disconnect {
351                                connect_disconnect.handle_suspend_imminent().await;
352                            }
353                        }
354                        UnclearPowerDemand(demand) => {
355                            if let Some(ref power_logger) = power_logger {
356                                power_logger.handle_unclear_power_demand(demand).await;
357                            }
358                        }
359                        ChipPowerUpFailure => {
360                            if let Some(ref power_logger) = power_logger {
361                                power_logger.handle_chip_power_up_failure().await;
362                            }
363                            if let Some(ref connect_disconnect) = connect_disconnect {
364                                connect_disconnect.handle_client_connections_failed_to_start().await;
365                            }
366                        }
367                        ChipPowerDownFailure => {
368                            if let Some(ref power_logger) = power_logger {
369                                power_logger.chip_power_down_failure().await;
370                            }
371                            if let Some(ref connect_disconnect) = connect_disconnect {
372                                connect_disconnect.handle_client_connections_failed_to_stop().await;
373                            }
374                        }
375                        BatteryChargeStatus(charge_status) => {
376                            if let Some(ref mut scan_logger) = scan_logger {
377                                scan_logger.handle_battery_charge_status(charge_status).await;
378                            }
379                            if let Some(ref mut toggle_logger) = toggle_logger {
380                                toggle_logger.handle_battery_charge_status(charge_status).await;
381                            }
382                        }
383                        RecoveryEvent => {
384                            if let Some(ref recovery_logger) = recovery_logger {
385                                recovery_logger.handle_recovery_event().await;
386                            }
387                        }
388                        RecoveryResult { result } => {
389                            if let Some(ref recovery_logger) = recovery_logger {
390                                recovery_logger.handle_recovery_result(result).await;
391                            }
392                        }
393                        SmeTimeout { source } => {
394                            if let Some(ref sme_timeout_logger) = sme_timeout_logger {
395                                sme_timeout_logger.handle_sme_timeout_event(source).await;
396                            }
397                        }
398                        ResetTxPowerScenario => {
399                            if let Some(ref tx_power_scenario_logger) = tx_power_scenario_logger {
400                                tx_power_scenario_logger.handle_sar_reset().await;
401                            }
402                        }
403                        SetTxPowerScenario {scenario} => {
404                            if let Some(ref tx_power_scenario_logger) = tx_power_scenario_logger {
405                                tx_power_scenario_logger.handle_set_sar(scenario).await;
406                            }
407                        }
408                        PnoScanFailure => {
409                            if let Some(ref connect_disconnect) = connect_disconnect {
410                                connect_disconnect.handle_pno_scan_failure().await;
411                            }
412                        }
413                        PnoScanEnabled => {
414                            if let Some(ref mut pno_scan_logger) = pno_scan_logger {
415                                let is_connected = connect_disconnect
416                                    .as_ref()
417                                    .map(|cd| cd.is_connected())
418                                    .unwrap_or(false);
419                                pno_scan_logger.handle_pno_scan_enabled(is_connected).await;
420                            }
421                        }
422                        PnoScanResultsReceived => {
423                            if let Some(ref mut pno_scan_logger) = pno_scan_logger {
424                                pno_scan_logger.handle_pno_scan_results_received().await;
425                            }
426                        }
427                        PnoScanDisabled { reason } => {
428                            if let Some(ref mut pno_scan_logger) = pno_scan_logger {
429                                pno_scan_logger.handle_pno_scan_disabled(reason).await;
430                            }
431                        }
432                    }
433                }
434                _ = telemetry_interval.next() => {
435                    if let Some(ref connect_disconnect) = connect_disconnect {
436                        connect_disconnect.handle_periodic_telemetry().await;
437                    }
438                    if let Some(ref client_iface_counters_logger) = client_iface_counters_logger {
439                        client_iface_counters_logger.handle_periodic_telemetry().await;
440                    }
441                    if let Some(ref mut pno_scan_logger) = pno_scan_logger {
442                        pno_scan_logger.handle_periodic_telemetry().await;
443                    }
444                }
445            }
446        }
447    };
448    (sender, fut)
449}