use std::time::{Duration, Instant};
use phoxal_api::v2 as api;
use phoxal_bus::{Bus, LogicalTime, OwnerCap, Publisher};
use sysinfo::{Pid, ProcessRefreshKind, ProcessesToUpdate, RefreshKind, System};
use tokio::sync::watch;
use tokio::task::JoinHandle;
pub(crate) type ProcessMetricsSample = api::telemetry::Process;
pub(crate) const PROCESS_METRICS_INTERVAL: Duration = Duration::from_secs(3);
pub(crate) struct ProcessMetricsPublisher {
publisher: Option<Publisher<api::telemetry::Process>>,
}
impl ProcessMetricsPublisher {
pub(crate) fn attach(bus: Bus) -> Self {
let topic = api::topic::internal::new(OwnerCap::__mint())
.telemetry()
.process();
let publisher = Publisher::new(bus, &topic)
.map_err(|error| {
tracing::warn!(
target: "phoxal.runtime",
error = %error,
"process telemetry publisher could not be created"
);
error
})
.ok();
Self { publisher }
}
#[cfg(test)]
fn disabled() -> Self {
Self { publisher: None }
}
pub(crate) fn publish(&self, at: LogicalTime, body: ProcessMetricsSample) {
let Some(publisher) = &self.publisher else {
return;
};
if let Err(error) = publisher.try_publish(at, body) {
tracing::warn!(
target: "phoxal.runtime",
error = %error,
"process telemetry publish failed"
);
}
}
}
pub(crate) fn spawn_sampler() -> (
watch::Receiver<Option<ProcessMetricsSample>>,
JoinHandle<()>,
) {
let (tx, rx) = watch::channel(None);
let handle = tokio::spawn(run_sampler(tx));
(rx, handle)
}
async fn run_sampler(tx: watch::Sender<Option<api::telemetry::Process>>) {
let pid = Pid::from_u32(std::process::id());
let (mut system, mut previous_refresh) = match tokio::task::spawn_blocking(|| {
let system = System::new_with_specifics(
RefreshKind::nothing()
.with_processes(ProcessRefreshKind::nothing().with_cpu().with_memory()),
);
(system, Instant::now())
})
.await
{
Ok(system) => system,
Err(error) => {
tracing::warn!(
target: "phoxal.runtime",
error = %error,
"process telemetry sampler initialization failed; stopping self-sampling"
);
return;
}
};
let mut interval = tokio::time::interval(PROCESS_METRICS_INTERVAL);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
interval.tick().await;
loop {
interval.tick().await;
let (returned, sample, refreshed_at) = match tokio::task::spawn_blocking(move || {
let (sample, refreshed_at) = sample_process(&mut system, pid, previous_refresh);
(system, sample, refreshed_at)
})
.await
{
Ok(pair) => pair,
Err(error) => {
tracing::warn!(
target: "phoxal.runtime",
error = %error,
"process telemetry sampler task failed; stopping self-sampling"
);
return;
}
};
system = returned;
previous_refresh = refreshed_at;
if let Some(body) = sample {
if tx.send(Some(body)).is_err() {
return;
}
}
}
}
fn sample_process(
system: &mut System,
pid: Pid,
previous_refresh: Instant,
) -> (Option<api::telemetry::Process>, Instant) {
system.refresh_processes_specifics(
ProcessesToUpdate::Some(&[pid]),
false,
ProcessRefreshKind::nothing().with_cpu().with_memory(),
);
let refreshed_at = Instant::now();
let sample = system.process(pid).map(|process| api::telemetry::Process {
cpu_pct: process.cpu_usage(),
rss_bytes: process.memory(),
window_ns: u64::try_from(
refreshed_at
.saturating_duration_since(previous_refresh)
.as_nanos(),
)
.unwrap_or(u64::MAX),
});
(sample, refreshed_at)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn disabled_process_metrics_publisher_is_a_noop() {
let metrics = ProcessMetricsPublisher::disabled();
metrics.publish(
LogicalTime::new(0, 1),
api::telemetry::Process {
cpu_pct: 0.0,
rss_bytes: 0,
window_ns: 0,
},
);
assert!(metrics.publisher.is_none());
}
#[serial_test::serial]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn process_metrics_publish_on_the_declared_topic() {
let bus = Bus::open(phoxal_bus::BusConfig::in_process("dev", "metrics-test"))
.await
.expect("open bus");
let latest = phoxal_bus::Latest::<api::telemetry::Process>::new(
&bus,
&api::topic::new().telemetry().process(),
)
.await
.expect("subscribe process telemetry");
let metrics = ProcessMetricsPublisher::attach(bus.clone());
metrics.publish(
LogicalTime::new(0, 1),
api::telemetry::Process {
cpu_pct: 12.5,
rss_bytes: 42,
window_ns: 3_000_000_000,
},
);
let mut observed = None;
for _ in 0..50 {
if let Some(sample) = latest.latest() {
observed = Some(sample);
break;
}
tokio::time::sleep(Duration::from_millis(20)).await;
}
let observed = observed.expect("process telemetry sample");
assert_eq!(observed.rss_bytes, 42);
assert_eq!(observed.window_ns, 3_000_000_000);
bus.close().await.expect("close bus");
}
#[test]
fn sample_process_reports_this_process_over_the_measured_window() {
let mut system = System::new_with_specifics(
RefreshKind::nothing()
.with_processes(ProcessRefreshKind::nothing().with_cpu().with_memory()),
);
let pid = Pid::from_u32(std::process::id());
let previous_refresh = Instant::now()
.checked_sub(Duration::from_secs(2))
.expect("two seconds before now");
let (sample, _) = sample_process(&mut system, pid, previous_refresh);
let sample =
sample.expect("the current process is always present in its own sysinfo refresh");
assert!(sample.window_ns >= 2_000_000_000);
assert!(sample.window_ns < 3_000_000_000);
assert!(sample.rss_bytes > 0);
}
}