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}