Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
98 changes: 80 additions & 18 deletions src/accelerate.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,43 +21,105 @@ pub fn id() -> usize {
}

/// Evict the accelerators.
pub fn evict() {
pub fn evict(release: bool) {
let mut accelerators = ACCELERATORS.write();
let (offset, vec) = &mut *accelerators;

// Update the offset.
*offset = ID.load(Ordering::SeqCst);
evict_pool(vec, release);
}

/// Retains used buffers for reuse, but releases idle or explicitly cleared storage.
fn evict_pool(vec: &mut Vec<Accelerator>, release: bool) {
if release {
*vec = Vec::new();
return;
}

// Clear all accelerators while keeping the memory allocated.
vec.iter_mut().for_each(|accelerator| accelerator.lock().clear())
for accelerator in vec {
let map = accelerator.get_mut();
if map.is_empty() {
*map = FxHashMap::default();
} else {
map.clear();
}
}
}

/// Get an accelerator by ID.
pub fn get(id: usize) -> Option<MappedRwLockReadGuard<'static, Accelerator>> {
// We always lock the accelerators, as we need to make sure that the
// accelerator is not removed while we are reading it.
let mut accelerators = ACCELERATORS.read();

let mut i = id.checked_sub(accelerators.0)?;
if i >= accelerators.1.len() {
loop {
let accelerators = ACCELERATORS.read();
let i = id.checked_sub(accelerators.0)?;
if i < accelerators.1.len() {
return Some(RwLockReadGuard::map(accelerators, move |(_, vec)| &vec[i]));
}
drop(accelerators);
resize(i + 1);
accelerators = ACCELERATORS.read();

// Because we release the lock before resizing the accelerator, we need
// to check again whether the ID is still valid because another thread
// might evicted the cache.
i = id.checked_sub(accelerators.0)?;
resize(id);
// Eviction can change both the offset and the pool length while the
// lock is released. Recheck both before indexing the pool.
}

Some(RwLockReadGuard::map(accelerators, move |(_, vec)| &vec[i]))
}

/// Adjusts the amount of accelerators.
#[cold]
fn resize(len: usize) {
fn resize(id: usize) {
let mut pair = ACCELERATORS.write();
let Some(i) = id.checked_sub(pair.0) else { return };
let len = i + 1;
if len > pair.1.len() {
pair.1.resize_with(len, || Mutex::new(FxHashMap::default()));
}
}

#[cfg(test)]
mod tests {
use super::*;

fn populated() -> Accelerator {
Mutex::new((0..128).map(|key| (key, key)).collect())
}

#[test]
fn eviction_releases_idle_and_explicitly_cleared_storage() {
let mut pool = (0..128).map(|_| populated()).collect::<Vec<_>>();
evict_pool(&mut pool, false);
assert!(pool[0].get_mut().is_empty());
assert!(pool[0].get_mut().capacity() >= 128);

// An unused epoch releases the empty tables, preserving pool slots.
evict_pool(&mut pool, false);
assert!(pool.iter_mut().all(|map| map.get_mut().capacity() == 0));
pool[0].get_mut().insert(1, 2);
evict_pool(&mut pool, true);
assert!(pool.is_empty());
assert_eq!(pool.capacity(), 0);
}

#[test]
fn concurrent_access_and_eviction() {
std::thread::scope(|scope| {
for _ in 0..4 {
scope.spawn(|| {
for _ in 0..2000 {
let id = id();
if let Some(accelerator) = get(id) {
let mut map = accelerator.lock();
map.insert(id as u128, 42);
assert_eq!(map.get(&(id as u128)), Some(&42));
}
}
});
}
for _ in 0..2000 {
evict(true);
std::thread::yield_now();
}
});
let expired = id();
evict(true);
assert!(get(expired).is_none());
}
}
4 changes: 2 additions & 2 deletions src/memoize.rs
Original file line number Diff line number Diff line change
Expand Up @@ -96,13 +96,13 @@ where
/// This removes all memoized results from the cache whose age is larger than or
/// equal to `max_age`. The age of a result grows by one during each eviction
/// and is reset to zero when the result produces a cache hit. Set `max_age` to
/// zero to completely clear the cache.
/// zero to completely clear the cache and release its backing storage.
pub fn evict(max_age: usize) {
for subevict in EVICTORS.read().iter() {
subevict(max_age);
}

accelerate::evict();
accelerate::evict(max_age == 0);
}

/// Register an eviction function in the global list.
Expand Down
22 changes: 22 additions & 0 deletions src/tree.rs
Original file line number Diff line number Diff line change
Expand Up @@ -191,6 +191,11 @@ impl<C, T> CallTree<C, T> {
// Prune edges.
self.edges.retain(|_, node| exists(*node));
self.start.retain(|_, node| exists(*node));

// Empty slabs and maps otherwise retain their high-water allocations.
if self.leaves.is_empty() {
*self = Self::new();
}
}

/// Checks a few invariants of the data structure.
Expand Down Expand Up @@ -286,6 +291,23 @@ mod tests {

use super::*;

#[test]
fn test_retain_releases_empty_storage() {
let mut tree = CallTree::new();
for key in 0..1024 {
tree.insert(key, [('a', key), ('b', key + 1)].into_iter().collect(), key)
.unwrap();
}
tree.retain(|_| false);
tree.assert_consistency();
assert_eq!(tree.inner.capacity(), 0);
assert_eq!(tree.leaves.capacity(), 0);
assert_eq!(tree.start.capacity(), 0);
assert_eq!(tree.edges.capacity(), 0);
tree.insert(0, [('a', 1)].into_iter().collect(), 42).unwrap();
assert_eq!(tree.get(0, |_| 1), Some(&42));
}

#[test]
fn test_call_tree() {
test_ops([
Expand Down
Loading