use std::alloc::{GlobalAlloc, Layout, System};
use std::cell::Cell;
use std::cmp::min;
use std::env::var;
use std::fmt::Write;
use std::mem::needs_drop;
use std::ptr::{copy_nonoverlapping, null_mut};
use std::sync::atomic::{AtomicUsize, Ordering::Relaxed, Ordering::SeqCst};
use backtrace::Backtrace;
use coarsetime::Clock;
use once_cell::sync::OnceCell;
pub mod callstack;
pub mod histogram;
mod sharded_map;
pub mod utils;
use callstack::{FriendlySymbol, StackStats, StdCallstack};
use sharded_map::ShardedMap;
use utils::{
ProfilerRunner, DEFAULT_CHECK_INTERVAL_SECS, DEFAULT_PCT_CHANGE_TRIGGER, DEFAULT_REPORTING_PATH,
};
const TOP_FRAMES_TO_SKIP: usize = 3;
const DEFAULT_GIANT_ALLOC_LIMIT: usize = 64 * 1024 * 1024 * 1024;
const MAP_CAPACITY_ENV_VAR: &str = "YING_TEST_MAP_CAPACITY";
pub(crate) type SymbolMap = ShardedMap<Vec<FriendlySymbol>>;
pub struct YingProfiler {
sampling_ratio: u32,
single_alloc_limit: usize,
state: OnceCell<YingState>,
}
static TOTAL_RETAINED: AtomicUsize = AtomicUsize::new(0);
static PROFILED_ALLOCATED: AtomicUsize = AtomicUsize::new(0);
static PROFILED_RETAINED: AtomicUsize = AtomicUsize::new(0);
impl YingProfiler {
pub const fn new(sampling_ratio: u32, single_alloc_limit: usize) -> Self {
Self {
sampling_ratio,
single_alloc_limit,
state: OnceCell::new(),
}
}
pub const fn default() -> Self {
Self {
sampling_ratio: 500,
single_alloc_limit: DEFAULT_GIANT_ALLOC_LIMIT,
state: OnceCell::new(),
}
}
pub fn start_profiling(&'static self) -> ProfilerRunner {
let runner = ProfilerRunner::new(
DEFAULT_CHECK_INTERVAL_SECS,
DEFAULT_PCT_CHANGE_TRIGGER,
DEFAULT_REPORTING_PATH,
false,
true,
false,
);
runner.spawn(self);
runner
}
#[inline]
pub fn total_retained_bytes() -> usize {
TOTAL_RETAINED.load(Relaxed)
}
#[inline]
pub fn profiled_bytes_allocated() -> usize {
PROFILED_ALLOCATED.load(Relaxed)
}
#[inline]
pub fn profiled_bytes_retained() -> usize {
PROFILED_RETAINED.load(Relaxed)
}
#[inline]
pub fn symbol_map_size(&self) -> usize {
self.lock_out_profiler(|| self.get_state().symbol_map.len())
}
#[inline]
pub fn num_outstanding_allocs(&self) -> usize {
self.lock_out_profiler(|| self.get_state().outstanding_allocs.len())
}
#[inline]
pub fn num_stack_traces(&self) -> usize {
self.lock_out_profiler(|| self.get_state().stack_stats.len())
}
pub fn top_k_stacks_by_allocated(&self, k: usize) -> Vec<StackStats> {
self.lock_out_profiler(|| {
let stacks_by_alloc = self.stack_list_allocated_bytes_desc();
stacks_by_alloc
.iter()
.take(k)
.filter_map(|&(stack_hash, _bytes_allocated)| {
self.get_stats_for_stack_hash(stack_hash)
})
.collect()
})
}
pub fn top_k_stacks_by_retained(&self, k: usize) -> Vec<StackStats> {
self.lock_out_profiler(|| {
let stacks_by_retained = self.stack_list_retained_bytes_desc();
stacks_by_retained
.iter()
.take(k)
.filter_map(|&(stack_hash, _bytes_retained)| {
self.get_stats_for_stack_hash(stack_hash)
})
.collect()
})
}
fn stack_list_allocated_bytes_desc(&self) -> Vec<(u64, u64)> {
let mut items = self
.get_state()
.stack_stats
.map_to_vec(|stack_hash, stats| (stack_hash, stats.allocated_bytes));
items.sort_unstable_by_key(|&(_, bytes)| std::cmp::Reverse(bytes));
items
}
fn stack_list_retained_bytes_desc(&self) -> Vec<(u64, u64)> {
let mut items = self
.get_state()
.stack_stats
.map_to_vec(|stack_hash, stats| (stack_hash, stats.retained_profiled_bytes()));
items.sort_unstable_by_key(|&(_, bytes)| std::cmp::Reverse(bytes));
items
}
pub fn reset_state_for_testing_only(&self) {
self.lock_out_profiler(|| {
let state = self.get_state();
state.stack_stats.clear();
state.outstanding_allocs.clear();
})
}
pub fn testing_only_guarantee_next_sample(&self) {
THREAD_STATE.with(YingThreadLocal::test_only_reset_sampling_counter)
}
#[inline]
fn check_and_deny_giant_allocations(&self, ptr: *mut u8, layout: Layout) -> *mut u8 {
if layout.size() >= self.single_alloc_limit && self.state.get().is_some() {
self.lock_out_profiler(|| {
println!(
"WARNING: Huge memory allocation of {} bytes denied by Ying profiler",
layout.size()
);
let mut bt = Backtrace::new_unresolved();
let stack = StdCallstack::from_backtrace_unresolved(&bt);
let state = self.get_state();
stack.populate_symbol_map(&mut bt, &state.symbol_map);
println!(
"Stack trace:\n{}",
stack.with_symbols_and_filename(&state.symbol_map, true)
);
});
null_mut::<u8>()
} else {
ptr
}
}
#[inline]
fn get_state(&self) -> &YingState {
self.lock_out_profiler(|| self.state.get_or_init(YingState::new))
}
#[inline]
fn get_stats_for_stack_hash(&self, stack_hash: u64) -> Option<StackStats> {
self.get_state().stack_stats.get_cloned(stack_hash)
}
#[inline]
fn lock_out_profiler<R>(&self, func: impl FnOnce() -> R) -> R {
THREAD_STATE.with(YingThreadLocal::set_allocator_lock);
let return_val = func();
exit_profiler();
return_val
}
}
struct YingState {
symbol_map: SymbolMap,
stack_stats: ShardedMap<StackStats>,
outstanding_allocs: ShardedMap<(u64, u64)>,
}
impl YingState {
pub fn new() -> Self {
let capacity_scale = var(MAP_CAPACITY_ENV_VAR)
.ok()
.and_then(|s| s.parse::<usize>().ok());
let symbol_map = SymbolMap::with_capacity(capacity_scale.unwrap_or(1000));
let stack_stats = ShardedMap::with_capacity(capacity_scale.unwrap_or(1000));
let outstanding_allocs = ShardedMap::with_capacity(capacity_scale.unwrap_or(5000));
Self {
symbol_map,
stack_stats,
outstanding_allocs,
}
}
}
thread_local! {
static THREAD_STATE: YingThreadLocal = const { YingThreadLocal::new() };
}
const _: () = assert!(
!needs_drop::<YingThreadLocal>(),
"YingThreadLocal must not need Drop, or thread local setup will allocate inside the allocator"
);
struct YingThreadLocal {
alloc_lock: Cell<u32>,
sample_count: Cell<u32>,
}
impl YingThreadLocal {
const fn new() -> Self {
Self {
alloc_lock: Cell::new(0),
sample_count: Cell::new(0),
}
}
#[inline]
fn is_allocator_locked(&self) -> bool {
self.alloc_lock.get() > 0
}
#[inline]
fn set_allocator_lock(&self) {
self.alloc_lock.set(self.alloc_lock.get().saturating_add(1));
}
#[inline]
fn release_allocator_lock(&self) {
self.alloc_lock.set(self.alloc_lock.get().saturating_sub(1));
}
#[inline]
fn should_sample(&self, ratio: u32) -> bool {
let count = self.sample_count.get().wrapping_add(1);
self.sample_count.set(count);
count % ratio == 0
}
#[inline]
fn test_only_reset_sampling_counter(&self) {
self.sample_count.set(0);
}
}
#[inline]
fn try_enter_profiler(sampling_ratio: Option<u32>) -> bool {
THREAD_STATE.with(|tl| {
if tl.is_allocator_locked() {
return false;
}
if let Some(ratio) = sampling_ratio {
if !tl.should_sample(ratio) {
return false;
}
}
tl.set_allocator_lock();
true
})
}
#[inline]
fn exit_profiler() {
THREAD_STATE.with(YingThreadLocal::release_allocator_lock)
}
unsafe impl GlobalAlloc for YingProfiler {
unsafe fn alloc(&self, layout: Layout) -> *mut u8 {
let alloc_ptr = self.check_and_deny_giant_allocations(System.alloc(layout), layout);
if !alloc_ptr.is_null() {
TOTAL_RETAINED.fetch_add(layout.size(), SeqCst);
if try_enter_profiler(Some(self.sampling_ratio)) {
PROFILED_ALLOCATED.fetch_add(layout.size(), SeqCst);
PROFILED_RETAINED.fetch_add(layout.size(), SeqCst);
let mut bt = Backtrace::new_unresolved();
let stack = StdCallstack::from_backtrace_unresolved(&bt);
let stack_hash = stack.compute_hash();
let state = self.get_state();
state.stack_stats.update_or_insert_with(
stack_hash,
|stats| {
stats.num_allocations += 1;
stats.allocated_bytes += layout.size() as u64;
},
|| {
stack.populate_symbol_map(&mut bt, &state.symbol_map);
StackStats::new(stack, Some(layout.size() as u64))
},
);
state.outstanding_allocs.update_or_insert_with(
alloc_ptr as u64,
|_existing| {},
|| (stack_hash, Clock::recent_since_epoch().as_millis()),
);
exit_profiler();
}
}
alloc_ptr
}
unsafe fn dealloc(&self, ptr: *mut u8, layout: Layout) {
System.dealloc(ptr, layout);
TOTAL_RETAINED.fetch_sub(layout.size(), SeqCst);
if self.state.get().is_none() {
return;
}
if !try_enter_profiler(None) {
return;
}
let state = self.get_state();
if state.outstanding_allocs.contains_key(ptr as u64) {
if let Some((stack_hash, alloc_ts)) = state.outstanding_allocs.remove(ptr as u64) {
PROFILED_RETAINED.fetch_sub(layout.size(), SeqCst);
let alloc_time_ms = Clock::recent_since_epoch()
.as_millis()
.saturating_sub(alloc_ts);
state.stack_stats.update(stack_hash, |stats| {
stats.update_free_stats(layout.size() as u64, alloc_time_ms)
});
}
}
exit_profiler();
}
unsafe fn realloc(&self, ptr: *mut u8, layout: Layout, new_size: usize) -> *mut u8 {
let old_size = layout.size();
let new_layout = Layout::from_size_align_unchecked(new_size, layout.align());
let new_ptr = self.check_and_deny_giant_allocations(System.alloc(new_layout), new_layout);
if !new_ptr.is_null() {
copy_nonoverlapping(ptr, new_ptr, min(old_size, new_size));
System.dealloc(ptr, layout);
if new_size > old_size {
TOTAL_RETAINED.fetch_add(new_size - old_size, SeqCst);
} else {
TOTAL_RETAINED.fetch_sub(old_size - new_size, SeqCst);
}
if self.state.get().is_none() {
return new_ptr;
}
if !try_enter_profiler(None) {
return new_ptr;
}
let state = self.get_state();
if state.outstanding_allocs.contains_key(ptr as u64) {
if let Some((stack_hash, alloc_ts)) = state.outstanding_allocs.remove(ptr as u64) {
if new_size > old_size {
PROFILED_RETAINED.fetch_add(new_size - old_size, SeqCst);
} else {
PROFILED_RETAINED.fetch_sub(old_size - new_size, SeqCst);
}
state
.outstanding_allocs
.insert(new_ptr as u64, (stack_hash, alloc_ts));
state.stack_stats.update(stack_hash, |stats| {
if new_size > old_size {
stats.allocated_bytes += (new_size - old_size) as u64;
} else {
stats.allocated_bytes -= (old_size - new_size) as u64;
}
});
}
}
exit_profiler();
}
new_ptr
}
}