Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
17 commits
Select commit Hold shift + click to select a range
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
30 changes: 24 additions & 6 deletions src/core/src/cache/budget.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,8 +6,8 @@ use crate::sync::{

#[derive(Debug)]
pub struct BudgetAccounting {
max_memory_bytes: usize,
max_disk_bytes: usize,
max_memory_bytes: AtomicUsize,
max_disk_bytes: AtomicUsize,
used_memory_bytes: AtomicUsize,
used_disk_bytes: AtomicUsize,
observer: Arc<Observer>,
Expand All @@ -20,8 +20,8 @@ impl BudgetAccounting {
observer: Arc<Observer>,
) -> Self {
Self {
max_memory_bytes,
max_disk_bytes,
max_memory_bytes: AtomicUsize::new(max_memory_bytes),
max_disk_bytes: AtomicUsize::new(max_disk_bytes),
used_memory_bytes: AtomicUsize::new(0),
used_disk_bytes: AtomicUsize::new(0),
observer,
Expand All @@ -33,11 +33,29 @@ impl BudgetAccounting {
self.used_disk_bytes.store(0, Ordering::Relaxed);
}

/// Dynamically update the max memory limit. Takes effect for new reservations.
pub fn set_max_memory_bytes(&self, new_limit: usize) {
self.max_memory_bytes.store(new_limit, Ordering::Relaxed);
}

/// Dynamically update the max disk limit. Takes effect for new reservations.
pub fn set_max_disk_bytes(&self, new_limit: usize) {
self.max_disk_bytes.store(new_limit, Ordering::Relaxed);
}

pub fn max_memory_bytes(&self) -> usize {
self.max_memory_bytes.load(Ordering::Relaxed)
}

pub fn max_disk_bytes(&self) -> usize {
self.max_disk_bytes.load(Ordering::Relaxed)
}

/// Try to reserve memory in the cache.
/// Returns ok if the memory was reserved, err if the memory budget is full.
pub(super) fn try_reserve_memory(&self, request_bytes: usize) -> Result<(), ()> {
let used = self.used_memory_bytes.load(Ordering::Relaxed);
if used + request_bytes > self.max_memory_bytes {
if used + request_bytes > self.max_memory_bytes.load(Ordering::Relaxed) {
return Err(());
}

Expand Down Expand Up @@ -80,7 +98,7 @@ impl BudgetAccounting {

pub(super) fn try_reserve_disk(&self, request_bytes: usize) -> Result<(), ()> {
let used = self.used_disk_bytes.load(Ordering::Relaxed);
if used + request_bytes > self.max_disk_bytes {
if used + request_bytes > self.max_disk_bytes.load(Ordering::Relaxed) {
self.observer.on_disk_reservation_failure();
return Err(());
}
Expand Down
4 changes: 2 additions & 2 deletions src/core/src/cache/core.rs
Original file line number Diff line number Diff line change
Expand Up @@ -112,8 +112,8 @@ impl LiquidCache {
memory_squeezed_liquid_bytes,
memory_usage_bytes,
disk_usage_bytes,
max_memory_bytes: self.config.max_memory_bytes(),
max_disk_bytes: self.config.max_disk_bytes(),
max_memory_bytes: self.budget.max_memory_bytes(),
max_disk_bytes: self.budget.max_disk_bytes(),
runtime,
}
}
Expand Down
2 changes: 1 addition & 1 deletion src/core/src/cache/observer/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -131,7 +131,7 @@ impl Observer {
}
}

pub(crate) fn runtime_stats(&self) -> &RuntimeStats {
pub fn runtime_stats(&self) -> &RuntimeStats {
&self.runtime
}
}
2 changes: 2 additions & 0 deletions src/core/src/cache/observer/stats.rs
Original file line number Diff line number Diff line change
Expand Up @@ -97,6 +97,8 @@ define_runtime_stats! {
(get, "Number of `get` calls issued via `CachedData`.", incr_get),
(get_with_selection, "Number of `get_with_selection` calls issued via `CachedData`.", incr_get_with_selection),
(eval_predicate, "Number of `eval_predicate` calls issued via `CachedData`.", incr_eval_predicate),
(cache_hit, "Number of cache hits (data found in cache).", incr_cache_hit),
(cache_miss, "Number of cache misses (data not in cache, fell back to Parquet).", incr_cache_miss),
(get_squeezed_success, "Number of Squeezed-Liquid full evaluations finished without IO.", incr_get_squeezed_success),
(get_squeezed_needs_io, "Number of Squeezed-Liquid full paths that required IO.", incr_get_squeezed_needs_io),
(try_read_liquid_calls, "Number of `try_read_liquid` calls issued via `CachedData`.", incr_try_read_liquid),
Expand Down
291 changes: 291 additions & 0 deletions src/core/src/cache/policies/cache/lru.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,291 @@
//! LRU (Least Recently Used) cache eviction policy.
//!
//! Uses 4 queues by entry type (Arrow, Liquid, Squeezed, Disk), same as LiquidPolicy.
//! Within each queue, entries are ordered by recency (moved to back on access).
//! Eviction priority: Arrow (largest) first, then Liquid, then Squeezed.
//! Within each queue, evicts the LRU entry (front of the queue).

use crate::cache::cached_batch::CachedBatchType;
use crate::cache::utils::EntryID;
use crate::sync::Mutex;
use ahash::AHashMap;
use std::ptr::NonNull;

use super::CachePolicy;
use super::doubly_linked_list::{DoublyLinkedList, DoublyLinkedNode, drop_boxed_node};

/// Which queue an entry belongs to.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum LruQueueKind {
Arrow,
Liquid,
Squeezed,
Disk,
}

/// LRU cache eviction policy with type-aware queues.
///
/// On insert: entry goes to the back of its type queue (most recently used).
/// On access: entry moves to the back of its type queue.
/// On eviction: picks the LRU entry from Arrow queue first, then Liquid, then Squeezed.
#[derive(Debug)]
pub struct LruPolicy {
inner: Mutex<LruInner>,
}

#[derive(Debug)]
struct LruInner {
arrow: DoublyLinkedList<EntryID>,
liquid: DoublyLinkedList<EntryID>,
squeezed: DoublyLinkedList<EntryID>,
disk: DoublyLinkedList<EntryID>,
/// Maps entry_id → (node_ptr, which queue it's in)
map: AHashMap<EntryID, (NonNull<DoublyLinkedNode<EntryID>>, LruQueueKind)>,
}

// Safety: We control access via Mutex and never expose raw pointers outside.
unsafe impl Send for LruInner {}
unsafe impl Sync for LruInner {}

impl LruInner {
fn queue_mut(&mut self, kind: LruQueueKind) -> &mut DoublyLinkedList<EntryID> {
match kind {
LruQueueKind::Arrow => &mut self.arrow,
LruQueueKind::Liquid => &mut self.liquid,
LruQueueKind::Squeezed => &mut self.squeezed,
LruQueueKind::Disk => &mut self.disk,
}
}

fn queue_kind_for(batch_type: CachedBatchType) -> LruQueueKind {
match batch_type {
CachedBatchType::MemoryArrow => LruQueueKind::Arrow,
CachedBatchType::MemoryLiquid => LruQueueKind::Liquid,
CachedBatchType::MemorySqueezedLiquid => LruQueueKind::Squeezed,
CachedBatchType::DiskLiquid | CachedBatchType::DiskArrow => LruQueueKind::Disk,
}
}

/// Pop the LRU entry (front) from a specific queue.
fn pop_lru(&mut self, kind: LruQueueKind) -> Option<EntryID> {
let list = self.queue_mut(kind);
let node_ptr = list.head()?;
let entry_id = unsafe { node_ptr.as_ref().data };

Check warning on line 73 in src/core/src/cache/policies/cache/lru.rs

View check run for this annotation

Codacy Production / Codacy Static Code Analysis

src/core/src/cache/policies/cache/lru.rs#L73

Detected 'unsafe' usage, please audit for secure usage
unsafe { list.unlink(node_ptr) };

Check warning on line 74 in src/core/src/cache/policies/cache/lru.rs

View check run for this annotation

Codacy Production / Codacy Static Code Analysis

src/core/src/cache/policies/cache/lru.rs#L74

Detected 'unsafe' usage, please audit for secure usage
self.map.remove(&entry_id);
unsafe { drop_boxed_node(node_ptr) };

Check warning on line 76 in src/core/src/cache/policies/cache/lru.rs

View check run for this annotation

Codacy Production / Codacy Static Code Analysis

src/core/src/cache/policies/cache/lru.rs#L76

Detected 'unsafe' usage, please audit for secure usage
Some(entry_id)
}
}

impl LruPolicy {
/// Create a new LRU policy.
pub fn new() -> Self {
Self {
inner: Mutex::new(LruInner {
arrow: DoublyLinkedList::new(),
liquid: DoublyLinkedList::new(),
squeezed: DoublyLinkedList::new(),
disk: DoublyLinkedList::new(),
map: AHashMap::new(),
}),
}
}
}

impl Default for LruPolicy {
fn default() -> Self {
Self::new()
}
}

impl CachePolicy for LruPolicy {
fn find_memory_victim(&self, cnt: usize) -> Vec<EntryID> {
let mut inner = self.inner.lock().unwrap();
let mut victims = Vec::with_capacity(cnt);

while victims.len() < cnt {
// Priority: evict Arrow (largest) first, then Liquid, then Squeezed
if let Some(entry) = inner.pop_lru(LruQueueKind::Arrow) {
victims.push(entry);
continue;
}
if let Some(entry) = inner.pop_lru(LruQueueKind::Liquid) {
victims.push(entry);
continue;
}
if let Some(entry) = inner.pop_lru(LruQueueKind::Squeezed) {
victims.push(entry);
continue;
}
break;
}

victims
}

fn find_disk_victim(&self, cnt: usize) -> Vec<EntryID> {
if cnt == 0 {
return vec![];
}
let mut inner = self.inner.lock().unwrap();
let mut victims = Vec::with_capacity(cnt);
while victims.len() < cnt {
if let Some(entry) = inner.pop_lru(LruQueueKind::Disk) {
victims.push(entry);
} else {
break;
}
}
victims
}

fn notify_insert(&self, entry_id: &EntryID, batch_type: CachedBatchType) {
let mut inner = self.inner.lock().unwrap();
let target = LruInner::queue_kind_for(batch_type);

// If already present, remove from old position/queue
if let Some((node_ptr, old_kind)) = inner.map.remove(entry_id) {
let old_list = inner.queue_mut(old_kind);
unsafe { old_list.unlink(node_ptr) };

Check warning on line 150 in src/core/src/cache/policies/cache/lru.rs

View check run for this annotation

Codacy Production / Codacy Static Code Analysis

src/core/src/cache/policies/cache/lru.rs#L150

Detected 'unsafe' usage, please audit for secure usage
unsafe { drop_boxed_node(node_ptr) };
}

// Insert at the back of the target queue (most recently used)
let node = DoublyLinkedNode::new(*entry_id);
let node_ptr = NonNull::new(Box::into_raw(node)).unwrap();
let list = inner.queue_mut(target);
unsafe { list.push_back(node_ptr) };

Check warning on line 158 in src/core/src/cache/policies/cache/lru.rs

View check run for this annotation

Codacy Production / Codacy Static Code Analysis

src/core/src/cache/policies/cache/lru.rs#L158

Detected 'unsafe' usage, please audit for secure usage
inner.map.insert(*entry_id, (node_ptr, target));
}

fn notify_access(&self, entry_id: &EntryID, _batch_type: CachedBatchType) {
let mut inner = self.inner.lock().unwrap();

let Some(&(node_ptr, kind)) = inner.map.get(entry_id) else {
return;
};

// Move to back of its current queue (most recently used)
let list = inner.queue_mut(kind);
unsafe {

Check warning on line 171 in src/core/src/cache/policies/cache/lru.rs

View check run for this annotation

Codacy Production / Codacy Static Code Analysis

src/core/src/cache/policies/cache/lru.rs#L171

Detected 'unsafe' usage, please audit for secure usage
list.unlink(node_ptr);
list.push_back(node_ptr);
}
}

fn notify_remove(&self, entry_id: &EntryID) {
let mut inner = self.inner.lock().unwrap();

if let Some((node_ptr, kind)) = inner.map.remove(entry_id) {
let list = inner.queue_mut(kind);
unsafe { list.unlink(node_ptr) };
unsafe { drop_boxed_node(node_ptr) };
}
}
}

impl Drop for LruPolicy {
fn drop(&mut self) {
let mut inner = self.inner.lock().unwrap();
unsafe {
inner.arrow.drop_all();
inner.liquid.drop_all();
inner.squeezed.drop_all();
inner.disk.drop_all();
}
inner.map.clear();
}
}


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

fn entry(id: usize) -> EntryID {
EntryID::from(id)
}

#[test]
fn test_lru_evicts_arrow_first() {
let policy = LruPolicy::new();

// Insert entries: some Arrow, some Liquid
policy.notify_insert(&entry(1), CachedBatchType::MemoryArrow);
policy.notify_insert(&entry(2), CachedBatchType::MemoryLiquid);
policy.notify_insert(&entry(3), CachedBatchType::MemoryArrow);
policy.notify_insert(&entry(4), CachedBatchType::MemoryLiquid);

// Evict 1 — should pick from Arrow queue first (LRU = entry 1)
let victims = policy.find_memory_victim(1);
assert_eq!(victims, vec![entry(1)]);

// Evict 1 more — next Arrow LRU = entry 3
let victims = policy.find_memory_victim(1);
assert_eq!(victims, vec![entry(3)]);

// Evict 1 more — Arrow empty, now Liquid LRU = entry 2
let victims = policy.find_memory_victim(1);
assert_eq!(victims, vec![entry(2)]);
}

#[test]
fn test_lru_access_moves_to_back() {
let policy = LruPolicy::new();

policy.notify_insert(&entry(1), CachedBatchType::MemoryArrow);
policy.notify_insert(&entry(2), CachedBatchType::MemoryArrow);
policy.notify_insert(&entry(3), CachedBatchType::MemoryArrow);

// Access entry 1 — moves it to back
policy.notify_access(&entry(1), CachedBatchType::MemoryArrow);

// Evict 1 — LRU is now entry 2 (entry 1 was moved to back)
let victims = policy.find_memory_victim(1);
assert_eq!(victims, vec![entry(2)]);
}

#[test]
fn test_lru_insert_moves_between_queues() {
let policy = LruPolicy::new();

// Insert as Arrow
policy.notify_insert(&entry(1), CachedBatchType::MemoryArrow);

// Re-insert same entry as Liquid (simulates transcode)
policy.notify_insert(&entry(1), CachedBatchType::MemoryLiquid);

// Arrow queue should be empty now
let victims = policy.find_memory_victim(1);
// Should come from Liquid queue
assert_eq!(victims, vec![entry(1)]);
}

#[test]
fn test_lru_disk_victims() {
let policy = LruPolicy::new();

policy.notify_insert(&entry(1), CachedBatchType::DiskLiquid);
policy.notify_insert(&entry(2), CachedBatchType::DiskLiquid);

// find_disk_victim should return LRU disk entries
let victims = policy.find_disk_victim(1);
assert_eq!(victims, vec![entry(1)]);
}

#[test]
fn test_lru_remove() {
let policy = LruPolicy::new();

policy.notify_insert(&entry(1), CachedBatchType::MemoryArrow);
policy.notify_insert(&entry(2), CachedBatchType::MemoryArrow);

// Remove entry 1
policy.notify_remove(&entry(1));

// Evict — should get entry 2 (entry 1 was removed)
let victims = policy.find_memory_victim(1);
assert_eq!(victims, vec![entry(2)]);
}
}
Loading