use std::alloc::{GlobalAlloc, Layout};
use std::fmt::Write;
use std::hint::black_box;
use std::io::{stderr, Write as IoWrite};
use std::panic::resume_unwind;
use std::process::abort;
use std::sync::atomic::{AtomicBool, Ordering::Relaxed};
use std::sync::mpsc::{channel, RecvTimeoutError};
use std::sync::Arc;
use std::thread;
use std::time::Duration;
use futures::future::join_all;
use moka::sync::Cache;
use rand::distributions::Alphanumeric;
use rand::{rngs::SmallRng, Rng, SeedableRng};
use serial_test::serial;
use ying_profiler::callstack::Measurement;
use ying_profiler::YingProfiler;
#[cfg(test)]
#[global_allocator]
static YING_ALLOC: YingProfiler = YingProfiler::new(5, 64 * 1024 * 1024 * 1024);
const NUM_ALLOCS: usize = 4000;
fn with_watchdog(name: &str, timeout: Duration, body: impl FnOnce() + Send + 'static) {
let (done_tx, done_rx) = channel();
let handle = thread::spawn(move || {
body();
let _ = done_tx.send(());
});
match done_rx.recv_timeout(timeout) {
Ok(()) | Err(RecvTimeoutError::Disconnected) => {
if let Err(panic) = handle.join() {
resume_unwind(panic)
}
}
Err(RecvTimeoutError::Timeout) => {
let mut err = stderr().lock();
let _ = err.write_all(b"\nWATCHDOG TIMEOUT in test: ");
let _ = err.write_all(name.as_bytes());
let _ =
err.write_all(b"\nlikely allocator deadlock; aborting to dump all thread stacks\n");
let _ = err.flush();
abort();
}
}
}
fn alloc_from_distinct_stack(depth: u32, path: u32) -> Vec<u8> {
if depth == 0 {
return vec![7u8; 256];
}
if path & 1 == 0 {
alloc_via_left(depth - 1, path >> 1)
} else {
alloc_via_right(depth - 1, path >> 1)
}
}
#[inline(never)]
fn alloc_via_left(depth: u32, path: u32) -> Vec<u8> {
let v = alloc_from_distinct_stack(depth, path);
black_box(v.len());
v
}
#[inline(never)]
fn alloc_via_right(depth: u32, path: u32) -> Vec<u8> {
let v = alloc_from_distinct_stack(depth, path);
black_box(v.capacity());
v
}
#[test]
#[serial]
fn basic_allocation_free_test() {
thread::sleep(Duration::from_millis(100));
YING_ALLOC.reset_state_for_testing_only();
let mut items: Vec<_> = (0..NUM_ALLOCS).map(|_n| Box::new([0u64; 64])).collect();
let allocated_now = YingProfiler::total_retained_bytes();
println!("allocated_now = {}", allocated_now);
let top_stacks = YING_ALLOC.top_k_stacks_by_allocated(5);
for s in &top_stacks {
println!("---\n{}\n", s.rich_report(&YING_ALLOC, false, true));
}
assert!(!top_stacks.is_empty());
let stat = &top_stacks[0];
assert_eq!(stat.freed_bytes, 0);
let allocated = stat.allocated_bytes;
assert_eq!(allocated / stat.num_allocations, 512);
items.truncate(NUM_ALLOCS / 2);
thread::sleep(Duration::from_millis(100));
let allocated2 = YingProfiler::total_retained_bytes();
println!("allocated2 = {}", allocated2);
assert!(allocated2 < allocated_now);
let top_stacks = YING_ALLOC.top_k_stacks_by_allocated(5);
assert!(!top_stacks.is_empty());
let stat = &top_stacks[0];
println!(
"\n---xxx after dropping xxx---\n{}",
stat.rich_report(&YING_ALLOC, false, true)
);
assert!(stat.freed_bytes > 0);
assert!(stat.retained_profiled_bytes() > 0);
}
#[test]
#[serial]
fn test_giant_allocation() {
thread::sleep(Duration::from_millis(100));
let layout = Layout::from_size_align(128 * 1024 * 1024 * 1024, 8).unwrap();
let ptr = unsafe { YingProfiler::alloc(&YING_ALLOC, layout) };
assert_eq!(ptr as u64, 0);
}
#[test]
#[serial]
fn test_print_allocations_deadlock() {
let _items: Vec<_> = (0..NUM_ALLOCS).map(|_n| Box::new([0u64; 64])).collect();
let top_stacks = YING_ALLOC.top_k_stacks_by_allocated(5);
println!("before potential deadlock");
YING_ALLOC.testing_only_guarantee_next_sample();
for s in &top_stacks {
println!("---\n{}\n", s.rich_report(&YING_ALLOC, false, false));
}
}
#[tokio::test]
#[serial]
async fn stress_test() {
YING_ALLOC.reset_state_for_testing_only();
let dump_allocs_handle = thread::spawn(|| {
for _ in 0..4000 {
let top_stacks = YING_ALLOC.top_k_stacks_by_allocated(10);
let mut report_str = String::new();
for s in &top_stacks {
writeln!(
&mut report_str,
"---\n{}\n",
s.rich_report(&YING_ALLOC, true, false)
)
.unwrap();
}
}
println!("Finished dumping reports...");
});
let cache = Cache::new(10_000);
let num_outer_loops = 50;
let num_inner_loops = 1000;
let rng = SmallRng::from_entropy();
for outer in 0isize..num_outer_loops {
let starting_num = outer * num_inner_loops;
let prev_num = (outer - 1) * num_inner_loops;
let handles: Vec<_> = (0..num_inner_loops)
.map(|n| {
let mut rng = rng.clone();
let cache = cache.clone();
tokio::task::spawn(async move {
let new_str: String =
(0..10).map(|_| rng.sample(Alphanumeric) as char).collect();
cache.insert(starting_num + n, new_str);
})
})
.collect();
if prev_num >= 0 {
for n in 0..1000 {
cache.invalidate(&(prev_num + n));
}
}
join_all(handles).await;
}
println!("Finished alloc/dealloc cycles");
dump_allocs_handle.join().expect("Cannot wait for thread");
let top_stacks = YING_ALLOC.top_k_stacks_by_allocated(5);
for s in &top_stacks {
println!(
"{}",
s.dtrace_report(&YING_ALLOC, Measurement::RetainedBytes)
);
}
assert!(!top_stacks.is_empty());
let total_allocs: u64 = top_stacks.iter().map(|s| s.num_allocations).sum();
let total_frees: u64 = top_stacks.iter().map(|s| s.num_frees).sum();
let total_expected_allocs = num_outer_loops * num_inner_loops / 5;
assert!(total_allocs >= total_expected_allocs as u64);
assert!(total_frees >= (total_expected_allocs * 9 / 10) as u64); }
#[test]
#[serial]
fn many_distinct_stacks_while_reporting() {
with_watchdog(
"many_distinct_stacks_while_reporting",
Duration::from_secs(120),
|| {
YING_ALLOC.reset_state_for_testing_only();
let stop = Arc::new(AtomicBool::new(false));
let reader_stop = stop.clone();
let reader = thread::spawn(move || {
let mut reports = 0u64;
while !reader_stop.load(Relaxed) {
for s in &YING_ALLOC.top_k_stacks_by_retained(20) {
let mut report = String::new();
write!(&mut report, "{}", s.rich_report(&YING_ALLOC, true, true)).unwrap();
}
reports += 1;
}
reports
});
let writers: Vec<_> = (0..4u32)
.map(|t| {
thread::spawn(move || {
for round in 0..8u32 {
for path in 0..1024u32 {
let v = alloc_from_distinct_stack(10, path ^ (t * 7) ^ round);
assert_eq!(v.len(), 256);
}
}
})
})
.collect();
for w in writers {
w.join().expect("writer thread panicked");
}
stop.store(true, Relaxed);
let reports = reader.join().expect("reader thread panicked");
println!(
"Generated {reports} reports over {} distinct stacks",
YING_ALLOC.num_stack_traces()
);
assert!(
YING_ALLOC.num_stack_traces() > 100,
"expected many distinct stacks, got {}",
YING_ALLOC.num_stack_traces()
);
},
);
}
#[test]
#[serial]
fn many_threads_allocating_and_reporting() {
with_watchdog(
"many_threads_allocating_and_reporting",
Duration::from_secs(120),
|| {
YING_ALLOC.reset_state_for_testing_only();
let stop = Arc::new(AtomicBool::new(false));
let reader_stop = stop.clone();
let reader = thread::spawn(move || {
while !reader_stop.load(Relaxed) {
YING_ALLOC.top_k_stacks_by_allocated(10);
YING_ALLOC.num_outstanding_allocs();
YING_ALLOC.symbol_map_size();
}
});
for wave in 0..4u32 {
let workers: Vec<_> = (0..48u32)
.map(|t| {
thread::spawn(move || {
let mut kept = Vec::new();
for i in 0..2000u32 {
kept.push(vec![0u8; 64 + (i % 97) as usize]);
if kept.len() > 64 {
kept.remove(0);
}
if i % 128 == 0 {
thread::sleep(Duration::from_micros(1));
}
}
alloc_from_distinct_stack(6, t ^ wave).len()
})
})
.collect();
for w in workers {
w.join().expect("worker thread panicked");
}
}
stop.store(true, Relaxed);
reader.join().expect("reader thread panicked");
},
);
}
#[test]
#[serial]
fn realloc_churn() {
with_watchdog("realloc_churn", Duration::from_secs(120), || {
YING_ALLOC.reset_state_for_testing_only();
let workers: Vec<_> = (0..8u32)
.map(|_| {
thread::spawn(|| {
for _ in 0..200 {
let mut grown = Vec::new();
for i in 0..2000u32 {
grown.push(i);
}
grown.truncate(10);
grown.shrink_to_fit();
assert_eq!(grown.len(), 10);
}
})
})
.collect();
for w in workers {
w.join().expect("worker thread panicked");
}
let top_stacks = YING_ALLOC.top_k_stacks_by_allocated(5);
assert!(!top_stacks.is_empty());
});
}