use std::any::Any;
use std::panic::AssertUnwindSafe;
use std::time::{Duration, Instant};
use futures_util::FutureExt;
use tracing::{debug, error, info};
use trusty_common::host_metrics::HostSampler;
use crate::server::AppState;
use crate::service_metrics::{ServiceMetricsSampler, resolve_pid};
fn panic_text(payload: &(dyn Any + Send)) -> String {
if let Some(s) = payload.downcast_ref::<&'static str>() {
(*s).to_string()
} else if let Some(s) = payload.downcast_ref::<String>() {
s.clone()
} else {
"<non-string panic payload>".to_string()
}
}
async fn guarded<F: Future<Output = ()>>(step: &'static str, body: F) -> bool {
match AssertUnwindSafe(body).catch_unwind().await {
Ok(()) => true,
Err(payload) => {
error!(
step,
panic = %panic_text(payload.as_ref()),
"machine_history: a sampler step panicked; the loop continues and this tick \
records nothing for that step"
);
false
}
}
}
async fn sample_services(
state: &AppState,
service_metrics: &mut ServiceMetricsSampler,
services: &[crate::connector::ServiceInfo],
) {
let pending = service_metrics.pending_lookups(services, Instant::now());
if !pending.is_empty() {
let found = tokio::task::spawn_blocking(move || {
pending
.into_iter()
.map(|id| {
let pid = resolve_pid(&id);
(id, pid)
})
.collect::<Vec<_>>()
})
.await
.unwrap_or_else(|e| {
debug!("machine_history: service pid lookup task failed: {e}");
Vec::new()
});
service_metrics.record_lookups(found, Instant::now());
}
let now_unix = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |d| d.as_secs());
let batch = service_metrics.sample(services, now_unix);
state.machine_history().record_service_samples(batch).await;
}
pub fn start(state: AppState, interval: Duration) {
start_with_disk_interval(
state,
interval,
Duration::from_secs(trusty_common::host_metrics::history::DISK_SAMPLE_INTERVAL_SECS),
);
}
pub fn start_with_disk_interval(state: AppState, interval: Duration, disk_interval: Duration) {
state
.machine_history()
.set_sample_interval(interval.as_secs());
let loop_task = tokio::spawn(async move {
info!(
"machine_history: sampling host metrics every {}s (disks every {}s, window={} points)",
interval.as_secs(),
disk_interval.as_secs(),
state.machine_history().snapshot().await.sample_capacity
);
let mut sampler = HostSampler::with_thresholds_and_disk_interval(
trusty_common::host_metrics::HostThresholds::default(),
disk_interval,
);
let mut service_metrics = ServiceMetricsSampler::new();
loop {
guarded("host", host_tick(&state, &mut sampler)).await;
guarded("services", service_tick(&state, &mut service_metrics)).await;
tokio::time::sleep(interval).await;
}
});
tokio::spawn(async move {
match loop_task.await {
Ok(()) => error!("machine_history: the sampler loop ended; history has stopped"),
Err(e) => error!(
error = %e,
"machine_history: the sampler task ended abnormally; history has stopped"
),
}
});
}
async fn host_tick(state: &AppState, sampler: &mut HostSampler) {
let metrics = sampler.sample();
debug!(
overall = ?metrics.overall_pressure,
cpu_pct = metrics.cpu.usage_pct,
mem_pct = metrics.memory.usage_pct,
"machine_history: sampled host metrics"
);
state.host_metrics_cache().set(metrics.clone()).await;
state.machine_history().record_sample(metrics).await;
}
async fn service_tick(state: &AppState, service_metrics: &mut ServiceMetricsSampler) {
let reports = state.collect_service_reports().await;
for change in state.machine_history().observe_services(&reports).await {
info!(
service = %change.service_id,
from = ?change.from,
to = ?change.to,
"machine_history: service changed state"
);
}
if let Some(snapshot) = state.poller_cache().snapshot().await {
sample_services(state, service_metrics, &snapshot.services).await;
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::connector::ServiceStatus;
use crate::machine_history::MachineHistory;
use crate::machine_history::service_samples::{ServiceSample, ServiceSampleBatch};
fn stub_sample(seq: u64) -> trusty_common::host_metrics::HostMetrics {
use trusty_common::host_metrics::{
CpuMetrics, DiskMetrics, HostMetrics, MemoryMetrics, NetworkMetrics, Pressure,
};
HostMetrics {
cpu: CpuMetrics {
usage_pct: 0.0,
logical_cores: 1,
physical_cores: None,
pressure: Pressure::Nominal,
},
memory: MemoryMetrics {
total_bytes: 1,
used_bytes: 0,
available_bytes: 1,
usage_pct: 0.0,
swap_total_bytes: 0,
swap_used_bytes: 0,
pressure: Pressure::Nominal,
},
disks: DiskMetrics {
aggregate_total_bytes: 1,
aggregate_available_bytes: 1,
aggregate_used_bytes: 0,
aggregate_usage_pct: 0.0,
pressure: Pressure::Nominal,
mounts: Vec::new(),
},
network: NetworkMetrics {
rx_bytes_per_sec: 0.0,
tx_bytes_per_sec: 0.0,
rx_total_bytes: 0,
tx_total_bytes: 0,
window_secs: 1.0,
},
overall_pressure: Pressure::Nominal,
sampled_at_unix: Some(seq),
}
}
#[tokio::test]
async fn a_panicking_step_is_contained_and_the_next_one_runs() {
let history = MachineHistory::new();
let survived = guarded("panicking", async { panic!("sampler exploded") }).await;
assert!(
!survived,
"a panicking step reports that it did not complete"
);
let survived = guarded("recording", async {
history.record_sample(stub_sample(1)).await;
})
.await;
assert!(survived, "the step after a panic still runs");
assert_eq!(
history.snapshot().await.samples.len(),
1,
"the loop kept recording after a panicking step"
);
assert_eq!(panic_text(&"static str payload"), "static str payload");
assert_eq!(panic_text(&"owned payload".to_string()), "owned payload");
assert_eq!(panic_text(&7u8), "<non-string panic payload>");
}
#[tokio::test]
async fn a_panicking_service_half_leaves_the_host_half_recording() {
let history = MachineHistory::new();
for seq in 1..=3u64 {
guarded("host", async {
history.record_sample(stub_sample(seq)).await;
})
.await;
guarded("services", async {
panic!("service sampling exploded on tick {seq}")
})
.await;
let snap = history.snapshot().await;
assert_eq!(
snap.samples.len() as u64,
seq,
"the host ring keeps growing while the service half panics"
);
assert!(
snap.service_samples.is_empty(),
"a panicking service half records nothing for that tick — the UI reads the gap \
as no measurement, never as 0.0"
);
}
guarded("services", async {
history
.record_service_samples(ServiceSampleBatch {
sampled_at_unix: 4,
services: vec![ServiceSample {
id: "trusty-search".to_string(),
status: ServiceStatus::Running,
cpu_pct: Some(2.0),
rss_bytes: Some(64 * 1024 * 1024),
}],
})
.await;
})
.await;
let snap = history.snapshot().await;
assert_eq!(snap.samples.len(), 3);
assert_eq!(
snap.service_samples["trusty-search"].len(),
1,
"the service half recovers once its next tick succeeds"
);
}
}