summaryrefslogtreecommitdiff
path: root/src
diff options
context:
space:
mode:
authorDennis Kobert <dennis@kobert.dev>2025-03-29 15:52:47 +0100
committerDennis Kobert <dennis@kobert.dev>2025-03-29 15:52:47 +0100
commit336eeb9b3e5f9cb9e8584719c32f02689300aee4 (patch)
treef944f47085ea71704d1dbf903967499002bfa138 /src
parentfd89fa452aa2c7a0124f9acb09fe0bbd8105b2d9 (diff)
Integrate energy model into perf estimator
Diffstat (limited to 'src')
-rw-r--r--src/benchmark.rs6
-rw-r--r--src/energy.rs161
-rw-r--r--src/energy/trackers.rs2
-rw-r--r--src/energy/trackers/kernel.rs6
-rw-r--r--src/energy/trackers/mock.rs4
-rw-r--r--src/energy/trackers/perf.rs107
-rw-r--r--src/main.rs13
-rw-r--r--src/scheduler.rs18
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);
}
}