1use std::io::Write as _;
2use std::path::PathBuf;
3use std::thread;
4use std::time::Duration;
5use std::{
6 fs::{create_dir_all, File},
7 io::Cursor,
8};
9
10use callstack::Measurement;
11use chrono::{offset::Local, SecondsFormat};
12use derive_builder::Builder;
13use inferno::collapse::{dtrace, Collapse};
14use inferno::flamegraph;
15use log::{error, info};
16
17use super::*;
18
19#[derive(Clone, Debug, PartialEq, Builder)]
55#[builder(setter(into))]
56pub struct ProfilerRunner {
57 #[builder(default = "DEFAULT_CHECK_INTERVAL_SECS")]
59 check_interval_secs: usize,
60 #[builder(default = "DEFAULT_PCT_CHANGE_TRIGGER")]
62 report_pct_change_trigger: usize,
63 #[builder(default)]
65 reporting_path: String,
66 #[builder(default = "false")]
68 expand_frames: bool,
69 #[builder(default = "false")]
71 gen_flamegraphs: bool,
72 #[builder(default = "false")]
74 measure_allocated_not_retained: bool,
75}
76
77const INITIAL_RETAINED_MEM_MB: usize = 20;
78
79pub const DEFAULT_REPORTING_PATH: &str = "ying-profiles";
82
83pub const DEFAULT_CHECK_INTERVAL_SECS: usize = 300;
85
86pub const DEFAULT_PCT_CHANGE_TRIGGER: usize = 10;
88
89impl Default for ProfilerRunner {
92 fn default() -> Self {
93 Self::new(
94 DEFAULT_CHECK_INTERVAL_SECS,
95 DEFAULT_PCT_CHANGE_TRIGGER,
96 "",
97 false,
98 false,
99 false,
100 )
101 }
102}
103
104impl ProfilerRunner {
105 pub fn new(
107 check_interval_secs: usize,
108 report_pct_change_trigger: usize,
109 reporting_path: &str,
110 expand_frames: bool,
111 gen_flamegraphs: bool,
112 measure_allocated_not_retained: bool,
113 ) -> Self {
114 Self {
115 check_interval_secs,
116 report_pct_change_trigger,
117 reporting_path: reporting_path.to_string(),
118 expand_frames,
119 gen_flamegraphs,
120 measure_allocated_not_retained,
121 }
122 }
123
124 pub fn spawn(&self, profiler: &'static YingProfiler) {
127 if !self.reporting_path.is_empty() {
128 if let Err(e) = create_dir_all(&self.reporting_path) {
129 error!(
130 "Ying: could not create reporting directory {:?}, reports will not be written: {}",
131 self.reporting_path, e
132 );
133 }
134 }
135
136 let check_interval_secs = self.check_interval_secs;
137 let report_pct_change_trigger = self.report_pct_change_trigger;
138 let reporting_path = PathBuf::from(self.reporting_path.clone());
139 let expand_frames = self.expand_frames;
140 let profiler2 = profiler;
141 let measurement = if self.measure_allocated_not_retained {
142 Measurement::AllocatedBytes
143 } else {
144 Measurement::RetainedBytes
145 };
146 let gen_flamegraphs = self.gen_flamegraphs;
147
148 thread::spawn(move || {
149 let mut last_retained_mem = INITIAL_RETAINED_MEM_MB as f64;
150
151 loop {
152 thread::sleep(Duration::from_secs(check_interval_secs as u64));
153
154 let new_allocated = YingProfiler::total_retained_bytes() as f64 / (1024.0 * 1024.0);
156 let ratio = (new_allocated - last_retained_mem) / last_retained_mem;
157
158 info!(
159 "Ying: total allocated memory is {:.2} MB and ratio to last = {}",
160 new_allocated, ratio
161 );
162
163 if (ratio.abs() * 100.0) >= report_pct_change_trigger as f64 {
165 #[cfg(feature = "profile-spans")]
166 info!(
167 "Significant memory change registered, dumping profile: new = {}, old = {}",
168 new_allocated, last_retained_mem
169 );
170
171 last_retained_mem = new_allocated;
172
173 let top_stacks = match measurement {
174 Measurement::AllocatedBytes => profiler2.top_k_stacks_by_allocated(10),
175 Measurement::RetainedBytes => profiler2.top_k_stacks_by_retained(10),
176 };
177 for s in &top_stacks {
178 println!("---\n{}\n", s.rich_report(profiler2, false, expand_frames));
180 }
181
182 let dt = Local::now();
184 let dt_str = dt.to_rfc3339_opts(SecondsFormat::Secs, true);
185 let dump_name = format!("ying.{}.{}MB.report", dt_str, new_allocated as i64);
186
187 let mut report_path = reporting_path.clone();
188 report_path.push(dump_name);
189 if let Ok(f) = File::create(&report_path) {
190 for s in &top_stacks {
191 let _ = writeln!(
192 &f,
193 "---\n{}\n",
194 s.rich_report(profiler2, false, expand_frames)
195 );
196 }
197 } else {
198 error!("Error: could not write memory report to {:?}", report_path);
199 }
200
201 if gen_flamegraphs {
202 let graph_name = format!("ying.{}.{}MB.svg", dt_str, new_allocated as i64);
203 let mut graph_path = reporting_path.clone();
204 graph_path.push(graph_name);
205 if let Err(e) = gen_flamegraph(profiler2, measurement, &graph_path) {
206 error!("Error generating flame graph to {:?}: {}", graph_path, e);
207 }
208 }
209 }
210 }
211 });
212 }
213}
214
215pub fn gen_flamegraph(
219 profiler: &YingProfiler,
220 measurement: Measurement,
221 path: &PathBuf,
222) -> Result<(), String> {
223 let mut report = String::new();
225 match measurement {
226 Measurement::RetainedBytes => {
227 let top_stacks = profiler.top_k_stacks_by_retained(50);
228 for s in &top_stacks {
229 writeln!(
230 &mut report,
231 "{}",
232 s.dtrace_report(profiler, Measurement::RetainedBytes)
233 )
234 .map_err(|e| e.to_string())?;
235 }
236 }
237 Measurement::AllocatedBytes => {
238 let top_stacks = profiler.top_k_stacks_by_allocated(50);
239 for s in &top_stacks {
240 writeln!(
241 &mut report,
242 "{}",
243 s.dtrace_report(profiler, Measurement::AllocatedBytes)
244 )
245 .map_err(|e| e.to_string())?;
246 }
247 }
248 }
249
250 let mut folder = dtrace::Folder::default();
252 let mut folded_buf = Vec::new();
253 let folded_out = Cursor::new(&mut folded_buf);
254 folder
255 .collapse(Cursor::new(report), folded_out)
256 .map_err(|e| e.to_string())?;
257
258 if let Ok(f) = File::create(path) {
260 flamegraph::from_reader(
261 &mut flamegraph::Options::default(),
262 Cursor::new(&folded_buf),
263 f,
264 )
265 .map_err(|e| e.to_string())?;
266 }
267 Ok(())
268}
269
270#[cfg(test)]
271mod tests {
272 use super::*;
273
274 #[test]
275 fn test_profiler_runner_builder() {
276 let runner = ProfilerRunnerBuilder::default()
277 .gen_flamegraphs(true)
278 .build()
279 .unwrap();
280 assert_eq!(runner.check_interval_secs, 300);
281 assert!(runner.gen_flamegraphs);
282 assert_eq!(runner.report_pct_change_trigger, 10);
283 assert_eq!(runner.reporting_path, "");
284 }
285
286 #[test]
287 fn test_profiler_runner_builder_overrides() {
288 let runner = ProfilerRunnerBuilder::default()
289 .check_interval_secs(60usize)
290 .report_pct_change_trigger(25usize)
291 .reporting_path("profiler_output/")
292 .expand_frames(true)
293 .measure_allocated_not_retained(true)
294 .build()
295 .unwrap();
296 assert_eq!(runner.check_interval_secs, 60);
297 assert_eq!(runner.report_pct_change_trigger, 25);
298 assert_eq!(runner.reporting_path, "profiler_output/");
299 assert!(runner.expand_frames);
300 assert!(runner.measure_allocated_not_retained);
301 assert!(!runner.gen_flamegraphs);
302 }
303}