fix: bound monitoring subprocesses and collect fresh system readings

This commit is contained in:
archipelago
2026-10-05 21:31:32 -04:00
parent 3d0c67eb9b
commit e97f958f45
4 changed files with 95 additions and 19 deletions
+55 -13
View File
@@ -13,11 +13,11 @@ pub async fn collect_snapshot() -> Result<MetricSnapshot> {
read_loadavg(), read_loadavg(),
); );
let cpu = cpu.unwrap_or(0.0); let cpu = cpu?;
let (mem_used, mem_total) = mem.unwrap_or((0, 0)); let (mem_used, mem_total) = mem?;
let (disk_used, disk_total) = disk.unwrap_or((0, 0)); let (disk_used, disk_total) = disk?;
let (net_rx, net_tx) = net.unwrap_or((0, 0)); let (net_rx, net_tx) = net?;
let (l1, l5, l15) = load.unwrap_or((0.0, 0.0, 0.0)); let (l1, l5, l15) = load?;
let system = SystemMetrics { let system = SystemMetrics {
cpu_percent: cpu, cpu_percent: cpu,
@@ -120,10 +120,9 @@ async fn read_disk_usage() -> Result<(u64, u64)> {
} else { } else {
"/" "/"
}; };
let output = tokio::process::Command::new("df") let mut command = tokio::process::Command::new("df");
.args(["--block-size=1", "--output=used,size", target]) command.args(["--block-size=1", "--output=used,size", target]);
.output() let output = bounded_output(command, std::time::Duration::from_secs(3)).await
.await
.context("Failed to run df")?; .context("Failed to run df")?;
if !output.status.success() { if !output.status.success() {
@@ -215,12 +214,23 @@ async fn read_network_totals() -> Result<(u64, u64)> {
Ok((rx_total, tx_total)) Ok((rx_total, tx_total))
} }
/// A wedged runtime or filesystem must not freeze every monitoring snapshot.
/// Dropping a timed-out child kills it, so repeated polls cannot leak processes.
async fn bounded_output(
mut command: tokio::process::Command,
timeout: std::time::Duration,
) -> Result<std::process::Output> {
command.kill_on_drop(true);
tokio::time::timeout(timeout, command.output())
.await.context("Metrics subprocess timed out")?
.context("Metrics subprocess failed")
}
/// Get per-container resource stats via `podman stats --no-stream --format json`. /// Get per-container resource stats via `podman stats --no-stream --format json`.
async fn read_container_stats() -> Result<Vec<ContainerMetrics>> { async fn read_container_stats() -> Result<Vec<ContainerMetrics>> {
let output = tokio::process::Command::new("podman") let mut command = tokio::process::Command::new("podman");
.args(["stats", "--no-stream", "--format", "json"]) command.args(["stats", "--no-stream", "--format", "json"]);
.output() let output = bounded_output(command, std::time::Duration::from_secs(8)).await
.await
.context("Failed to run podman stats")?; .context("Failed to run podman stats")?;
if !output.status.success() { if !output.status.success() {
@@ -391,3 +401,35 @@ mod tests {
assert_eq!(parse_bytes_field(&obj, "mem"), Some(268435456)); assert_eq!(parse_bytes_field(&obj, "mem"), Some(268435456));
} }
} }
#[cfg(test)]
mod subprocess_deadline_tests {
use super::*;
#[tokio::test]
async fn a_stalled_metrics_command_is_bounded_and_killed() {
let dir = tempfile::tempdir().unwrap();
let pid_file = dir.path().join("pid");
let mut command = tokio::process::Command::new("sh");
command.arg("-c").arg("echo $$ > \"$1\"; exec sleep 30").arg("metrics-test").arg(&pid_file);
let start = std::time::Instant::now();
let error = bounded_output(command, std::time::Duration::from_millis(500)).await.unwrap_err();
assert!(error.to_string().contains("timed out"));
assert!(start.elapsed() < std::time::Duration::from_secs(3));
let pid = tokio::fs::read_to_string(pid_file).await.unwrap();
for _ in 0..40 {
if !std::path::Path::new(&format!("/proc/{}", pid.trim())).exists() { return; }
tokio::time::sleep(std::time::Duration::from_millis(25)).await;
}
panic!("Timed-out metrics subprocess was not reaped");
}
#[tokio::test]
async fn successful_metrics_output_is_preserved() {
let mut command = tokio::process::Command::new("printf");
command.arg("metrics-ok");
let output = bounded_output(command, std::time::Duration::from_secs(1)).await.unwrap();
assert!(output.status.success());
assert_eq!(output.stdout, b"metrics-ok");
}
}
+6 -6
View File
@@ -14,21 +14,21 @@ use std::path::PathBuf;
use std::sync::Arc; use std::sync::Arc;
use tracing::{debug, warn}; use tracing::{debug, warn};
/// Spawn the background metrics collector (runs every 300 seconds / 5 minutes). /// Spawn the background metrics collector at the store's one-minute resolution.
/// Evaluates alert rules on each snapshot and dispatches notifications. /// Evaluates alert rules on each snapshot and dispatches notifications.
/// Note: health_monitor.rs handles container state polling at 120s intervals. /// Note: health_monitor.rs handles container state polling at 120s intervals.
/// This collector handles system-level metrics (CPU, disk, network) and only /// Runtime commands have deadlines; unavailable container stats cannot hold
/// calls podman stats every 5 minutes to avoid duplicate subprocess overhead. /// system readings indefinitely. Missed ticks are skipped, never replayed.
pub fn spawn_metrics_collector( pub fn spawn_metrics_collector(
store: Arc<MetricsStore>, store: Arc<MetricsStore>,
state: Option<Arc<crate::state::StateManager>>, state: Option<Arc<crate::state::StateManager>>,
data_dir: Option<PathBuf>, data_dir: Option<PathBuf>,
) { ) {
tokio::spawn(async move { tokio::spawn(async move {
// Wait 60s for system to stabilize after boot // Start promptly without competing with the very first boot tasks.
tokio::time::sleep(std::time::Duration::from_secs(60)).await; tokio::time::sleep(std::time::Duration::from_secs(5)).await;
let mut interval = tokio::time::interval(std::time::Duration::from_secs(300)); let mut interval = tokio::time::interval(std::time::Duration::from_secs(60));
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop { loop {
+23
View File
@@ -38,3 +38,26 @@ mobile/desktop browser layout; upgraded sender and receiver comparing local
Monitoring to received Fleet values; restart/reconnect, offline ageing and real Monitoring to received Fleet values; restart/reconnect, offline ageing and real
transport evidence. Legacy senders still advertising zeros need the sender fix; transport evidence. Legacy senders still advertising zeros need the sender fix;
the receiver cannot reliably distinguish a legacy fabricated zero from idle CPU. the receiver cannot reliably distinguish a legacy fabricated zero from idle CPU.
## Dev deployment and collector follow-up
Dev qualification deployment uses backend041f1fa2 (SHA256
11b62f697a71b762bf8638063d1a858a68eea6d7268e300f6320c99068efa78a)
and frontend3d0c67eb. Health and unchanged app-container identities/start times
pass. Live federation CPU/memory/disk values exactly match local Monitoring;
Web5 mobile390px/desktop1440px checks pass. Rollback retained on the node at
/var/lib/archipelago/support/followup-20261005-2120. Yaya remains on the prior
backend; no receiver/fleet-wide acceptance is claimed.
Live testing caught an unbounded podman stats subprocess delaying the first
snapshot, and a300-second collection interval conflicting with180-second Fleet
freshness. The next candidate starts after5seconds and collects once per minute,
with3-second df and8-second podman deadlines and kill-on-drop cleanup. Failed
system reads no longer become invented zero samples. Container-stat failure
leaves container readings unavailable while retaining valid system readings.
The actual subprocess timeout/reaping regression passes; all1,686 backend tests
pass (four ignored). This collector correction is not deployed yet.
Framework SSH and backend health work, but its stored dashboard session returns
401. Authenticated Framework Monitoring acceptance remains open. No authentication
boundary was bypassed to produce an apparent pass.
@@ -24,8 +24,16 @@ type ProviderWindow = Window & {
describe('nostr-provider identity selection', () => { describe('nostr-provider identity selection', () => {
let providerWindow: ProviderWindow let providerWindow: ProviderWindow
let listeners: Array<[string, EventListenerOrEventListenerObject, boolean | AddEventListenerOptions | undefined]>
beforeEach(() => { beforeEach(() => {
vi.useFakeTimers()
listeners = []
const add = window.addEventListener.bind(window)
vi.spyOn(window, 'addEventListener').mockImplementation((type, listener, options) => {
if (listener) listeners.push([type, listener as EventListenerOrEventListenerObject, options])
add(type, listener, options)
})
sessionStorage.clear() sessionStorage.clear()
providerWindow = window as ProviderWindow providerWindow = window as ProviderWindow
delete providerWindow.__archipelagoNostr delete providerWindow.__archipelagoNostr
@@ -37,6 +45,9 @@ describe('nostr-provider identity selection', () => {
}) })
afterEach(() => { afterEach(() => {
for (const [type, listener, options] of listeners) window.removeEventListener(type, listener, options)
vi.clearAllTimers()
vi.unstubAllGlobals()
vi.useRealTimers() vi.useRealTimers()
vi.restoreAllMocks() vi.restoreAllMocks()
Reflect.deleteProperty(document, 'readyState') Reflect.deleteProperty(document, 'readyState')