Skip to main content

ying_profiler/
utils.rs

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/// A background thread that dumps out stats, flamegraphs, and does other periodic cleanup.
20/// 1. Dumps out top retained memory stats to both logs and disk
21/// 2. Can optionally dump out flamegraphs
22/// 3. Lets you choose between dumping retained or allocated reports/flamegraphs
23///
24/// The runner thread will check the amount of retained memory as estimated by this profiler,
25/// every `check_interval_secs`.  If the retained memory changes from the previous time by more than
26/// `report_pct_change_trigger`, then it will dump out a memory report of the top either retained
27/// or allocated stack traces as a file to the chosen `reporting_path` directory on disk.  The file will
28/// have the ISO8601 timestamp and the amount of retained memory in the filename for convenience.
29/// Optionally, a flamegraph will also be dumped.
30///
31/// Note that ProfilerRunner uses log framework to periodically dump out logs.  The app is responsible
32/// for initializing the logging infrastructure.
33///
34/// To run.  `no_run` because spawning the reporting thread and creating its output directory are
35/// side effects that do not belong in a doctest:
36/// ```no_run
37///     use ying_profiler::{YingProfiler, utils::ProfilerRunner};
38///     static YING_ALLOC: YingProfiler = YingProfiler::default();
39///     ProfilerRunner::default().spawn(&YING_ALLOC);
40/// ```
41///
42/// Builder pattern can also be used.
43/// ```no_run
44///     use ying_profiler::{YingProfiler, utils::ProfilerRunnerBuilder};
45///     static YING_ALLOC: YingProfiler = YingProfiler::default();
46///     let runner = ProfilerRunnerBuilder::default()
47///         .check_interval_secs(60usize)
48///         .gen_flamegraphs(true)
49///         .reporting_path("profiler_output/")
50///         .build()
51///         .unwrap();
52///      runner.spawn(&YING_ALLOC);
53/// ```
54#[derive(Clone, Debug, PartialEq, Builder)]
55#[builder(setter(into))]
56pub struct ProfilerRunner {
57    /// Number of seconds in between memory checks
58    #[builder(default = "DEFAULT_CHECK_INTERVAL_SECS")]
59    check_interval_secs: usize,
60    /// Percent change in retained memory to trigger a report
61    #[builder(default = "DEFAULT_PCT_CHANGE_TRIGGER")]
62    report_pct_change_trigger: usize,
63    /// Path to write top retained memory reports to
64    #[builder(default)]
65    reporting_path: String,
66    /// Expand inlined call stack symbols with a > ?
67    #[builder(default = "false")]
68    expand_frames: bool,
69    /// Generate flamegraphs at reporting_path
70    #[builder(default = "false")]
71    gen_flamegraphs: bool,
72    /// True=measure allocated memory instead of False=measure retained memory
73    #[builder(default = "false")]
74    measure_allocated_not_retained: bool,
75}
76
77const INITIAL_RETAINED_MEM_MB: usize = 20;
78
79/// Directory that [`YingProfiler::start_profiling`](crate::YingProfiler::start_profiling) writes
80/// reports and flamegraphs to.
81pub const DEFAULT_REPORTING_PATH: &str = "ying-profiles";
82
83/// Default seconds between memory checks.
84pub const DEFAULT_CHECK_INTERVAL_SECS: usize = 300;
85
86/// Default percent change in memory that triggers a report.
87pub const DEFAULT_PCT_CHANGE_TRIGGER: usize = 10;
88
89/// Creates a new ProfilerRunner with default values.  Writes reports to current directory, does not expand frames,
90/// every 5 minute checks on memory, 10% change triggers report.  No flamegraphs, retained memory.
91impl 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    /// Creates a new ProfilerRunner with specific parameters
106    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    /// Spawn a new background thread to run profiler and get stats.
125    /// Creates `reporting_path` if it does not already exist.
126    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                // Check and compare memory
155                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                // Threshold for change exceeded, do report
164                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                        // In case the app does not use log, we still output to STDOUT the report
179                        println!("---\n{}\n", s.rich_report(profiler2, false, expand_frames));
180                    }
181
182                    // Formulate profiling filename based on ISO8601 timestamp and number of MBs
183                    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
215/// Function to produce a FlameGraph file to a specific path.
216/// Specify whether to measure retained or allocated bytes, and the path to write flamegraph file to
217/// (should probably end in .svg)
218pub fn gen_flamegraph(
219    profiler: &YingProfiler,
220    measurement: Measurement,
221    path: &PathBuf,
222) -> Result<(), String> {
223    // Generate dtrace-compatible output
224    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    // Fold/collapse output to folded lines
251    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    // Now, generate the flamegraph from folded lines
259    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}