use std::collections::HashMap;
use std::path::PathBuf;
use std::time::{Duration, Instant};
use trusty_common::sys_metrics::ProcessCpuSampler;
use crate::connector::{ServiceInfo, ServiceStatus};
use crate::machine_history::service_samples::{ServiceSample, ServiceSampleBatch};
const LOOKUP_BACKOFF: Duration = Duration::from_secs(15);
#[must_use]
pub fn resolve_pid(service_id: &str) -> Option<u32> {
match service_id {
"trusty-mpm" => mpm_lock_pid(&mpm_lock_path()?),
"trusty-search" => socket_peer_pid(crate::search_uds::socket_path().ok()?),
"trusty-memory" => {
socket_peer_pid(trusty_common::daemon_socket_path("trusty-memory").ok()?)
}
"trusty-console" => Some(std::process::id()),
_ => None,
}
}
fn mpm_lock_path() -> Option<PathBuf> {
Some(dirs::home_dir()?.join(".trusty-mpm").join("daemon.lock"))
}
fn mpm_lock_pid(path: &std::path::Path) -> Option<u32> {
let body = std::fs::read_to_string(path).ok()?;
for line in body.lines() {
let Some((key, value)) = line.split_once('=') else {
continue;
};
if key.trim() != "pid" {
continue;
}
if let Ok(pid) = value.trim().trim_matches('"').parse::<u32>()
&& pid > 0
{
return Some(pid);
}
}
None
}
fn socket_peer_pid(socket: PathBuf) -> Option<u32> {
let std_stream = std::os::unix::net::UnixStream::connect(&socket).ok()?;
std_stream.set_nonblocking(true).ok()?;
let stream = tokio::net::UnixStream::from_std(std_stream).ok()?;
trusty_common::uds::peer_pid(&stream)
}
pub struct ServiceMetricsSampler {
cpu: ProcessCpuSampler,
pids: HashMap<String, u32>,
retry_after: HashMap<String, Instant>,
}
impl Default for ServiceMetricsSampler {
fn default() -> Self {
Self::new()
}
}
impl ServiceMetricsSampler {
#[must_use]
pub fn new() -> Self {
Self {
cpu: ProcessCpuSampler::new(),
pids: HashMap::new(),
retry_after: HashMap::new(),
}
}
#[must_use]
pub fn pending_lookups(&self, services: &[ServiceInfo], now: Instant) -> Vec<String> {
services
.iter()
.filter(|s| is_live(&s.status))
.filter(|s| !self.pids.contains_key(&s.id))
.filter(|s| self.retry_after.get(&s.id).is_none_or(|at| now >= *at))
.map(|s| s.id.clone())
.collect()
}
pub fn record_lookups(&mut self, found: Vec<(String, Option<u32>)>, now: Instant) {
for (id, pid) in found {
match pid {
Some(pid) => {
self.cpu.track(pid);
self.pids.insert(id.clone(), pid);
self.retry_after.remove(&id);
}
None => {
self.retry_after.insert(id, now + LOOKUP_BACKOFF);
}
}
}
}
#[must_use]
pub fn sample(&mut self, services: &[ServiceInfo], sampled_at_unix: u64) -> ServiceSampleBatch {
self.cpu.refresh();
let mut samples = Vec::with_capacity(services.len());
for service in services {
let (cpu_pct, rss_bytes) = match self.pids.get(&service.id) {
Some(pid) => (self.cpu.cpu_pct(*pid), self.cpu.rss_bytes(*pid)),
None => (None, None),
};
if cpu_pct.is_none()
&& rss_bytes.is_none()
&& let Some(pid) = self.pids.remove(&service.id)
{
self.cpu.untrack(pid);
}
samples.push(ServiceSample {
id: service.id.clone(),
status: service.status.clone(),
cpu_pct,
rss_bytes,
});
}
ServiceSampleBatch {
sampled_at_unix,
services: samples,
}
}
#[must_use]
pub fn pid_of(&self, service_id: &str) -> Option<u32> {
self.pids.get(service_id).copied()
}
}
pub async fn apply_metrics_overlay(
services: &mut [ServiceInfo],
history: &crate::machine_history::MachineHistory,
) {
let latest = history.latest_service_metrics().await;
for service in services.iter_mut() {
let newest = latest.get(&service.id);
service.cpu_pct = newest.and_then(|s| s.cpu_pct);
service.rss_bytes = newest.and_then(|s| s.rss_bytes);
}
}
fn is_live(status: &ServiceStatus) -> bool {
matches!(status, ServiceStatus::Running | ServiceStatus::Degraded)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::connector::ServiceLifecycle;
fn info(id: &str, status: ServiceStatus) -> ServiceInfo {
ServiceInfo {
id: id.to_string(),
display_name: id.to_string(),
status,
version: None,
url: None,
hint: None,
lifecycle: ServiceLifecycle::Daemon,
cpu_pct: None,
rss_bytes: None,
}
}
fn spawn_sleeper() -> std::process::Child {
std::process::Command::new("sleep")
.arg("30")
.spawn()
.expect("spawn a sleeping child")
}
#[test]
fn every_service_gets_a_sample() {
let services = vec![
info("trusty-search", ServiceStatus::Running),
info("trusty-review", ServiceStatus::Available),
info("trusty-analyze", ServiceStatus::Absent),
info("trusty-agents", ServiceStatus::Degraded),
];
let batch = ServiceMetricsSampler::new().sample(&services, 42);
assert_eq!(batch.sampled_at_unix, 42);
assert_eq!(batch.services.len(), 4);
for sample in &batch.services {
assert_eq!(
sample.cpu_pct, None,
"{} has no pid, so it must report no measurement rather than 0.0",
sample.id
);
}
}
#[tokio::test]
async fn console_row_gets_cpu_and_rss_on_the_first_tick() {
let services = vec![info("trusty-console", ServiceStatus::Running)];
let mut sampler = ServiceMetricsSampler::new();
let pending = sampler.pending_lookups(&services, Instant::now());
assert_eq!(
pending,
vec!["trusty-console".to_string()],
"the console row must be offered for a pid lookup on the first tick"
);
let found: Vec<(String, Option<u32>)> = pending
.into_iter()
.map(|id| {
let pid = resolve_pid(&id);
(id, pid)
})
.collect();
assert_eq!(
found[0].1,
Some(std::process::id()),
"resolve_pid must answer the console with this very process"
);
sampler.record_lookups(found, Instant::now());
let history = crate::machine_history::MachineHistory::new();
history
.record_service_samples(sampler.sample(&services, 1))
.await;
let snap = history.snapshot().await;
let ring = snap
.service_samples
.get("trusty-console")
.expect("the console must have a per-service ring after one tick");
assert!(
!ring.is_empty(),
"the ring must hold the first tick's sample"
);
let newest = ring.last().expect("a non-empty ring has a last entry");
assert!(
newest.cpu_pct.is_some(),
"the console row must carry a CPU measurement, not null: {newest:?}"
);
assert!(
newest.rss_bytes.is_some_and(|b| b > 0),
"the console row must carry a non-zero RSS measurement: {newest:?}"
);
}
#[test]
fn console_pid_is_this_process() {
assert_eq!(resolve_pid("trusty-console"), Some(std::process::id()));
assert_eq!(resolve_pid("trusty-consoleX"), None);
}
#[test]
fn pending_lookups_skips_a_service_that_is_not_running() {
let services = vec![
info("trusty-search", ServiceStatus::Running),
info("trusty-memory", ServiceStatus::Available),
info("trusty-review", ServiceStatus::Absent),
info("trusty-mpm", ServiceStatus::Degraded),
];
let mut pending = ServiceMetricsSampler::new().pending_lookups(&services, Instant::now());
pending.sort();
assert_eq!(pending, vec!["trusty-mpm", "trusty-search"]);
}
#[test]
fn a_failed_lookup_is_not_retried_immediately() {
let services = vec![info("trusty-agents", ServiceStatus::Running)];
let now = Instant::now();
let mut sampler = ServiceMetricsSampler::new();
assert_eq!(sampler.pending_lookups(&services, now).len(), 1);
sampler.record_lookups(vec![("trusty-agents".to_string(), None)], now);
assert!(
sampler.pending_lookups(&services, now).is_empty(),
"a failed lookup must back off, not retry on the next tick"
);
assert_eq!(
sampler
.pending_lookups(&services, now + LOOKUP_BACKOFF)
.len(),
1,
"the backoff must expire so a daemon that starts later is found"
);
}
#[test]
fn a_resolved_pid_is_tracked() {
let mut child = spawn_sleeper();
let pid = child.id();
let services = vec![info("trusty-search", ServiceStatus::Running)];
let mut sampler = ServiceMetricsSampler::new();
sampler.record_lookups(
vec![("trusty-search".to_string(), Some(pid))],
Instant::now(),
);
let batch = sampler.sample(&services, 1);
let recorded = sampler.pid_of("trusty-search");
let cpu = batch.services[0].cpu_pct;
let _ = child.kill();
let _ = child.wait();
assert_eq!(recorded, Some(pid));
assert!(
cpu.is_some(),
"a live tracked process must produce a measurement"
);
}
#[test]
fn a_live_pid_samples_cpu_and_memory_on_one_tick() {
let mut child = spawn_sleeper();
let pid = child.id();
let services = vec![info("trusty-search", ServiceStatus::Running)];
let mut sampler = ServiceMetricsSampler::new();
sampler.record_lookups(
vec![("trusty-search".to_string(), Some(pid))],
Instant::now(),
);
let batch = sampler.sample(&services, 1);
let sample = batch.services[0].clone();
let _ = child.kill();
let _ = child.wait();
assert!(sample.cpu_pct.is_some(), "the CPU half of the tick");
let rss = sample
.rss_bytes
.expect("the memory half of the SAME tick — one refresh serves both");
assert!(rss > 0, "a live process occupies memory, got {rss} bytes");
assert!(
rss < 1024 * 1024 * 1024 * 1024,
"implausibly large ({rss}) — the unit must be bytes"
);
}
#[test]
fn a_vanished_pid_samples_as_none_and_the_tick_continues() {
let mut child = spawn_sleeper();
let pid = child.id();
let services = vec![
info("trusty-search", ServiceStatus::Running),
info("trusty-agents", ServiceStatus::Running),
];
let mut sampler = ServiceMetricsSampler::new();
sampler.record_lookups(
vec![("trusty-search".to_string(), Some(pid))],
Instant::now(),
);
let _ = sampler.sample(&services, 1);
child.kill().expect("kill the sleeper");
child.wait().expect("reap the sleeper");
let batch = sampler.sample(&services, 2);
assert_eq!(
batch.services.len(),
2,
"the tick must not stop or truncate"
);
assert_eq!(
batch.services[0].cpu_pct, None,
"a vanished process reports no measurement, never 0.0"
);
assert_eq!(batch.services[1].cpu_pct, None);
assert_eq!(
sampler.pid_of("trusty-search"),
None,
"the dead pid must be forgotten so a restart can be rediscovered"
);
assert_eq!(sampler.sample(&services, 3).services.len(), 2);
}
#[tokio::test]
async fn the_overlay_stamps_only_the_services_with_a_measurement() {
let history = crate::machine_history::MachineHistory::new();
history
.record_service_samples(ServiceSampleBatch {
sampled_at_unix: 1,
services: vec![
ServiceSample {
id: "trusty-search".to_string(),
status: ServiceStatus::Running,
cpu_pct: Some(7.5),
rss_bytes: Some(148_897_792),
},
ServiceSample {
id: "trusty-review".to_string(),
status: ServiceStatus::Available,
cpu_pct: None,
rss_bytes: None,
},
],
})
.await;
let mut services = vec![
info("trusty-search", ServiceStatus::Running),
info("trusty-review", ServiceStatus::Available),
info("trusty-agents", ServiceStatus::Absent),
];
apply_metrics_overlay(&mut services, &history).await;
assert_eq!(services[0].cpu_pct, Some(7.5));
assert_eq!(services[0].rss_bytes, Some(148_897_792));
assert_eq!(
services[1].cpu_pct, None,
"a sampled-but-unmeasurable service stays null"
);
assert_eq!(
services[1].rss_bytes, None,
"a sampled-but-unmeasurable service stays null in memory too"
);
assert_eq!(
services[2].cpu_pct, None,
"a service the history never saw stays null"
);
assert_eq!(services[2].rss_bytes, None);
}
#[test]
fn resolve_pid_is_none_for_a_service_with_no_source() {
assert_eq!(resolve_pid("trusty-agents"), None);
assert_eq!(resolve_pid("trusty-review"), None);
assert_eq!(resolve_pid("not-a-service"), None);
}
#[test]
fn mpm_pid_is_read_from_the_lock_file() {
let tmp = tempfile::tempdir().expect("tempdir");
let lock = tmp.path().join("daemon.lock");
std::fs::write(
&lock,
"pid = 4242\naddr = \"http://127.0.0.1:7880\"\nstarted_at = \"x\"\n",
)
.expect("write");
assert_eq!(mpm_lock_pid(&lock), Some(4242));
}
#[test]
fn mpm_pid_rejects_a_prefixed_key() {
let tmp = tempfile::tempdir().expect("tempdir");
let lock = tmp.path().join("daemon.lock");
std::fs::write(&lock, "pid_file = \"/tmp/x\"\npidx = 7\n").expect("write");
assert_eq!(mpm_lock_pid(&lock), None);
std::fs::write(&lock, "pid = 0\n").expect("write");
assert_eq!(mpm_lock_pid(&lock), None, "pid 0 is not a process");
}
#[test]
fn mpm_pid_is_none_for_a_missing_file() {
let tmp = tempfile::tempdir().expect("tempdir");
assert_eq!(mpm_lock_pid(&tmp.path().join("absent.lock")), None);
}
#[tokio::test]
async fn socket_peer_pid_reads_this_process_from_its_own_socket() {
let tmp = tempfile::tempdir().expect("tempdir");
let path = tmp.path().join("peer.sock");
let _listener = tokio::net::UnixListener::bind(&path).expect("bind");
#[cfg(any(target_os = "linux", target_os = "macos"))]
assert_eq!(socket_peer_pid(path), Some(std::process::id()));
#[cfg(not(any(target_os = "linux", target_os = "macos")))]
assert_eq!(socket_peer_pid(path), None);
}
#[test]
fn socket_peer_pid_is_none_when_nothing_is_listening() {
let tmp = tempfile::tempdir().expect("tempdir");
assert_eq!(socket_peer_pid(tmp.path().join("nothing.sock")), None);
}
}