use std::time::Duration;
use sysinfo::Disks;
use crate::api::FrameBus;
use crate::api::handlers::SharedState;
use crate::app_state::AppState;
use crate::device::{
ChassisInfo, CpuInfo, GpuInfo, MemoryInfo, ProcessInfo, create_chassis_reader, get_cpu_readers,
get_gpu_readers, get_memory_readers,
};
use crate::snapshot::{SNAPSHOT_SCHEMA_VERSION, Snapshot};
use crate::storage::info::StorageInfo;
use crate::utils::{filter_docker_aware_disks, get_hostname};
pub async fn run_collection_loop(
state: SharedState,
bus: FrameBus,
interval_secs: u64,
processes_enabled: bool,
) {
let gpu_readers = get_gpu_readers();
let cpu_readers = get_cpu_readers();
let memory_readers = get_memory_readers();
let chassis_reader = create_chassis_reader();
let mut disks = Disks::new_with_refreshed_list();
let hostname = get_hostname();
let interval = Duration::from_secs(interval_secs);
loop {
let all_gpu_info: Vec<GpuInfo> = gpu_readers
.iter()
.flat_map(|reader| reader.get_gpu_info())
.collect();
let all_cpu_info: Vec<CpuInfo> = cpu_readers
.iter()
.flat_map(|reader| reader.get_cpu_info())
.collect();
let all_memory_info: Vec<MemoryInfo> = memory_readers
.iter()
.flat_map(|reader| reader.get_memory_info())
.collect();
let all_processes: Vec<ProcessInfo> = if processes_enabled {
gpu_readers
.iter()
.flat_map(|reader| reader.get_process_info())
.collect()
} else {
Vec::new()
};
let chassis_info: Vec<ChassisInfo> = chassis_reader
.get_chassis_info()
.into_iter()
.map(|mut ci| {
if ci.total_power_watts.is_none() {
let total_gpu_power =
crate::metrics::gpu_readings::total_power_watts(&all_gpu_info);
if total_gpu_power > 0.0 {
ci.total_power_watts = Some(total_gpu_power);
}
}
ci
})
.collect();
disks.refresh(true);
let storage_info = collect_storage_info_from(&disks, &hostname);
let frame = Snapshot {
schema: SNAPSHOT_SCHEMA_VERSION,
timestamp: chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
hostname: hostname.clone(),
gpus: Some(all_gpu_info.clone()),
cpus: Some(all_cpu_info.clone()),
memory: Some(all_memory_info.clone()),
chassis: Some(chassis_info.clone()),
processes: if processes_enabled {
Some(all_processes.clone())
} else {
None
},
storage: Some(storage_info.clone()),
errors: Vec::new(),
};
{
let mut guard = state.write().await;
guard.gpu_info = all_gpu_info;
guard.cpu_info = all_cpu_info;
guard.memory_info = all_memory_info;
guard.process_info = all_processes;
guard.chassis_info = chassis_info;
guard.storage_info = storage_info;
if guard.loading {
guard.loading = false;
}
integrate_power_samples(&mut guard);
}
bus.publish(frame).await;
tokio::time::sleep(interval).await;
}
}
pub(crate) fn integrate_power_samples(state: &mut AppState) {
use std::time::Instant;
use crate::metrics::energy::EnergyKey;
let now = Instant::now();
let mut samples: Vec<(EnergyKey, f64)> =
Vec::with_capacity(state.gpu_info.len() + state.cpu_info.len() + state.chassis_info.len());
for gpu in &state.gpu_info {
if let Some(watts) = gpu.power_consumption_reading() {
samples.push((
EnergyKey::gpu(gpu.hostname.clone(), gpu.uuid.clone()),
watts,
));
}
}
for cpu in &state.cpu_info {
if let Some(power) = cpu.power_consumption {
samples.push((EnergyKey::cpu(cpu.hostname.clone()), power));
}
}
for chassis in &state.chassis_info {
if let Some(power) = chassis.total_power_watts {
samples.push((EnergyKey::chassis(chassis.hostname.clone()), power));
}
}
let wal_index = &mut state.energy_wal_replay;
let integrator = state.energy.integrator_mut();
for (key, watts) in samples {
if !integrator.has_samples(&key) && !wal_index.is_empty() {
wal_index.seed_if_matches(&key, integrator);
}
integrator.record_sample(key, now, watts);
}
}
fn collect_storage_info_from(disks: &Disks, hostname: &str) -> Vec<StorageInfo> {
let mut storage_info = Vec::new();
let mut filtered_disks = filter_docker_aware_disks(disks);
filtered_disks.sort_by(|a, b| {
a.mount_point()
.to_string_lossy()
.cmp(&b.mount_point().to_string_lossy())
});
for (index, disk) in filtered_disks.iter().enumerate() {
let mount_point_str = disk.mount_point().to_string_lossy();
storage_info.push(StorageInfo {
mount_point: mount_point_str.to_string(),
total_bytes: disk.total_space(),
available_bytes: disk.available_space(),
host_id: hostname.to_string(),
hostname: hostname.to_string(),
index: index as u32,
});
}
storage_info
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn collect_storage_info_respects_hostname() {
let disks = Disks::new_with_refreshed_list();
let _ = collect_storage_info_from(&disks, "h");
}
}