1pub mod processors;
6
7use crate::telemetry::processors::network_properties::NetworkPropertiesProcessor;
8use anyhow::Error;
9use fidl_fuchsia_net_policy_socketproxy as fnp_socketproxy;
10use fuchsia_sync::Mutex;
11use futures::channel::mpsc;
12use futures::{Future, StreamExt};
13use log::{info, warn};
14use std::sync::Arc;
15use std::sync::atomic::{AtomicBool, Ordering};
16
17#[derive(Clone, Debug)]
18pub struct NetworkEventMetadata {
19 pub id: u64,
20 pub name: Option<String>,
21 pub transport: fnp_socketproxy::NetworkType,
22 pub is_fuchsia_provisioned: bool,
23 pub connectivity_state: Option<fnp_socketproxy::ConnectivityState>,
24}
25
26#[derive(Debug)]
27pub enum TelemetryEvent {
28 DefaultNetworkChanged(NetworkEventMetadata),
29 DefaultNetworkLost,
30 NetworkChanged(NetworkEventMetadata),
31}
32
33#[derive(Clone, Debug)]
34pub struct TelemetrySender {
35 sender: Arc<Mutex<mpsc::Sender<TelemetryEvent>>>,
36 sender_is_blocked: Arc<AtomicBool>,
37}
38
39impl TelemetrySender {
40 pub fn new(sender: mpsc::Sender<TelemetryEvent>) -> Self {
41 Self {
42 sender: Arc::new(Mutex::new(sender)),
43 sender_is_blocked: Arc::new(AtomicBool::new(false)),
44 }
45 }
46
47 pub fn send(&self, event: TelemetryEvent) {
48 match self.sender.lock().try_send(event) {
49 Ok(_) => {
50 if self
51 .sender_is_blocked
52 .compare_exchange(true, false, Ordering::SeqCst, Ordering::SeqCst)
53 .is_ok()
54 {
55 info!("TelemetrySender recovered and resumed sending");
56 }
57 }
58 Err(_) => {
59 if self
60 .sender_is_blocked
61 .compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst)
62 .is_ok()
63 {
64 warn!(
65 "TelemetrySender dropped a msg: either buffer is full or no receiver is waiting"
66 );
67 }
68 }
69 }
70 }
71}
72
73const TELEMETRY_EVENT_BUFFER_SIZE: usize = 100;
74
75pub fn serve_telemetry(
76 telemetry_node: fuchsia_inspect::Node,
77 telemetry_path: &str,
78) -> (TelemetrySender, impl Future<Output = Result<(), Error>> + use<>) {
79 let time_series_node = telemetry_node.create_child("time_series");
80 let client =
81 windowed_stats::experimental::inspect::TimeMatrixClient::new(time_series_node.clone_weak());
82
83 let processor = NetworkPropertiesProcessor::new(&telemetry_node, telemetry_path, &client);
84 telemetry_node.record(time_series_node);
85
86 let (sender, mut receiver) = mpsc::channel::<TelemetryEvent>(TELEMETRY_EVENT_BUFFER_SIZE);
87 let sender = TelemetrySender::new(sender);
88
89 let fut = async move {
90 let mut processor = processor;
91 while let Some(event) = receiver.next().await {
92 match event {
93 TelemetryEvent::DefaultNetworkChanged(metadata) => {
94 processor.log_default_network_changed(metadata);
95 }
96 TelemetryEvent::DefaultNetworkLost => {
97 processor.log_default_network_lost();
98 }
99 TelemetryEvent::NetworkChanged(metadata) => {
100 processor.log_network_changed(metadata, &client);
101 }
102 }
103 }
104 Ok(())
105 };
106 (sender, fut)
107}