monitord/
machines.rs

1use std::collections::HashMap;
2use std::collections::HashSet;
3use std::sync::Arc;
4
5use thiserror::Error;
6use tokio::sync::RwLock;
7use tracing::{debug, error, warn};
8
9use crate::MachineStats;
10use crate::MonitordStats;
11
12#[derive(Error, Debug)]
13pub enum MonitordMachinesError {
14    #[error("Machines D-Bus error: {0}")]
15    ZbusError(#[from] zbus::Error),
16}
17
18pub fn filter_machines(
19    machines: Vec<crate::dbus::zbus_machines::ListedMachine>,
20    allowlist: &HashSet<String>,
21    blocklist: &HashSet<String>,
22) -> Vec<crate::dbus::zbus_machines::ListedMachine> {
23    machines
24        .into_iter()
25        .filter(|c| c.class == "container")
26        .filter(|c| !blocklist.contains(&c.name))
27        .filter(|c| allowlist.is_empty() || allowlist.contains(&c.name))
28        .collect()
29}
30
31pub async fn get_machines(
32    connection: &zbus::Connection,
33    config: &crate::config::Config,
34) -> Result<HashMap<String, u32>, MonitordMachinesError> {
35    let c = crate::dbus::zbus_machines::ManagerProxy::new(connection).await?;
36    let mut results = HashMap::<String, u32>::new();
37
38    let machines = c.list_machines().await?;
39
40    for machine in filter_machines(
41        machines,
42        &config.machines.allowlist,
43        &config.machines.blocklist,
44    ) {
45        let m = c.get_machine(&machine.name).await?;
46        let leader_pid = m.leader().await?;
47        results.insert(machine.name, leader_pid);
48    }
49
50    Ok(results)
51}
52
53pub async fn update_machines_stats(
54    config: Arc<crate::config::Config>,
55    connection: zbus::Connection,
56    locked_monitord_stats: Arc<RwLock<MonitordStats>>,
57) -> anyhow::Result<()> {
58    let locked_machine_stats: Arc<RwLock<MachineStats>> =
59        Arc::new(RwLock::new(MachineStats::default()));
60
61    for (machine, leader_pid) in get_machines(&connection, &config).await?.into_iter() {
62        debug!(
63            "Collecting container: machine: {} leader_pid: {}",
64            machine, leader_pid
65        );
66        let container_address = format!(
67            "unix:path=/proc/{}/root/run/dbus/system_bus_socket",
68            leader_pid
69        );
70        let sdc = zbus::connection::Builder::address(container_address.as_str())?
71            .method_timeout(std::time::Duration::from_secs(config.monitord.dbus_timeout))
72            .build()
73            .await?;
74        let mut join_set = tokio::task::JoinSet::new();
75
76        if config.pid1.enabled {
77            join_set.spawn(crate::pid1::update_pid1_stats(
78                leader_pid as i32,
79                locked_machine_stats.clone(),
80            ));
81        }
82
83        if config.networkd.enabled {
84            join_set.spawn(crate::networkd::update_networkd_stats(
85                config.networkd.link_state_dir.clone(),
86                None,
87                sdc.clone(),
88                locked_machine_stats.clone(),
89            ));
90        }
91
92        if config.system_state.enabled {
93            join_set.spawn(crate::system::update_system_stats(
94                sdc.clone(),
95                locked_machine_stats.clone(),
96            ));
97        }
98
99        join_set.spawn(crate::system::update_version(
100            sdc.clone(),
101            locked_machine_stats.clone(),
102        ));
103
104        if config.units.enabled {
105            if config.varlink.enabled {
106                let config_clone = Arc::clone(&config);
107                let sdc_clone = sdc.clone();
108                let stats_clone = locked_machine_stats.clone();
109                let container_socket_path = format!(
110                    "/proc/{}/root{}",
111                    leader_pid,
112                    crate::varlink_units::METRICS_SOCKET_PATH
113                );
114                join_set.spawn(async move {
115                    match crate::varlink_units::update_unit_stats(
116                        Arc::clone(&config_clone),
117                        stats_clone.clone(),
118                        container_socket_path,
119                    )
120                    .await
121                    {
122                        Ok(()) => Ok(()),
123                        Err(err) => {
124                            warn!(
125                                "Varlink units stats failed, falling back to D-Bus: {:?}",
126                                err
127                            );
128                            crate::units::update_unit_stats(config_clone, sdc_clone, stats_clone)
129                                .await
130                        }
131                    }
132                });
133            } else {
134                join_set.spawn(crate::units::update_unit_stats(
135                    Arc::clone(&config),
136                    sdc.clone(),
137                    locked_machine_stats.clone(),
138                ));
139            }
140        }
141
142        if config.dbus_stats.enabled {
143            join_set.spawn(crate::dbus_stats::update_dbus_stats(
144                Arc::clone(&config),
145                sdc.clone(),
146                locked_machine_stats.clone(),
147            ));
148        }
149
150        while let Some(res) = join_set.join_next().await {
151            match res {
152                Ok(r) => match r {
153                    Ok(_) => (),
154                    Err(e) => {
155                        error!(
156                            "Collection specific failure (container {}): {:?}",
157                            machine, e
158                        );
159                    }
160                },
161                Err(e) => {
162                    error!("Join error (container {}): {:?}", machine, e);
163                }
164            }
165        }
166
167        {
168            let mut monitord_stats = locked_monitord_stats.write().await;
169            let machine_stats = locked_machine_stats.read().await;
170            monitord_stats
171                .machines
172                .insert(machine, machine_stats.clone());
173        }
174    }
175
176    Ok(())
177}
178
179#[cfg(test)]
180mod tests {
181    use std::collections::HashSet;
182    use zbus::zvariant::OwnedObjectPath;
183
184    #[test]
185    fn test_filter_machines() {
186        let machines = vec![
187            crate::dbus::zbus_machines::ListedMachine {
188                name: "foo".to_string(),
189                class: "container".to_string(),
190                service: "".to_string(),
191                path: OwnedObjectPath::try_from("/sample/object").unwrap(),
192            },
193            crate::dbus::zbus_machines::ListedMachine {
194                name: "bar".to_string(),
195                class: "container".to_string(),
196                service: "".to_string(),
197                path: OwnedObjectPath::try_from("/sample/object").unwrap(),
198            },
199            crate::dbus::zbus_machines::ListedMachine {
200                name: "baz".to_string(),
201                class: "container".to_string(),
202                service: "".to_string(),
203                path: OwnedObjectPath::try_from("/sample/object").unwrap(),
204            },
205        ];
206        let allowlist = HashSet::from(["foo".to_string(), "baz".to_string()]);
207        let blocklist = HashSet::from(["bar".to_string()]);
208
209        let filtered = super::filter_machines(machines, &allowlist, &blocklist);
210
211        assert_eq!(filtered.len(), 2);
212        assert_eq!(filtered[0].name, "foo");
213        assert_eq!(filtered[1].name, "baz");
214    }
215}