use std::io::Write as _;
use std::path::PathBuf;
use std::thread;
use std::time::Duration;
use std::{
fs::{create_dir_all, File},
io::Cursor,
};
use callstack::Measurement;
use chrono::{offset::Local, SecondsFormat};
use derive_builder::Builder;
use inferno::collapse::{dtrace, Collapse};
use inferno::flamegraph;
use log::{error, info};
use super::*;
#[derive(Clone, Debug, PartialEq, Builder)]
#[builder(setter(into))]
pub struct ProfilerRunner {
#[builder(default = "DEFAULT_CHECK_INTERVAL_SECS")]
check_interval_secs: usize,
#[builder(default = "DEFAULT_PCT_CHANGE_TRIGGER")]
report_pct_change_trigger: usize,
#[builder(default)]
reporting_path: String,
#[builder(default = "false")]
expand_frames: bool,
#[builder(default = "false")]
gen_flamegraphs: bool,
#[builder(default = "false")]
measure_allocated_not_retained: bool,
}
const INITIAL_RETAINED_MEM_MB: usize = 20;
pub const DEFAULT_REPORTING_PATH: &str = "ying-profiles";
pub const DEFAULT_CHECK_INTERVAL_SECS: usize = 300;
pub const DEFAULT_PCT_CHANGE_TRIGGER: usize = 10;
impl Default for ProfilerRunner {
fn default() -> Self {
Self::new(
DEFAULT_CHECK_INTERVAL_SECS,
DEFAULT_PCT_CHANGE_TRIGGER,
"",
false,
false,
false,
)
}
}
impl ProfilerRunner {
pub fn new(
check_interval_secs: usize,
report_pct_change_trigger: usize,
reporting_path: &str,
expand_frames: bool,
gen_flamegraphs: bool,
measure_allocated_not_retained: bool,
) -> Self {
Self {
check_interval_secs,
report_pct_change_trigger,
reporting_path: reporting_path.to_string(),
expand_frames,
gen_flamegraphs,
measure_allocated_not_retained,
}
}
pub fn spawn(&self, profiler: &'static YingProfiler) {
if !self.reporting_path.is_empty() {
if let Err(e) = create_dir_all(&self.reporting_path) {
error!(
"Ying: could not create reporting directory {:?}, reports will not be written: {}",
self.reporting_path, e
);
}
}
let check_interval_secs = self.check_interval_secs;
let report_pct_change_trigger = self.report_pct_change_trigger;
let reporting_path = PathBuf::from(self.reporting_path.clone());
let expand_frames = self.expand_frames;
let profiler2 = profiler;
let measurement = if self.measure_allocated_not_retained {
Measurement::AllocatedBytes
} else {
Measurement::RetainedBytes
};
let gen_flamegraphs = self.gen_flamegraphs;
thread::spawn(move || {
let mut last_retained_mem = INITIAL_RETAINED_MEM_MB as f64;
loop {
thread::sleep(Duration::from_secs(check_interval_secs as u64));
let new_allocated = YingProfiler::total_retained_bytes() as f64 / (1024.0 * 1024.0);
let ratio = (new_allocated - last_retained_mem) / last_retained_mem;
info!(
"Ying: total allocated memory is {:.2} MB and ratio to last = {}",
new_allocated, ratio
);
if (ratio.abs() * 100.0) >= report_pct_change_trigger as f64 {
#[cfg(feature = "profile-spans")]
info!(
"Significant memory change registered, dumping profile: new = {}, old = {}",
new_allocated, last_retained_mem
);
last_retained_mem = new_allocated;
let top_stacks = match measurement {
Measurement::AllocatedBytes => profiler2.top_k_stacks_by_allocated(10),
Measurement::RetainedBytes => profiler2.top_k_stacks_by_retained(10),
};
for s in &top_stacks {
println!("---\n{}\n", s.rich_report(profiler2, false, expand_frames));
}
let dt = Local::now();
let dt_str = dt.to_rfc3339_opts(SecondsFormat::Secs, true);
let dump_name = format!("ying.{}.{}MB.report", dt_str, new_allocated as i64);
let mut report_path = reporting_path.clone();
report_path.push(dump_name);
if let Ok(f) = File::create(&report_path) {
for s in &top_stacks {
let _ = writeln!(
&f,
"---\n{}\n",
s.rich_report(profiler2, false, expand_frames)
);
}
} else {
error!("Error: could not write memory report to {:?}", report_path);
}
if gen_flamegraphs {
let graph_name = format!("ying.{}.{}MB.svg", dt_str, new_allocated as i64);
let mut graph_path = reporting_path.clone();
graph_path.push(graph_name);
if let Err(e) = gen_flamegraph(profiler2, measurement, &graph_path) {
error!("Error generating flame graph to {:?}: {}", graph_path, e);
}
}
}
}
});
}
}
pub fn gen_flamegraph(
profiler: &YingProfiler,
measurement: Measurement,
path: &PathBuf,
) -> Result<(), String> {
let mut report = String::new();
match measurement {
Measurement::RetainedBytes => {
let top_stacks = profiler.top_k_stacks_by_retained(50);
for s in &top_stacks {
writeln!(
&mut report,
"{}",
s.dtrace_report(profiler, Measurement::RetainedBytes)
)
.map_err(|e| e.to_string())?;
}
}
Measurement::AllocatedBytes => {
let top_stacks = profiler.top_k_stacks_by_allocated(50);
for s in &top_stacks {
writeln!(
&mut report,
"{}",
s.dtrace_report(profiler, Measurement::AllocatedBytes)
)
.map_err(|e| e.to_string())?;
}
}
}
let mut folder = dtrace::Folder::default();
let mut folded_buf = Vec::new();
let folded_out = Cursor::new(&mut folded_buf);
folder
.collapse(Cursor::new(report), folded_out)
.map_err(|e| e.to_string())?;
if let Ok(f) = File::create(path) {
flamegraph::from_reader(
&mut flamegraph::Options::default(),
Cursor::new(&folded_buf),
f,
)
.map_err(|e| e.to_string())?;
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_profiler_runner_builder() {
let runner = ProfilerRunnerBuilder::default()
.gen_flamegraphs(true)
.build()
.unwrap();
assert_eq!(runner.check_interval_secs, 300);
assert!(runner.gen_flamegraphs);
assert_eq!(runner.report_pct_change_trigger, 10);
assert_eq!(runner.reporting_path, "");
}
#[test]
fn test_profiler_runner_builder_overrides() {
let runner = ProfilerRunnerBuilder::default()
.check_interval_secs(60usize)
.report_pct_change_trigger(25usize)
.reporting_path("profiler_output/")
.expand_frames(true)
.measure_allocated_not_retained(true)
.build()
.unwrap();
assert_eq!(runner.check_interval_secs, 60);
assert_eq!(runner.report_pct_change_trigger, 25);
assert_eq!(runner.reporting_path, "profiler_output/");
assert!(runner.expand_frames);
assert!(runner.measure_allocated_not_retained);
assert!(!runner.gen_flamegraphs);
}
}