monitord/lib.rs
1//! # monitord Crate
2//!
3//! `monitord` is a library to gather statistics about systemd.
4
5use std::sync::Arc;
6
7use std::collections::HashMap;
8use std::time::Duration;
9use std::time::Instant;
10
11use thiserror::Error;
12use tokio::sync::RwLock;
13use tracing::debug;
14use tracing::error;
15use tracing::info;
16use tracing::warn;
17use tracing::Instrument;
18
19#[derive(Error, Debug)]
20pub enum MonitordError {
21 #[error("D-Bus connection error: {0}")]
22 ZbusError(#[from] zbus::Error),
23}
24
25impl MonitordError {
26 /// Unwrap the inner zbus error. Exhaustive today: connection setup is
27 /// the only thing that produces this error, so there is exactly one
28 /// variant to unwrap.
29 pub fn into_zbus(self) -> zbus::Error {
30 match self {
31 MonitordError::ZbusError(inner) => inner,
32 }
33 }
34}
35
36pub mod boot;
37pub mod cgroup;
38pub mod config;
39pub(crate) mod dbus;
40pub mod dbus_stats;
41pub mod json;
42pub mod logging;
43pub mod machines;
44pub mod networkd;
45pub mod pid1;
46pub mod system;
47pub mod timer;
48pub mod unit_constants;
49pub mod units;
50pub mod varlink;
51pub mod varlink_boot;
52pub mod varlink_fallback;
53pub mod varlink_networkd;
54pub mod varlink_system;
55pub mod varlink_unit;
56pub mod varlink_units;
57pub mod varlink_verify;
58pub mod verify;
59
60pub const DEFAULT_DBUS_ADDRESS: &str = "unix:path=/run/dbus/system_bus_socket";
61
62/// Per-collector timing for a single stat collection run.
63///
64/// `start_offset_ms` is the wall time between the top of the collection cycle and
65/// the moment this collector's future was first polled. A non-trivial offset
66/// indicates the spawn/scheduling loop or the runtime is delaying first poll,
67/// which means collectors are not starting in parallel as intended.
68///
69/// `elapsed_ms` is the wall time between first poll and completion.
70#[derive(serde::Serialize, serde::Deserialize, Clone, Debug, Default, PartialEq)]
71pub struct CollectorTiming {
72 /// Name of the collector (e.g. "units", "pid1", "dbus_stats")
73 pub name: String,
74 /// Milliseconds from top of the run until the spawned future's first poll.
75 /// Should be small (< a few ms) when collectors are truly running in parallel.
76 pub start_offset_ms: f64,
77 /// Milliseconds from first poll to future completion.
78 pub elapsed_ms: f64,
79 /// Whether the collector returned Ok.
80 pub success: bool,
81}
82
83/// Which API transport served a collector on the last run.
84///
85/// Recorded per enabled collector so operators can watch varlink adoption
86/// climb as the fleet's systemd upgrades past each endpoint's minimum
87/// version: a collector reports `Varlink` when its varlink attempt succeeded
88/// and `Dbus` when it collected the legacy way — either because varlink is
89/// disabled in config or because the varlink attempt failed and fell back.
90/// `Dbus` also covers the file-based networkd fallback: the value answers
91/// "did this collector use varlink", not "which fallback served it".
92/// Collectors with no varlink path (`pid1`, `dbus_stats`) and disabled
93/// collectors stay `None` and emit no gauge, so the present gauges are
94/// exactly the enabled set — no separate enabled-collectors counter needed.
95#[derive(
96 serde_repr::Serialize_repr,
97 serde_repr::Deserialize_repr,
98 Clone,
99 Copy,
100 Debug,
101 Default,
102 PartialEq,
103 Eq,
104)]
105#[repr(u8)]
106pub enum CollectorTransport {
107 #[default]
108 Dbus = 0,
109 Varlink = 1,
110}
111
112impl CollectorTransport {
113 /// Gauge value for flat JSON and metric exporters: 1 when varlink served
114 /// the collector, 0 for the D-Bus/file fallback.
115 ///
116 /// The `repr(u8)` discriminants are the gauge values, so every output
117 /// format (json, json-pretty, json-flat) reports the same integer —
118 /// matching the other stat enums (`SystemdSystemState`, the unit state
119 /// enums) rather than serializing as a string in some formats.
120 pub fn as_u64(self) -> u64 {
121 self as u64
122 }
123}
124
125/// Per-collector transport record for the last collection run.
126///
127/// Every field is `Some` when its collector ran and `None` when it did not
128/// (disabled in config). Absent-when-disabled is what makes the adoption
129/// ratio work: `count()` over the gauges is the enabled set.
130///
131/// NOTE for future collectors: the host `MachineStats` is constructed once
132/// outside the daemon loop, not fresh each iteration — so a collector with
133/// an early-return path that skips its gauge assignment would silently
134/// inherit the previous run's value instead of dropping the gauge. Always
135/// assign on every path, including cache hits and config-disabled fallbacks.
136#[derive(serde::Serialize, serde::Deserialize, Clone, Debug, Default, PartialEq)]
137pub struct VarlinkUsage {
138 /// Always collected; follows `[system-state] varlink` via the shared
139 /// `Manager.Describe` call. Normally agrees with `system_state` (one
140 /// call serves both), but persists when `[system-state]` is disabled —
141 /// making it the cheapest varlink canary on such hosts.
142 #[serde(skip_serializing_if = "Option::is_none")]
143 pub version: Option<CollectorTransport>,
144 #[serde(skip_serializing_if = "Option::is_none")]
145 pub system_state: Option<CollectorTransport>,
146 /// `Varlink` means the bulk enumeration came from the metrics stream and
147 /// per-unit details from `io.systemd.Unit.List`, on the host and inside
148 /// containers alike.
149 #[serde(skip_serializing_if = "Option::is_none")]
150 pub units: Option<CollectorTransport>,
151 #[serde(skip_serializing_if = "Option::is_none")]
152 pub networkd: Option<CollectorTransport>,
153 /// Host-side machine enumeration: machined's `io.systemd.Machine.List`
154 /// (systemd v257+) or its D-Bus API. Per-container transports land on
155 /// each machine's own `MachineStats` instead.
156 #[serde(skip_serializing_if = "Option::is_none")]
157 pub machines: Option<CollectorTransport>,
158 #[serde(skip_serializing_if = "Option::is_none")]
159 pub boot_blame: Option<CollectorTransport>,
160 #[serde(skip_serializing_if = "Option::is_none")]
161 pub verify: Option<CollectorTransport>,
162}
163
164/// Stats collected for a single systemd-nspawn container or VM managed by systemd-machined
165#[derive(serde::Serialize, serde::Deserialize, Clone, Debug, Default, PartialEq)]
166pub struct MachineStats {
167 /// systemd-networkd interface states inside the container
168 pub networkd: networkd::NetworkdState,
169 /// PID 1 process stats from procfs (using the container's leader PID)
170 pub pid1: Option<pid1::Pid1Stats>,
171 /// Overall systemd system state (e.g. running, degraded) inside the container
172 pub system_state: system::SystemdSystemState,
173 /// Aggregated systemd unit counts and per-service/timer stats inside the container
174 pub units: units::SystemdUnitStats,
175 /// systemd version running inside the container
176 pub version: system::SystemdVersion,
177 /// D-Bus daemon/broker statistics inside the container
178 pub dbus_stats: Option<dbus_stats::DBusStats>,
179 /// Boot blame statistics: slowest units at boot with activation times in seconds
180 #[serde(skip_serializing_if = "Option::is_none")]
181 pub boot_blame: Option<boot::BootBlameStats>,
182 /// Unit verification error statistics
183 pub verify_stats: Option<verify::VerifyStats>,
184 /// Which transport served each collector on the last run.
185 pub varlink_usage: VarlinkUsage,
186}
187
188/// Root struct containing all enabled monitord metrics for the host system and containers
189#[derive(serde::Serialize, serde::Deserialize, Debug, Default, PartialEq)]
190pub struct MonitordStats {
191 /// systemd-networkd interface states and managed interface count
192 pub networkd: networkd::NetworkdState,
193 /// PID 1 (systemd) process stats from procfs: CPU, memory, FDs, tasks
194 pub pid1: Option<pid1::Pid1Stats>,
195 /// Overall systemd manager state (e.g. running, degraded, initializing)
196 pub system_state: system::SystemdSystemState,
197 /// Aggregated systemd unit counts by type/state and per-service/timer detailed metrics
198 pub units: units::SystemdUnitStats,
199 /// Installed systemd version (major.minor.revision.os)
200 pub version: system::SystemdVersion,
201 /// D-Bus daemon/broker statistics (connections, bus names, match rules, per-peer accounting)
202 pub dbus_stats: Option<dbus_stats::DBusStats>,
203 /// Per-container stats keyed by machine name, collected via systemd-machined
204 pub machines: HashMap<String, MachineStats>,
205 /// Boot blame statistics: slowest units at boot with activation times in seconds
206 #[serde(skip_serializing_if = "Option::is_none")]
207 pub boot_blame: Option<boot::BootBlameStats>,
208 /// Unit verification error statistics
209 pub verify_stats: Option<verify::VerifyStats>,
210 /// End-to-end duration of the last stat collection run in milliseconds.
211 pub stat_collection_run_time_ms: f64,
212 /// Per-collector timings from the last run, sorted slowest first. Empty
213 /// before the first run completes. Callers compute parallelism ratio
214 /// (sum of `elapsed_ms` / `stat_collection_run_time_ms`) and identify the
215 /// gating collector (first entry) directly from this vector.
216 pub collector_timings: Vec<CollectorTiming>,
217 /// Which transport served each collector on the last run.
218 pub varlink_usage: VarlinkUsage,
219}
220
221/// Print statistics in the format set in configuration
222pub fn print_stats(
223 key_prefix: &str,
224 output_format: &config::MonitordOutputFormat,
225 stats: &MonitordStats,
226) {
227 match output_format {
228 config::MonitordOutputFormat::Json => println!(
229 "{}",
230 serde_json::to_string(&stats).expect("Invalid JSON serialization")
231 ),
232 config::MonitordOutputFormat::JsonFlat => println!(
233 "{}",
234 json::flatten(stats, key_prefix).expect("Invalid JSON serialization")
235 ),
236 config::MonitordOutputFormat::JsonPretty => println!(
237 "{}",
238 serde_json::to_string_pretty(&stats).expect("Invalid JSON serialization")
239 ),
240 }
241}
242
243fn set_stat_collection_run_time(stats: &mut MonitordStats, elapsed_runtime: Duration) {
244 stats.stat_collection_run_time_ms = elapsed_runtime.as_secs_f64() * 1000.0;
245}
246
247/// Output produced by every spawned collector future after wrapping with timing.
248type TimedCollectorOutput = (String, anyhow::Result<()>, Duration, Duration);
249
250/// Spawn a collector future onto the join set with timing instrumentation.
251///
252/// The wrapping closure records the moment the future is first polled (relative
253/// to `collect_start`) and the elapsed wall time until it completes. Both
254/// durations and the collector name are returned alongside the original result.
255///
256/// `tokio::task::JoinSet::spawn` runs the future on a new task, and tracing
257/// spans do not cross task boundaries automatically — without explicitly
258/// capturing and re-attaching the caller's span here, every collector (and
259/// anything it spawns in turn, e.g. units.rs's per-unit tasks) would show up
260/// as an unrelated root trace instead of a child of the current collection run.
261fn spawn_timed<F>(
262 join_set: &mut tokio::task::JoinSet<TimedCollectorOutput>,
263 name: &'static str,
264 collect_start: Instant,
265 fut: F,
266) where
267 F: std::future::Future<Output = anyhow::Result<()>> + Send + 'static,
268{
269 let parent_span = tracing::Span::current();
270 let span = tracing::debug_span!(
271 parent: &parent_span,
272 "collector",
273 name,
274 elapsed_ms = tracing::field::Empty,
275 success = tracing::field::Empty,
276 );
277 let recording_span = span.clone();
278 join_set.spawn(
279 async move {
280 let task_first_poll = Instant::now();
281 let start_offset = task_first_poll.duration_since(collect_start);
282 let result = fut.await;
283 let elapsed = task_first_poll.elapsed();
284 recording_span.record("elapsed_ms", elapsed.as_secs_f64() * 1000.0);
285 recording_span.record("success", result.is_ok());
286 (name.to_string(), result, start_offset, elapsed)
287 }
288 .instrument(span),
289 );
290}
291
292/// Lazily-created system bus connection shared by every collector that
293/// needs D-Bus — also passed into collectors that decide their transport
294/// internally (`verify`, `boot_blame`), so those resolve it only on their
295/// D-Bus code paths rather than up front (see finding 1 on #224).
296///
297/// The cell starts empty (or seeded from `maybe_connection` for library
298/// callers) and is only connected on first use, so a run whose enabled
299/// collectors all succeed over varlink/fs/procfs never touches the bus at
300/// all — the prerequisite for the zero-D-Bus claim in #37.
301///
302/// `tokio::sync::OnceCell` serializes initializers behind a semaphore:
303/// concurrent first-users wait for the in-flight attempt and share its
304/// result — there is no race and no spare connection. A failed attempt is
305/// *not* cached (a waiter starts a fresh attempt instead), so a dead bus
306/// costs one connect attempt per D-Bus collector, sequentially. `ENOENT`
307/// fails fast and the run finishes in milliseconds, but `method_timeout`
308/// only bounds method calls, not `Builder::build()` — a socket that
309/// exists and hangs (wedged broker) would serialize one blocking connect
310/// per D-Bus collector per cycle, where the old code blocked exactly once
311/// at startup. Worth knowing when pointing at an unfamiliar bus address.
312pub(crate) type DbusCell = Arc<tokio::sync::OnceCell<zbus::Connection>>;
313
314/// Resolve the shared D-Bus connection, connecting on first use.
315///
316/// Only called on D-Bus code paths (varlink/fs-first collectors reach here
317/// solely through their fallbacks), so a missing bus surfaces as that
318/// collector's ordinary error — logged per-collector, never fatal to the
319/// run. Daemon-mode reuse is preserved: the same connection serves every
320/// cycle until a failed cycle drops it (see below).
321pub(crate) async fn dbus_connection(
322 cell: &DbusCell,
323 dbus_timeout: u64,
324) -> Result<zbus::Connection, MonitordError> {
325 cell.get_or_try_init(|| async {
326 zbus::connection::Builder::system()?
327 .method_timeout(std::time::Duration::from_secs(dbus_timeout))
328 .build()
329 .await
330 })
331 .await
332 .cloned()
333 .map_err(MonitordError::ZbusError)
334}
335
336/// Main statistic collection function running what's required by configuration in parallel
337/// Takes an optional locked stats struct to update and to output stats to STDOUT or not.
338/// Takes an optional D-Bus connection. Returns `Some(connection)` if the
339/// collection cycle completed without errors (meaning the connection is reusable),
340/// `None` if errors occurred.
341pub async fn stat_collector(
342 config: config::Config,
343 maybe_locked_stats: Option<Arc<RwLock<MonitordStats>>>,
344 output_stats: bool,
345 maybe_connection: Option<zbus::Connection>,
346) -> Result<Option<zbus::Connection>, MonitordError> {
347 let mut collect_interval_ms: u128 = 0;
348 if config.monitord.daemon {
349 collect_interval_ms = (config.monitord.daemon_stats_refresh_secs * 1000).into();
350 }
351
352 let config = Arc::new(config);
353 let locked_monitord_stats: Arc<RwLock<MonitordStats>> =
354 maybe_locked_stats.unwrap_or(Arc::new(RwLock::new(MonitordStats::default())));
355 let locked_machine_stats: Arc<RwLock<MachineStats>> =
356 Arc::new(RwLock::new(MachineStats::default()));
357 let cached_machine_connections: Arc<tokio::sync::Mutex<machines::MachineConnections>> =
358 Arc::new(tokio::sync::Mutex::new(HashMap::new()));
359
360 std::env::set_var("DBUS_SYSTEM_BUS_ADDRESS", &config.monitord.dbus_address);
361 // Seeded when the library caller passes a connection, otherwise
362 // connected lazily by the first collector that needs the bus. Mutable
363 // so a failed daemon cycle can swap in a fresh cell (see below) rather
364 // than letting the next cycle reuse a bus that may have gone away.
365 let mut dbus_cell: DbusCell = Arc::new(tokio::sync::OnceCell::new());
366 if let Some(conn) = maybe_connection {
367 // Seeding cannot fail: a fresh cell is always empty.
368 let _ = dbus_cell.set(conn);
369 }
370 let dbus_timeout = config.monitord.dbus_timeout;
371 let mut join_set: tokio::task::JoinSet<TimedCollectorOutput> = tokio::task::JoinSet::new();
372 let mut had_error;
373
374 loop {
375 let collect_start_time = Instant::now();
376 // Kept alive for the whole iteration so its reported duration covers
377 // the full run (spawn + drain), not just the synchronous spawn phase
378 // below where `run_guard` is held.
379 let run_span = tracing::info_span!("stat_collector_run");
380 let run_guard = run_span.enter();
381 info!("Starting stat collection run");
382
383 // One Manager.Describe serves both the version and system state
384 // collectors: PID 1 handles varlink requests one at a time, so a second
385 // call would sit behind the units collector's whole metrics stream.
386 //
387 // Which request reaches PID 1 first is still up to the scheduler, so
388 // under CPU pressure this call can land behind that stream anyway. That
389 // is the old cost for one call rather than two, not a new failure mode.
390 // Making it deterministic would mean holding the units collector until
391 // this resolves, which would put the collector that gates the cycle
392 // behind an unrelated socket.
393 let manager_describe = config.use_varlink(&[config.system_state.varlink]).then(|| {
394 crate::varlink_system::shared_describe(
395 crate::varlink_system::MANAGER_SOCKET_PATH.into(),
396 )
397 });
398
399 // Always collect systemd version
400 {
401 let dbus_cell = Arc::clone(&dbus_cell);
402 let stats_clone = locked_machine_stats.clone();
403 let describe = manager_describe.clone();
404 let no_fallback = config.varlink.no_fallback;
405 spawn_timed(&mut join_set, "version", collect_start_time, async move {
406 if let Some(describe) = describe {
407 match crate::varlink_system::update_version(describe, stats_clone.clone()).await
408 {
409 Ok(()) => {
410 stats_clone.write().await.varlink_usage.version =
411 Some(CollectorTransport::Varlink);
412 return Ok(());
413 }
414 Err(err) => {
415 crate::varlink_fallback::report_varlink_failure(
416 no_fallback,
417 "version",
418 "D-Bus",
419 err,
420 )?;
421 }
422 }
423 }
424 stats_clone.write().await.varlink_usage.version = Some(CollectorTransport::Dbus);
425 let conn = dbus_connection(&dbus_cell, dbus_timeout).await?;
426 crate::system::update_version(conn, stats_clone.clone()).await
427 });
428 }
429
430 // Collect pid1 procfs stats
431 if config.pid1.enabled {
432 spawn_timed(
433 &mut join_set,
434 "pid1",
435 collect_start_time,
436 crate::pid1::update_pid1_stats(1, locked_machine_stats.clone()),
437 );
438 }
439
440 // Run networkd collector if enabled
441 if config.networkd.enabled {
442 let config_clone = Arc::clone(&config);
443 let dbus_cell = Arc::clone(&dbus_cell);
444 let stats_clone = locked_machine_stats.clone();
445 let no_fallback = config.varlink.no_fallback;
446 spawn_timed(&mut join_set, "networkd", collect_start_time, async move {
447 if config_clone.use_varlink(&[config_clone.networkd.varlink]) {
448 let endpoint = crate::varlink_networkd::NETWORK_SOCKET_PATH.into();
449 match crate::varlink_networkd::get_networkd_state(&endpoint).await {
450 Ok(networkd_stats) => {
451 let mut machine_stats = stats_clone.write().await;
452 machine_stats.networkd = networkd_stats;
453 machine_stats.varlink_usage.networkd =
454 Some(CollectorTransport::Varlink);
455 return Ok(());
456 }
457 Err(err) => {
458 crate::varlink_fallback::report_varlink_failure(
459 no_fallback,
460 "networkd",
461 "file-based",
462 err,
463 )?;
464 }
465 }
466 }
467 stats_clone.write().await.varlink_usage.networkd = Some(CollectorTransport::Dbus);
468 // No eager connect: the cell flows into the collector and
469 // is resolved only if sysfs yields no usable map — the
470 // same lazy pattern as boot_blame.
471 crate::networkd::update_networkd_stats(
472 config_clone.networkd.link_state_dir.clone(),
473 None,
474 std::path::PathBuf::from("/sys"),
475 Some((dbus_cell, dbus_timeout)),
476 stats_clone,
477 )
478 .await
479 });
480 }
481
482 // Run system running (SystemState) state collector
483 if config.system_state.enabled {
484 let dbus_cell = Arc::clone(&dbus_cell);
485 let stats_clone = locked_machine_stats.clone();
486 let describe = manager_describe.clone();
487 let no_fallback = config.varlink.no_fallback;
488 spawn_timed(
489 &mut join_set,
490 "system_state",
491 collect_start_time,
492 async move {
493 if let Some(describe) = describe {
494 match crate::varlink_system::update_system_stats(
495 describe,
496 stats_clone.clone(),
497 )
498 .await
499 {
500 Ok(()) => {
501 stats_clone.write().await.varlink_usage.system_state =
502 Some(CollectorTransport::Varlink);
503 return Ok(());
504 }
505 Err(err) => {
506 crate::varlink_fallback::report_varlink_failure(
507 no_fallback,
508 "system state",
509 "D-Bus",
510 err,
511 )?;
512 }
513 }
514 }
515 stats_clone.write().await.varlink_usage.system_state =
516 Some(CollectorTransport::Dbus);
517 let conn = dbus_connection(&dbus_cell, dbus_timeout).await?;
518 crate::system::update_system_stats(conn, stats_clone.clone()).await
519 },
520 );
521 }
522
523 // Run service collectors if there are services listed in config
524 if config.units.enabled {
525 let config_clone = Arc::clone(&config);
526 let dbus_cell = Arc::clone(&dbus_cell);
527 let stats_clone = locked_machine_stats.clone();
528 let no_fallback = config.varlink.no_fallback;
529 spawn_timed(&mut join_set, "units", collect_start_time, async move {
530 if config_clone.use_varlink(&[config_clone.units.varlink]) {
531 match crate::varlink_units::update_unit_stats(
532 Arc::clone(&config_clone),
533 stats_clone.clone(),
534 crate::varlink_units::METRICS_SOCKET_PATH.into(),
535 )
536 .await
537 {
538 Ok(timer_names) => {
539 // Per-service stats, timer properties and service
540 // types all come from io.systemd.Unit.List. If that
541 // socket is unusable the whole units collection is
542 // redone over D-Bus: restoring pieces of it would
543 // leave [services] and timers partially filled from
544 // the metrics, which is worse than either path alone.
545 if let Err(err) = crate::varlink_units::apply_unit_details(
546 &crate::varlink_unit::MANAGER_SOCKET_PATH.into(),
547 &stats_clone,
548 &config_clone,
549 "",
550 &timer_names,
551 )
552 .await
553 {
554 crate::varlink_fallback::report_varlink_failure(
555 no_fallback,
556 "unit details",
557 "D-Bus",
558 err,
559 )?;
560 stats_clone.write().await.varlink_usage.units =
561 Some(CollectorTransport::Dbus);
562 let conn = dbus_connection(&dbus_cell, dbus_timeout).await?;
563 return crate::units::update_unit_stats(
564 config_clone,
565 conn,
566 stats_clone,
567 String::new(),
568 )
569 .await;
570 }
571 if config_clone.units.unit_files {
572 let unit_files = crate::units::collect_unit_files_stats("").await;
573 let mut ms = stats_clone.write().await;
574 ms.units.unit_files = unit_files;
575 }
576 stats_clone.write().await.varlink_usage.units =
577 Some(CollectorTransport::Varlink);
578 return Ok(());
579 }
580 Err(err) => {
581 crate::varlink_fallback::report_varlink_failure(
582 no_fallback,
583 "units",
584 "D-Bus",
585 err,
586 )?;
587 }
588 }
589 }
590 stats_clone.write().await.varlink_usage.units = Some(CollectorTransport::Dbus);
591 let conn = dbus_connection(&dbus_cell, dbus_timeout).await?;
592 crate::units::update_unit_stats(config_clone, conn, stats_clone, String::new())
593 .await
594 });
595 }
596
597 if config.machines.enabled {
598 let stats_clone = locked_machine_stats.clone();
599 let config_clone = Arc::clone(&config);
600 let dbus_cell = Arc::clone(&dbus_cell);
601 let monitord_stats_clone = locked_monitord_stats.clone();
602 let connections_clone = cached_machine_connections.clone();
603 spawn_timed(&mut join_set, "machines", collect_start_time, async move {
604 // Enumeration records its transport (varlink or D-Bus) on
605 // the host stats; per-container collection records its own
606 // transports on each machine's stats.
607 crate::machines::update_machines_stats(
608 config_clone,
609 dbus_cell,
610 stats_clone,
611 monitord_stats_clone,
612 connections_clone,
613 )
614 .await
615 });
616 }
617
618 if config.dbus_stats.enabled {
619 let dbus_cell = Arc::clone(&dbus_cell);
620 let config_clone = Arc::clone(&config);
621 let stats_clone = locked_machine_stats.clone();
622 spawn_timed(
623 &mut join_set,
624 "dbus_stats",
625 collect_start_time,
626 async move {
627 let conn = dbus_connection(&dbus_cell, dbus_timeout).await?;
628 crate::dbus_stats::update_dbus_stats(config_clone, conn, stats_clone).await
629 },
630 );
631 }
632
633 // `verify` and `boot_blame` decide their transport internally:
634 // they take the shared cell and resolve it only on a D-Bus code
635 // path, so a varlink success or cache hit never connects.
636 if config.boot_blame.enabled {
637 let dbus_cell = Arc::clone(&dbus_cell);
638 let config_clone = Arc::clone(&config);
639 let stats_clone = locked_machine_stats.clone();
640 spawn_timed(
641 &mut join_set,
642 "boot_blame",
643 collect_start_time,
644 async move {
645 let no_fallback = config_clone.varlink.no_fallback;
646 crate::boot::update_boot_blame_stats(
647 config_clone,
648 dbus_cell,
649 dbus_timeout,
650 stats_clone,
651 no_fallback,
652 )
653 .await
654 },
655 );
656 }
657
658 if config.verify.enabled {
659 let dbus_cell = Arc::clone(&dbus_cell);
660 let config_clone = Arc::clone(&config);
661 let stats_clone = locked_machine_stats.clone();
662 spawn_timed(&mut join_set, "verify", collect_start_time, async move {
663 crate::verify::update_verify_stats(
664 dbus_cell,
665 dbus_timeout,
666 stats_clone,
667 config_clone.verify.allowlist.clone(),
668 config_clone.verify.blocklist.clone(),
669 config_clone.use_varlink(&[config_clone.verify.varlink]),
670 config_clone.varlink.no_fallback,
671 )
672 .await
673 });
674 }
675
676 if join_set.len() == 1 {
677 warn!("No collectors except systemd version scheduled to run. Exiting");
678 }
679
680 // All collectors above were spawned with `run_span` captured as their
681 // parent; the guard must be dropped before the first `.await` below
682 // since span guards are not valid to hold across an await point.
683 drop(run_guard);
684
685 // Drain join_set, collect per-collector timings + log per-collector failures
686 had_error = false;
687 let mut timings: Vec<CollectorTiming> = Vec::new();
688 while let Some(res) = join_set.join_next().await {
689 match res {
690 Ok((name, collector_result, start_offset, elapsed)) => {
691 let success = collector_result.is_ok();
692 if let Err(e) = collector_result {
693 had_error = true;
694 error!("Collector '{}' failure: {:?}", name, e);
695 }
696 timings.push(CollectorTiming {
697 name,
698 start_offset_ms: start_offset.as_secs_f64() * 1000.0,
699 elapsed_ms: elapsed.as_secs_f64() * 1000.0,
700 success,
701 });
702 }
703 Err(e) => {
704 had_error = true;
705 error!("Join error: {:?}", e);
706 }
707 }
708 }
709
710 let elapsed_runtime = collect_start_time.elapsed();
711 let elapsed_runtime_ms = elapsed_runtime.as_millis();
712
713 // Sort timings by elapsed desc so the slowest collector is first in the JSON output
714 timings.sort_by(|a, b| {
715 b.elapsed_ms
716 .partial_cmp(&a.elapsed_ms)
717 .unwrap_or(std::cmp::Ordering::Equal)
718 });
719
720 // Per-collector lines log at debug! to keep daemon-mode noise low.
721 // The same data is on MonitordStats::collector_timings for callers that need it.
722 for t in &timings {
723 debug!(
724 "collector '{}' start_offset={:.1}ms elapsed={:.1}ms{}",
725 t.name,
726 t.start_offset_ms,
727 t.elapsed_ms,
728 if t.success { "" } else { " (FAILED)" },
729 );
730 }
731
732 {
733 // Update monitord stats with machine stats
734 let mut monitord_stats = locked_monitord_stats.write().await;
735 let machine_stats = locked_machine_stats.read().await;
736 monitord_stats.pid1 = machine_stats.pid1.clone();
737 monitord_stats.networkd = machine_stats.networkd.clone();
738 monitord_stats.system_state = machine_stats.system_state;
739 monitord_stats.version = machine_stats.version.clone();
740 monitord_stats.units = machine_stats.units.clone();
741 monitord_stats.dbus_stats = machine_stats.dbus_stats.clone();
742 monitord_stats.boot_blame = machine_stats.boot_blame.clone();
743 monitord_stats.verify_stats = machine_stats.verify_stats.clone();
744 monitord_stats.varlink_usage = machine_stats.varlink_usage.clone();
745 set_stat_collection_run_time(&mut monitord_stats, elapsed_runtime);
746 monitord_stats.collector_timings = timings;
747 }
748
749 info!("stat collection run took {}ms", elapsed_runtime_ms);
750 if output_stats {
751 let monitord_stats = locked_monitord_stats.read().await;
752 print_stats(
753 &config.monitord.key_prefix,
754 &config.monitord.output_format,
755 &monitord_stats,
756 );
757 }
758 if !config.monitord.daemon {
759 break;
760 }
761 if had_error {
762 // Drop the shared connection so the next cycle reconnects
763 // rather than reusing a bus that may have gone away mid-run.
764 // Each iteration clones the cell into its tasks at spawn time,
765 // so reassigning the local is enough for the next cycle's
766 // spawns to pick up the fresh cell.
767 dbus_cell = Arc::new(tokio::sync::OnceCell::new());
768 }
769 let sleep_time_ms = collect_interval_ms - elapsed_runtime_ms;
770 info!("stat collection sleeping for {}s 😴", sleep_time_ms / 1000);
771 tokio::time::sleep(Duration::from_millis(
772 sleep_time_ms
773 .try_into()
774 .expect("Sleep time does not fit into a u64 :O"),
775 ))
776 .await;
777 }
778 // A failed cycle returns no connection so a repeated one-shot caller
779 // reconnects instead of reusing a broken bus; a clean cycle hands the
780 // shared connection back for reuse. (In daemon mode the reset above
781 // already swapped in a fresh cell, so `get()` here is empty by design
782 // and `None` is the only honest answer.)
783 let conn = if had_error {
784 None
785 } else {
786 dbus_cell.get().cloned()
787 };
788 Ok(conn)
789}
790
791#[cfg(test)]
792mod tests {
793 use super::*;
794
795 /// The lazy-connect contract finding 6 on #224 pins down: a failed
796 /// connect is not cached, so consecutive failures each attempt (and
797 /// report) a fresh error rather than inheriting a poisoned cell.
798 #[tokio::test]
799 async fn test_dbus_connection_failure_is_not_cached() {
800 std::env::set_var(
801 "DBUS_SYSTEM_BUS_ADDRESS",
802 "unix:path=/nonexistent/monitord-test-bus-socket",
803 );
804 let cell: DbusCell = Arc::new(tokio::sync::OnceCell::new());
805 assert!(dbus_connection(&cell, 1).await.is_err());
806 assert!(
807 cell.get().is_none(),
808 "failed connect must not populate the cell"
809 );
810 assert!(dbus_connection(&cell, 1).await.is_err());
811 assert!(cell.get().is_none());
812 }
813
814 #[test]
815 fn test_stat_collection_run_time_ms_conversion() {
816 let mut stats = MonitordStats::default();
817 set_stat_collection_run_time(&mut stats, Duration::from_millis(5));
818 assert_eq!(stats.stat_collection_run_time_ms, 5.0);
819
820 set_stat_collection_run_time(&mut stats, Duration::from_micros(500));
821 assert!((stats.stat_collection_run_time_ms - 0.5).abs() < f64::EPSILON);
822 }
823}