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