#![allow(unsafe_code)]
use std::alloc::Layout;
use std::ptr::NonNull;
use kovan_queue::array_queue::ArrayQueue;
use crate::sync::{Arc, AtomicUsize, Mutex, Ordering};
pub(crate) const CHUNK_ALIGN: usize = 16;
const MIN_CHUNK: usize = 64;
const MAX_CHUNKS_PER_CLASS: usize = 256;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct ArenaProfile {
pub initial_chunk_size: usize,
pub max_chunk_size: usize,
}
impl ArenaProfile {
pub const SERVER: Self = Self {
initial_chunk_size: 64 * 1024,
max_chunk_size: 1024 * 1024,
};
pub const EMBEDDED: Self = Self {
initial_chunk_size: 4 * 1024,
max_chunk_size: 64 * 1024,
};
pub fn is_valid(&self) -> bool {
self.initial_chunk_size.is_power_of_two()
&& self.max_chunk_size.is_power_of_two()
&& self.initial_chunk_size <= self.max_chunk_size
}
}
impl Default for ArenaProfile {
fn default() -> Self {
Self::SERVER
}
}
fn prev_power_of_two(n: usize) -> usize {
if n == 0 {
return 1;
}
1usize << (usize::BITS - 1 - n.leading_zeros())
}
fn base_class(profile: ArenaProfile, budget: usize) -> usize {
let base = profile
.initial_chunk_size
.min(prev_power_of_two(budget))
.max(MIN_CHUNK);
base.min(profile.max_chunk_size.max(MIN_CHUNK))
}
struct ChunkHandle {
ptr: NonNull<u8>,
size: usize,
}
unsafe impl Send for ChunkHandle {}
unsafe impl Sync for ChunkHandle {}
impl ChunkHandle {
fn alloc(size: usize) -> Option<Self> {
let layout = Layout::from_size_align(size, CHUNK_ALIGN).ok()?;
let ptr = unsafe { std::alloc::alloc(layout) };
NonNull::new(ptr).map(|ptr| Self { ptr, size })
}
}
impl Drop for ChunkHandle {
fn drop(&mut self) {
unsafe {
let layout = Layout::from_size_align_unchecked(self.size, CHUNK_ALIGN);
std::alloc::dealloc(self.ptr.as_ptr(), layout);
}
}
}
pub(crate) struct ChunkPool {
classes: Box<[ArrayQueue<ChunkHandle>]>,
base: usize,
parked: AtomicUsize,
budget: usize,
#[cfg(test)]
misses: AtomicUsize,
}
impl ChunkPool {
pub(crate) fn new(
profile: ArenaProfile,
write_buffer_size: usize,
max_write_buffer_number: usize,
) -> Self {
let profile = if profile.is_valid() {
profile
} else {
ArenaProfile::SERVER
};
let base = base_class(profile, write_buffer_size);
let max = profile.max_chunk_size.max(base);
let budget = write_buffer_size.saturating_mul(max_write_buffer_number.max(1));
let class_count = (max.trailing_zeros() - base.trailing_zeros()) as usize + 1;
let classes: Vec<ArrayQueue<ChunkHandle>> = (0..class_count)
.map(|i| {
let size = base << i;
let cap = (budget / size).clamp(1, MAX_CHUNKS_PER_CLASS);
ArrayQueue::new(cap)
})
.collect();
Self {
classes: classes.into_boxed_slice(),
base,
parked: AtomicUsize::new(0),
budget,
#[cfg(test)]
misses: AtomicUsize::new(0),
}
}
fn class_index(&self, size: usize) -> Option<usize> {
if !size.is_power_of_two() || size < self.base {
return None;
}
let idx = (size.trailing_zeros() - self.base.trailing_zeros()) as usize;
(idx < self.classes.len()).then_some(idx)
}
fn take(&self, size: usize) -> Option<ChunkHandle> {
let chunk = self
.class_index(size)
.and_then(|idx| self.classes[idx].pop());
match chunk {
Some(chunk) => {
self.parked.fetch_sub(chunk.size, Ordering::Relaxed);
Some(chunk)
}
None => {
#[cfg(test)]
self.misses.fetch_add(1, Ordering::Relaxed);
None
}
}
}
#[cfg(test)]
fn misses(&self) -> usize {
self.misses.load(Ordering::Relaxed)
}
fn give(&self, chunk: ChunkHandle) {
let Some(idx) = self.class_index(chunk.size) else {
return;
};
let size = chunk.size;
let mut parked = self.parked.load(Ordering::Relaxed);
loop {
if parked.saturating_add(size) > self.budget {
return;
}
match self.parked.compare_exchange_weak(
parked,
parked + size,
Ordering::AcqRel,
Ordering::Relaxed,
) {
Ok(_) => break,
Err(observed) => parked = observed,
}
}
if self.classes[idx].push(chunk).is_err() {
self.parked.fetch_sub(size, Ordering::Relaxed);
}
}
pub(crate) fn parked_bytes(&self) -> usize {
self.parked.load(Ordering::Relaxed)
}
pub(crate) fn budget(&self) -> usize {
self.budget
}
}
struct ArenaState {
chunks: Vec<ChunkHandle>,
offset: usize,
}
pub(crate) struct Arena {
state: Mutex<ArenaState>,
reserved: AtomicUsize,
used: AtomicUsize,
pool: Arc<ChunkPool>,
budget: usize,
base: usize,
max_chunk_size: usize,
}
unsafe impl Send for Arena {}
unsafe impl Sync for Arena {}
impl Arena {
pub(crate) fn new(pool: Arc<ChunkPool>, budget: usize, profile: ArenaProfile) -> Self {
let profile = if profile.is_valid() {
profile
} else {
ArenaProfile::SERVER
};
let base = base_class(profile, budget);
Self {
state: Mutex::new(ArenaState {
chunks: Vec::new(),
offset: 0,
}),
reserved: AtomicUsize::new(0),
used: AtomicUsize::new(0),
pool,
budget,
base,
max_chunk_size: profile.max_chunk_size.max(base),
}
}
pub(crate) fn alloc(&self, size: usize, align: usize) -> Option<NonNull<u8>> {
debug_assert!(align.is_power_of_two() && align <= CHUNK_ALIGN);
if size == 0 {
return None;
}
let mut state = self.state.lock();
if let Some(ptr) = Self::bump(&mut state, size, align) {
self.used.fetch_add(size, Ordering::Relaxed);
return Some(ptr);
}
let chunk_size = self.next_chunk_size(state.chunks.len(), size);
let chunk = match self.pool.take(chunk_size) {
Some(chunk) => chunk,
None => ChunkHandle::alloc(chunk_size)?,
};
self.reserved.fetch_add(chunk.size, Ordering::Relaxed);
state.chunks.push(chunk);
state.offset = 0;
let ptr = Self::bump(&mut state, size, align)?;
self.used.fetch_add(size, Ordering::Relaxed);
Some(ptr)
}
fn bump(state: &mut ArenaState, size: usize, align: usize) -> Option<NonNull<u8>> {
let chunk = state.chunks.last()?;
let base = chunk.ptr.as_ptr() as usize;
let start = base.checked_add(state.offset)?.next_multiple_of(align);
let end = start.checked_add(size)?;
if end > base.checked_add(chunk.size)? {
return None;
}
state.offset = end - base;
let ptr = unsafe { chunk.ptr.as_ptr().add(start - base) };
NonNull::new(ptr)
}
fn next_chunk_size(&self, chunk_index: usize, need: usize) -> usize {
if need > self.max_chunk_size {
return need.next_multiple_of(CHUNK_ALIGN);
}
let max_shift = self.base.leading_zeros() as usize;
let class = (self.base << chunk_index.min(max_shift)).min(self.max_chunk_size);
let remaining = self
.budget
.saturating_sub(self.reserved.load(Ordering::Relaxed));
let mut size = class;
while size > self.base && size > remaining {
size >>= 1;
}
size.max(need.next_power_of_two()).min(self.max_chunk_size)
}
pub(crate) fn reserved_bytes(&self) -> usize {
self.reserved.load(Ordering::Relaxed)
}
pub(crate) fn used_bytes(&self) -> usize {
self.used.load(Ordering::Relaxed)
}
#[cfg(test)]
pub(crate) fn chunk_count(&self) -> usize {
self.state.lock().chunks.len()
}
#[cfg(test)]
pub(crate) fn chunk_sizes(&self) -> Vec<usize> {
self.state.lock().chunks.iter().map(|c| c.size).collect()
}
}
impl Drop for Arena {
fn drop(&mut self) {
let state = self.state.get_mut();
for chunk in state.chunks.drain(..) {
self.pool.give(chunk);
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use proptest::prelude::*;
fn pool(profile: ArenaProfile, budget: usize, count: usize) -> Arc<ChunkPool> {
Arc::new(ChunkPool::new(profile, budget, count))
}
#[cfg(not(miri))]
const LADDER_BUDGET: usize = 64 * 1024 * 1024;
#[cfg(miri)]
const LADDER_BUDGET: usize = 8 * 1024 * 1024;
#[cfg(not(miri))]
const RECYCLE_CYCLES: usize = 50;
#[cfg(miri)]
const RECYCLE_CYCLES: usize = 8;
#[test]
fn empty_arena_reserves_nothing() {
let arena = Arena::new(
pool(ArenaProfile::SERVER, 64 * 1024, 2),
64 * 1024,
ArenaProfile::SERVER,
);
assert_eq!(arena.reserved_bytes(), 0);
assert_eq!(arena.used_bytes(), 0);
assert_eq!(arena.chunk_count(), 0);
}
#[test]
fn alloc_is_aligned_and_in_range() {
let profile = ArenaProfile::EMBEDDED;
let arena = Arena::new(pool(profile, 64 * 1024, 2), 64 * 1024, profile);
let mut seen: Vec<(usize, usize)> = Vec::new();
for i in 1..200usize {
let size = i * 7;
let ptr = arena.alloc(size, 8).expect("arena has budget");
assert_eq!(ptr.as_ptr() as usize % 8, 0, "8-aligned");
let start = ptr.as_ptr() as usize;
for (other, other_len) in &seen {
assert!(
start + size <= *other || *other + *other_len <= start,
"allocations must not overlap"
);
}
seen.push((start, size));
}
assert!(arena.used_bytes() >= seen.iter().map(|(_, l)| l).sum::<usize>());
}
#[test]
fn chunk_sizes_follow_the_growth_rule() {
let profile = ArenaProfile::EMBEDDED;
let budget = 256 * 1024;
let arena = Arena::new(pool(profile, budget, 2), budget, profile);
while arena.used_bytes() < budget {
arena.alloc(1024, 8).expect("allocator");
}
let sizes = arena.chunk_sizes();
assert_eq!(
sizes,
vec![4096, 8192, 16384, 32768, 65536, 65536, 65536, 4096],
"4 KiB doubling up to the 64 KiB cap, stepped down at the budget"
);
assert_eq!(sizes.iter().sum::<usize>(), budget);
}
#[test]
fn the_class_ladder_saturates_instead_of_wrapping() {
let profile = ArenaProfile::SERVER;
let budget = 64 * 1024 * 1024;
let arena = Arena::new(pool(profile, budget, 2), budget, profile);
for index in [0usize, 4, 47, 48, 64, 200, usize::MAX] {
let size = arena.next_chunk_size(index, 512);
assert!(
size >= arena.base && size <= profile.max_chunk_size,
"chunk {index} sized {size}, outside [{}, {}]",
arena.base,
profile.max_chunk_size
);
}
assert_eq!(arena.next_chunk_size(48, 512), profile.max_chunk_size);
}
#[test]
fn a_large_budget_keeps_full_size_chunks_all_the_way_up() {
let profile = ArenaProfile::SERVER;
let budget = LADDER_BUDGET;
let arena = Arena::new(pool(profile, budget, 2), budget, profile);
while arena.used_bytes() < budget {
arena.alloc(4096, 8).expect("allocator");
}
let sizes = arena.chunk_sizes();
assert!(
sizes.len() < 80,
"a 64 MiB memtable must be a few dozen chunks, got {}",
sizes.len()
);
let at_cap = sizes
.iter()
.filter(|size| **size == profile.max_chunk_size)
.count();
assert!(
at_cap >= sizes.len() - 6,
"all but the ladder's first few chunks must be at the cap: {sizes:?}"
);
}
#[test]
fn growth_never_overshoots_the_budget_by_more_than_one_class() {
for budget in [1usize, 100, 4096, 6000, 1024 * 1024, LADDER_BUDGET] {
let profile = ArenaProfile::SERVER;
let arena = Arena::new(pool(profile, budget, 2), budget, profile);
while arena.used_bytes() < budget {
arena.alloc(512, 8).expect("allocator");
}
let bound = budget + profile.max_chunk_size;
assert!(
arena.reserved_bytes() <= bound,
"budget {budget}: reserved {} > {bound}",
arena.reserved_bytes()
);
}
}
#[test]
fn a_small_budget_does_not_reserve_a_full_initial_chunk() {
let profile = ArenaProfile::SERVER;
let arena = Arena::new(pool(profile, 4096, 2), 4096, profile);
arena.alloc(64, 8).expect("allocator");
assert_eq!(arena.reserved_bytes(), 4096);
}
#[test]
fn oversized_entry_gets_a_dedicated_chunk() {
let profile = ArenaProfile::EMBEDDED;
let budget = 256 * 1024;
let arena = Arena::new(pool(profile, budget, 2), budget, profile);
let big = profile.max_chunk_size * 3 + 5;
arena.alloc(big, 8).expect("allocator");
let sizes = arena.chunk_sizes();
assert_eq!(sizes.len(), 1);
assert_eq!(sizes[0], big.next_multiple_of(CHUNK_ALIGN));
let pool = Arc::clone(&arena.pool);
drop(arena);
assert_eq!(pool.parked_bytes(), 0);
}
#[test]
fn pool_bound_is_respected() {
let profile = ArenaProfile::EMBEDDED;
let pool = pool(profile, 4096, 1);
assert_eq!(pool.budget(), 4096);
for _ in 0..64 {
pool.give(ChunkHandle::alloc(4096).expect("allocator"));
assert!(pool.parked_bytes() <= pool.budget());
}
assert_eq!(pool.parked_bytes(), 4096);
assert!(pool.take(4096).is_some());
assert_eq!(pool.parked_bytes(), 0);
assert!(pool.take(4096).is_none());
}
#[test]
fn pool_refuses_sizes_that_are_not_classes() {
let profile = ArenaProfile::EMBEDDED;
let pool = pool(profile, 256 * 1024, 2);
pool.give(ChunkHandle::alloc(3000).expect("allocator"));
pool.give(ChunkHandle::alloc(1024 * 1024).expect("allocator"));
assert_eq!(pool.parked_bytes(), 0, "neither size is a class");
pool.give(ChunkHandle::alloc(8192).expect("allocator"));
assert_eq!(pool.parked_bytes(), 8192);
}
#[test]
fn steady_state_recycling_touches_the_allocator_zero_times() {
let profile = ArenaProfile::EMBEDDED;
let budget = 64 * 1024;
let pool = pool(profile, budget, 2);
let fill = |arena: &Arena| {
while arena.used_bytes() < budget {
arena.alloc(256, 8).expect("allocator");
}
};
for _ in 0..2 {
let arena = Arena::new(Arc::clone(&pool), budget, profile);
fill(&arena);
}
assert!(pool.parked_bytes() > 0, "pool must hold the retired chunks");
for cycle in 3..8 {
let before = pool.misses();
let arena = Arena::new(Arc::clone(&pool), budget, profile);
fill(&arena);
let chunks = arena.chunk_count();
drop(arena);
let missed = pool.misses() - before;
assert_eq!(
missed, 0,
"cycle {cycle}: {missed} of {chunks} chunks came from the global allocator"
);
}
}
#[test]
fn high_water_mark_is_bounded_across_many_cycles() {
let profile = ArenaProfile::EMBEDDED;
let w = 64 * 1024usize;
let m = 2usize;
let c = profile.max_chunk_size;
let pool = pool(profile, w, m);
let bound = 2 * m * (w + c) + m * w;
let mut peak = 0usize;
let mut live: Vec<Arena> = Vec::new();
for _ in 0..RECYCLE_CYCLES {
let arena = Arena::new(Arc::clone(&pool), w, profile);
while arena.used_bytes() < w {
arena.alloc(128, 8).expect("allocator");
}
live.push(arena);
if live.len() > 2 * m {
live.remove(0);
}
let total: usize =
live.iter().map(|a| a.reserved_bytes()).sum::<usize>() + pool.parked_bytes();
peak = peak.max(total);
}
assert!(peak <= bound, "peak {peak} exceeds bound {bound}");
}
proptest! {
#[test]
fn allocations_never_overlap_or_leave_their_chunk(
sizes in proptest::collection::vec(1usize..3000, 1..200),
budget_kib in 1usize..64,
) {
let profile = ArenaProfile::EMBEDDED;
let budget = budget_kib * 1024;
let arena = Arena::new(pool(profile, budget, 2), budget, profile);
let mut spans: Vec<(usize, usize)> = Vec::new();
for size in sizes {
let ptr = arena.alloc(size, 8).expect("allocator");
let start = ptr.as_ptr() as usize;
prop_assert_eq!(start % 8, 0);
for (other, len) in &spans {
prop_assert!(start + size <= *other || other + len <= start);
}
spans.push((start, size));
}
prop_assert!(arena.used_bytes() >= spans.iter().map(|(_, l)| l).sum::<usize>());
prop_assert!(arena.reserved_bytes() >= arena.used_bytes());
}
#[test]
fn base_class_is_always_a_valid_pool_class(budget in 0usize..(1 << 28)) {
for profile in [ArenaProfile::SERVER, ArenaProfile::EMBEDDED] {
let base = base_class(profile, budget);
prop_assert!(base.is_power_of_two());
prop_assert!(base <= profile.max_chunk_size);
let pool = ChunkPool::new(profile, budget, 2);
prop_assert_eq!(pool.class_index(base), Some(0));
}
}
}
}