Skip to main content

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}