diff options
| author | Dennis Kobert <dennis@kobert.dev> | 2025-03-29 15:52:47 +0100 |
|---|---|---|
| committer | Dennis Kobert <dennis@kobert.dev> | 2025-03-29 15:52:47 +0100 |
| commit | 336eeb9b3e5f9cb9e8584719c32f02689300aee4 (patch) | |
| tree | f944f47085ea71704d1dbf903967499002bfa138 /src | |
| parent | fd89fa452aa2c7a0124f9acb09fe0bbd8105b2d9 (diff) | |
Integrate energy model into perf estimator
Diffstat (limited to 'src')
| -rw-r--r-- | src/benchmark.rs | 6 | ||||
| -rw-r--r-- | src/energy.rs | 161 | ||||
| -rw-r--r-- | src/energy/trackers.rs | 2 | ||||
| -rw-r--r-- | src/energy/trackers/kernel.rs | 6 | ||||
| -rw-r--r-- | src/energy/trackers/mock.rs | 4 | ||||
| -rw-r--r-- | src/energy/trackers/perf.rs | 107 | ||||
| -rw-r--r-- | src/main.rs | 13 | ||||
| -rw-r--r-- | src/scheduler.rs | 18 |
8 files changed, 203 insertions, 114 deletions
diff --git a/src/benchmark.rs b/src/benchmark.rs index 6701f0c..2640ba4 100644 --- a/src/benchmark.rs +++ b/src/benchmark.rs @@ -354,7 +354,7 @@ fn initialize_csv_writer( Ok(csv_writer) } -fn read_cpu_frequency(cpu_id: u32) -> Option<f64> { +pub fn read_cpu_frequency(cpu_id: u32) -> Option<f64> { // Try to read frequency from sysfs let freq_path = format!( "/sys/devices/system/cpu/cpu{}/cpufreq/scaling_cur_freq", @@ -403,6 +403,10 @@ fn define_available_events() -> Vec<(String, Event)> { "ref_cpu_cycles".to_string(), Event::Hardware(Hardware::REF_CPU_CYCLES), ), + ( + "task_clock".to_string(), + Event::Software(Software::TASK_CLOCK), + ), // ( // "stalled-cycles-frontend".to_string(), // Event::Hardware(Hardware::STALLED_CYCLES_FRONTEND), diff --git a/src/energy.rs b/src/energy.rs index c989685..6490efd 100644 --- a/src/energy.rs +++ b/src/energy.rs @@ -22,15 +22,17 @@ pub enum Request { LogEnergyTask(Pid), } +#[derive(Clone)] pub struct ProcessInfo { - energy: u64, + energy: f64, + tree_energy: f64, last_update: std::time::Instant, parent: Pid, } pub struct EnergyLog { - pub energy_of_exited_tasks: u64, - pub energy_of_running_tasks: u64, + pub energy_of_exited_tasks: f64, + pub energy_of_running_tasks: f64, } pub struct EnergyService { @@ -46,6 +48,8 @@ pub struct EnergyService { shared_cpu_current_frequencies: Arc<RwLock<Vec<FrequencyKHZ>>>, energy_logging_children_to_root_of_tree: HashMap<Pid, Pid>, logged_energy: Arc<RwLock<HashMap<Pid, EnergyLog>>>, + rapl_offset: f64, + old_rapl: f64, } impl EnergyService { @@ -73,6 +77,8 @@ impl EnergyService { shared_cpu_current_frequencies, energy_logging_children_to_root_of_tree: HashMap::new(), logged_energy, + rapl_offset: rapl::read_package_energy().unwrap(), + old_rapl: 0., } } @@ -96,55 +102,65 @@ impl EnergyService { fn handle_requests(&mut self) { while let Ok(request) = self.request_receiver.try_recv() { - match request { - Request::NewTask(pid) => { - self.estimator.start_trace(pid as u64); - self.active_processes.insert(pid); - let parent = (|| { - let process = procfs::process::Process::new(pid)?; - process.stat().map(|stat| stat.ppid) - })() - .unwrap_or_default(); - self.process_info.insert( - pid, - ProcessInfo { - energy: 0, - last_update: std::time::Instant::now(), - parent, - }, - ); - // Initialize with default budget - self.shared_budgets.insert(pid, u64::MAX); + self.handle_request(request); + } + } + + fn handle_request(&mut self, request: Request) { + match request { + Request::NewTask(pid) => { + self.estimator.start_trace(pid as u64); + self.active_processes.insert(pid); + let parent = (|| { + let process = procfs::process::Process::new(pid)?; + process.stat().map(|stat| stat.ppid) + })() + .unwrap_or_default(); + self.process_info.insert( + pid, + ProcessInfo { + energy: 0., + tree_energy: 0., + last_update: std::time::Instant::now(), + parent, + }, + ); + // Initialize with default budget + self.shared_budgets.insert(pid, u64::MAX); + if !self.process_info.contains_key(&parent) && parent != 0 { + self.handle_request(Request::NewTask(parent)); } - Request::RemoveTask(pid) => { - self.estimator.stop_trace(pid as u64); - if let Some(root_pid) = - self.energy_logging_children_to_root_of_tree.remove(&pid) - { - if let Some(info) = self.process_info.remove(&pid) { - self.logged_energy - .write() - .unwrap() - .get_mut(&root_pid) - .unwrap() - .energy_of_exited_tasks += info.energy; - } - } - self.active_processes.remove(&pid); - self.process_info.remove(&pid); - self.shared_budgets.remove(&pid); + } + Request::RemoveTask(pid) => { + if procfs::process::Process::new(pid).is_ok() { + return; } - Request::LogEnergyTask(pid) => { - self.energy_logging_children_to_root_of_tree - .insert(pid, pid); - self.logged_energy.write().unwrap().insert( - pid, - EnergyLog { - energy_of_exited_tasks: 0, - energy_of_running_tasks: 0, - }, - ); + + self.estimator.stop_trace(pid as u64); + if let Some(root_pid) = self.energy_logging_children_to_root_of_tree.remove(&pid) { + if let Some(info) = self.process_info.remove(&pid) { + self.logged_energy + .write() + .unwrap() + .get_mut(&root_pid) + .unwrap() + .energy_of_exited_tasks += info.energy; + } } + self.active_processes.remove(&pid); + self.process_info.remove(&pid); + self.shared_budgets.remove(&pid); + } + Request::LogEnergyTask(pid) => { + self.energy_logging_children_to_root_of_tree + .insert(pid, pid); + self.logged_energy.write().unwrap().insert( + pid, + EnergyLog { + energy_of_exited_tasks: 0., + energy_of_running_tasks: 0., + }, + ); } } } @@ -152,31 +168,36 @@ impl EnergyService { fn update_measurements(&mut self) { let mut logged_energy = self.logged_energy.write().unwrap(); for (_, value) in logged_energy.iter_mut() { - value.energy_of_running_tasks = 0u64; + value.energy_of_running_tasks = 0f64; } + let old_energy = self.process_info.get(&1).unwrap().tree_energy; for &pid in &self.active_processes { if let Some(info) = self.process_info.get_mut(&pid) { - let energy = self.estimator.read_consumption(pid as u64); - if let Some(root_pid) = self.energy_logging_children_to_root_of_tree.get(&pid) { - logged_energy - .get_mut(root_pid) - .unwrap() - .energy_of_running_tasks += energy; - } else if let Some(&root_pid) = self - .energy_logging_children_to_root_of_tree - .get(&info.parent) - { - self.energy_logging_children_to_root_of_tree - .insert(pid, root_pid); - logged_energy - .get_mut(&root_pid) - .unwrap() - .energy_of_running_tasks += energy; + if let Some(energy) = self.estimator.read_consumption(pid as u64) { + info.energy += energy; + let mut parent = info.parent; + while let Some(info) = self.process_info.get_mut(&parent) { + info.tree_energy += energy; + info.last_update = std::time::Instant::now(); + parent = info.parent; + } } - info.energy = energy; - info.last_update = std::time::Instant::now(); } } + if let Some(init) = self.process_info.get(&1) { + let rapl = rapl::read_package_energy().unwrap() - self.rapl_offset; + let rapl_diff = rapl - self.old_rapl; + self.old_rapl = rapl; + //let v = logged_energy.get(&1).unwrap(); + println!( + "Energy estimation: {:.1} rapl: {:.1}, est diff: {:.1} rapl diff: {:.1}", + init.tree_energy, + rapl, + (init.tree_energy - old_energy), + rapl_diff, + //(rapl_diff - (init.tree_energy - old_energy)) / rapl_diff * 100., + ); + } } fn update_budgets(&mut self) { @@ -202,11 +223,11 @@ impl EnergyService { rapl::read_package_energy().unwrap() } - pub fn process_energy(&self, pid: Pid) -> Option<u64> { + pub fn process_energy(&self, pid: Pid) -> Option<f64> { self.process_info.get(&pid).map(|info| info.energy) } - pub fn all_process_energies(&self) -> HashMap<Pid, u64> { + pub fn all_process_energies(&self) -> HashMap<Pid, f64> { self.process_info .iter() .map(|(&pid, info)| (pid, info.energy)) diff --git a/src/energy/trackers.rs b/src/energy/trackers.rs index 31e2f3e..1665177 100644 --- a/src/energy/trackers.rs +++ b/src/energy/trackers.rs @@ -10,5 +10,5 @@ pub use perf::*; pub trait Estimator: Send + 'static { fn start_trace(&mut self, pid: u64); fn stop_trace(&mut self, pid: u64); - fn read_consumption(&mut self, pid: u64) -> u64; + fn read_consumption(&mut self, pid: u64) -> Option<f64>; } diff --git a/src/energy/trackers/kernel.rs b/src/energy/trackers/kernel.rs index b71fee6..2a19e4c 100644 --- a/src/energy/trackers/kernel.rs +++ b/src/energy/trackers/kernel.rs @@ -38,12 +38,12 @@ impl Estimator for KernelDriver { let _ = STOP_TRACE.ioctl(&mut self.file, &pid); } - fn read_consumption(&mut self, pid: u64) -> u64 { + fn read_consumption(&mut self, pid: u64) -> Option<f64> { let mut arg = pid; if READ_POWER.ioctl(&mut self.file, &mut arg).is_ok() { - arg + Some(arg as f64) } else { - 0 + None } } } diff --git a/src/energy/trackers/mock.rs b/src/energy/trackers/mock.rs index eb23421..e3ef377 100644 --- a/src/energy/trackers/mock.rs +++ b/src/energy/trackers/mock.rs @@ -8,7 +8,7 @@ impl Estimator for MockEstimator { fn stop_trace(&mut self, _pid: u64) {} - fn read_consumption(&mut self, _pid: u64) -> u64 { - 14 + fn read_consumption(&mut self, _pid: u64) -> Option<f64> { + Some(14.) } } diff --git a/src/energy/trackers/perf.rs b/src/energy/trackers/perf.rs index 17bc693..d340e1f 100644 --- a/src/energy/trackers/perf.rs +++ b/src/energy/trackers/perf.rs @@ -1,43 +1,76 @@ use std::collections::HashMap; +use burn::tensor::Tensor; use perf_event::{ - events::{Event, Hardware}, + events::{Event, Hardware, Software}, Builder, Counter, Group, }; -use crate::energy::Estimator; +use crate::energy::{rapl, Estimator}; +type ArrayBackend = burn_ndarray::NdArray<f32>; -#[derive(Default)] pub struct PerfEstimator { registry: HashMap<u64, Counters>, + model: crate::model::Net<ArrayBackend>, + device: <ArrayBackend as burn::prelude::Backend>::Device, +} + +impl Default for PerfEstimator { + fn default() -> Self { + let model = crate::model::load_model(); + Self { + registry: Default::default(), + model, + device: Default::default(), + } + } } struct Counters { group: Group, counters: Vec<Counter>, + old_time: u64, + old_total_energy: f64, } -static EVENT_TYPES: &[(f32, Event)] = &[ - (1.0, Event::Hardware(Hardware::CPU_CYCLES)), - (2.0, Event::Hardware(Hardware::INSTRUCTIONS)), +static EVENT_TYPES: &[Event] = &[ + Event::Hardware(Hardware::BRANCH_INSTRUCTIONS), + Event::Hardware(Hardware::BRANCH_MISSES), + Event::Hardware(Hardware::CACHE_MISSES), + Event::Hardware(Hardware::CACHE_REFERENCES), + Event::Hardware(Hardware::CPU_CYCLES), + Event::Hardware(Hardware::INSTRUCTIONS), + Event::Hardware(Hardware::REF_CPU_CYCLES), + Event::Software(Software::TASK_CLOCK), ]; impl Estimator for PerfEstimator { fn start_trace(&mut self, pid: u64) { - let Ok(mut group) = Group::new_with_pid_and_cpu(-1, 0) else { - eprintln!("Failed to create performance counter group for PID {}", pid); - return; + //let Ok(mut group) = Group::new_with_pid_and_cpu(-1, 0) else { + let mut group = match Group::new_with_pid_and_cpu(pid as i32, -1) { + Ok(counters) => counters, + Err(e) => { + eprintln!( + "Failed to create performance counter group for PID {}: {}", + pid, e + ); + return; + } }; + //let Ok(mut group) = Group::new_with_pid_and_cpu(pid as i32, -1) else { + // eprintln!("Failed to create performance counter group for PID {}", pid); + // return; + //}; let counters: Result<Vec<_>, _> = EVENT_TYPES .iter() - .map(|(_, kind)| { + .map(|kind| { Builder::new() .group(&mut group) .kind(kind.clone()) - // .observe_pid(pid as i32) - .observe_pid(-1) - .one_cpu(0) + .observe_pid(pid as i32) + //.observe_pid(-1) + //.one_cpu(0) .build() }) .collect(); @@ -57,8 +90,14 @@ impl Estimator for PerfEstimator { eprintln!("Failed to enable performance counters: {}", e); return; } + group.reset(); - let counters = Counters { counters, group }; + let counters = Counters { + counters, + group, + old_time: 0, + old_total_energy: 0., + }; self.registry.insert(pid, counters); } @@ -66,21 +105,47 @@ impl Estimator for PerfEstimator { self.registry.remove(&pid); } - fn read_consumption(&mut self, pid: u64) -> u64 { + fn read_consumption(&mut self, pid: u64) -> Option<f64> { let Some(counters) = self.registry.get_mut(&pid) else { - return 0; + return None; }; let counts = match counters.group.read() { Ok(counts) => counts, - Err(_) => return 0, + Err(e) => { + println!("failed to read group: {e}"); + return None; + } }; + //let task_clock = counts[&counters.counters[7]]; + let task_clock = counts[&counters.counters[7]]; + + if task_clock == 0 { + return None; + } + let current_time = std::time::Instant::now(); + let time_running_ns = counts.time_running(); + let correction_factor = 10_000_000. / (time_running_ns - counters.old_time) as f64; + counters.old_time = time_running_ns; + //dbg!(time_running_ns); + //dbg!(counters.old_time.elapsed().as_millis()); - let mut sum = 0; - for ((factor, _ty), count) in EVENT_TYPES.iter().zip(counts.iter()) { - sum += *factor as u64 * count.1; + // TODO:: Respect per core frequency for the running task + let mut values = vec![crate::benchmark::read_cpu_frequency(0).unwrap()]; + for ty in counters.counters.iter().take(7) { + let count: u64 = counts[&ty]; + values.push((count as f64) * correction_factor); } + //dbg!(&values); + let result = self + .model + .forward(Tensor::from_floats(&values.as_slice()[0..8], &self.device)); - sum + let energy = result.into_scalar() as f64; + //dbg!(energy, correction_factor); + counters.old_total_energy += energy / correction_factor; + counters.group.reset().unwrap(); + //dbg!(counters.old_total_energy); + Some(energy / correction_factor) } } diff --git a/src/main.rs b/src/main.rs index 0caae55..6242917 100644 --- a/src/main.rs +++ b/src/main.rs @@ -56,19 +56,6 @@ fn main() -> Result<()> { ) .get_matches(); - let model = model::load_model(); - let device = Default::default(); - let result = model.forward(Tensor::from_floats( - [ - // 2899.97, 21420886., 59226., 148084., 301003., 36244800., 115862107., 43766905., - 1533, 1077473, 9448, 4269, 52805, 2456984, 5867215, 3954587, - ], - // [1., 1., 1., 1., 1., 1., 1., 1.], - &device, - )); - - println!("result: {}", result); - let power_cap = *matches.get_one::<u64>("power_cap").unwrap_or(&u64::MAX); let use_mocking = matches.get_flag("mock"); let benchmark = matches.get_flag("benchmark"); diff --git a/src/scheduler.rs b/src/scheduler.rs index 42d480e..0a83b12 100644 --- a/src/scheduler.rs +++ b/src/scheduler.rs @@ -216,10 +216,22 @@ impl<'a> Scheduler<'a> { // Low budget tasks go to e-cores let cpu = self.e_core_selector.next_core(task.cpu); + if cpu >= 0 { - dispatched_task.cpu = cpu; + //dispatched_task.cpu = cpu; + // TODO: Fix efficiency core allocation + //dispatched_task.flags |= RL_CPU_ANY as u64; } else { - dispatched_task.flags |= RL_CPU_ANY as u64; + //dispatched_task.flags |= RL_CPU_ANY as u64; + } + // If it's our own process, schedule it to core 1 + if task.pid == self.own_pid { + dispatched_task.cpu = 1; + // dispatched_task.flags |= RL_CPU_ANY as u64; + } else { + // Schedule all other tasks on core 0 + // dispatched_task.flags |= RL_CPU_ANY as u64; + dispatched_task.cpu = 0; } // Scheduler tasks get longer slices @@ -242,7 +254,7 @@ impl<'a> Scheduler<'a> { fn cleanup_old_tasks(&mut self) { let current = Instant::now(); for (pid, last_scheduled) in &self.managed_tasks { - if current - *last_scheduled > Duration::from_secs(5) { + if current - *last_scheduled > Duration::from_secs(500) { self.to_remove.push(*pid); } } |
