use std::env;
use std::time::{Duration, Instant};
use futures::future::join_all;
use moka::sync::Cache;
use rand::distributions::Alphanumeric;
use rand::{rngs::SmallRng, Rng, SeedableRng};
#[cfg(feature = "profile-spans")]
use tracing::instrument;
use ying_profiler::utils::{ProfilerRunnerBuilder, DEFAULT_REPORTING_PATH};
use ying_profiler::YingProfiler;
#[global_allocator]
static YING_ALLOC: YingProfiler = YingProfiler::new(5, 64 * 1024 * 1024 * 1024);
const CACHE_MAX_ENTRIES: u64 = 2_000;
const TASKS_PER_ROUND: u64 = 1_000;
const ROUND_SLEEP: Duration = Duration::from_secs(2);
fn leak_bytes_per_round(interval_secs: usize) -> usize {
(4 * 1024 * 1024 / (interval_secs / 2).max(1)).max(512 * 1024)
}
fn env_val<T: std::str::FromStr>(name: &str, default: T) -> T {
env::var(name)
.ok()
.and_then(|v| v.parse().ok())
.unwrap_or(default)
}
#[tokio::main]
async fn main() {
tracing_subscriber::fmt::init();
let interval_secs: usize = env_val("YING_EXAMPLE_INTERVAL_SECS", 300);
let runtime_secs: usize = env_val("YING_EXAMPLE_RUNTIME_SECS", interval_secs * 2 + 90);
let reporting_path: String =
env::var("YING_EXAMPLE_REPORTING_PATH").unwrap_or_else(|_| DEFAULT_REPORTING_PATH.into());
let runner = ProfilerRunnerBuilder::default()
.check_interval_secs(interval_secs)
.report_pct_change_trigger(10usize)
.reporting_path(reporting_path)
.gen_flamegraphs(true)
.build()
.unwrap();
runner.spawn(&YING_ALLOC);
println!(
"Ying example running for {runtime_secs}s, checking every {interval_secs}s; \
a report and flamegraph are written whenever retained memory moves 10%."
);
let _leak = cache_update_loop(
Duration::from_secs(runtime_secs as u64),
leak_bytes_per_round(interval_secs),
)
.await;
println!("\n===== Done. Final profiler stats =====");
println!(
"Total bytes retained: {}",
YingProfiler::total_retained_bytes()
);
println!(
"Profiled bytes allocated: {}",
YingProfiler::profiled_bytes_allocated()
);
println!("Size of symbol map: {}", YING_ALLOC.symbol_map_size());
println!("\n===== Top stacks by TOTAL ALLOCATED (churn should dominate) =====");
for s in &YING_ALLOC.top_k_stacks_by_allocated(5) {
println!("---\n{}\n", s.rich_report(&YING_ALLOC, false, false));
}
println!("\n===== Top stacks by RETAINED (the leak should dominate) =====");
for s in &YING_ALLOC.top_k_stacks_by_retained(5) {
println!("---\n{}\n", s.rich_report(&YING_ALLOC, false, false));
}
}
async fn cache_update_loop(runtime: Duration, leak_bytes: usize) -> Vec<Vec<u8>> {
let cache: Cache<String, String> = Cache::new(CACHE_MAX_ENTRIES);
let mut leak: Vec<Vec<u8>> = Vec::new();
let mut leaked_bytes = 0usize;
let rng = SmallRng::from_entropy();
let deadline = Instant::now() + runtime;
let mut round = 0u64;
while Instant::now() < deadline {
round += 1;
let handles: Vec<_> = (0..TASKS_PER_ROUND)
.map(|_n| {
let mut rng = rng.clone();
let cache = cache.clone();
tokio::task::spawn(async move {
let new_str: String = (0..(1024 * (1 + rng.gen_range(0..4))))
.map(|_| rng.sample(Alphanumeric) as char)
.collect();
insert_one(&cache, new_str).await;
})
})
.collect();
join_all(handles).await;
leak.push(vec![0u8; leak_bytes]);
leaked_bytes += leak_bytes;
println!(
"round {round}: leaked {} MB so far",
leaked_bytes / (1024 * 1024)
);
tokio::time::sleep(ROUND_SLEEP).await;
}
leak
}
#[cfg(feature = "profile-spans")]
#[instrument(level = "info", skip_all)]
async fn insert_one(cache: &Cache<String, String>, s: String) {
cache.insert(s.clone(), s);
}
#[cfg(not(feature = "profile-spans"))]
async fn insert_one(cache: &Cache<String, String>, s: String) {
cache.insert(s.clone(), s);
}