Skip to main content

feoxdb/core/
cache.rs

1use crate::constants::*;
2use crate::core::record::Record;
3use crate::stats::Statistics;
4use crate::utils::hash::murmur3_32;
5use bytes::Bytes;
6use parking_lot::{Mutex, RwLock, RwLockWriteGuard};
7use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
8use std::sync::{Arc, Weak};
9
10/// CLOCK algorithm cache implementation
11/// Uses reference bits and circular scanning for eviction
12pub struct ClockCache {
13    /// Cache entries organized in buckets for better locality
14    buckets: Vec<RwLock<Vec<CacheEntry>>>,
15
16    /// Global CLOCK hand position for eviction scanning
17    clock_hand: AtomicUsize,
18
19    /// High watermark for triggering eviction (bytes)
20    high_watermark: AtomicUsize,
21
22    /// Low watermark to evict down to (bytes)
23    low_watermark: AtomicUsize,
24
25    /// Lock for eviction process
26    eviction_lock: Mutex<()>,
27
28    /// Shared statistics
29    stats: Arc<Statistics>,
30}
31
32struct CacheEntry {
33    key: Vec<u8>,
34    value: Bytes,
35    record: Option<Weak<Record>>,
36
37    /// Reference bit for CLOCK algorithm (accessed recently)
38    reference_bit: AtomicBool,
39
40    /// Size of this entry in bytes
41    size: usize,
42}
43
44pub(crate) struct RecordCacheEntry<'a> {
45    bucket: RwLockWriteGuard<'a, Vec<CacheEntry>>,
46    position: Option<usize>,
47    stats: &'a Statistics,
48}
49
50impl RecordCacheEntry<'_> {
51    pub(crate) fn value(&self) -> Option<Bytes> {
52        self.position
53            .map(|position| self.bucket[position].value.clone())
54    }
55
56    pub(crate) fn remove(mut self) {
57        if let Some(position) = self.position {
58            let entry = self.bucket.remove(position);
59            self.stats
60                .cache_memory
61                .fetch_sub(entry.size, Ordering::Relaxed);
62        }
63    }
64}
65
66impl ClockCache {
67    pub fn new(stats: Arc<Statistics>) -> Self {
68        let buckets = (0..CACHE_BUCKETS)
69            .map(|_| RwLock::new(Vec::new()))
70            .collect();
71
72        Self {
73            buckets,
74            clock_hand: AtomicUsize::new(0),
75            high_watermark: AtomicUsize::new(CACHE_HIGH_WATERMARK_MB * MB),
76            low_watermark: AtomicUsize::new(CACHE_LOW_WATERMARK_MB * MB),
77            eviction_lock: Mutex::new(()),
78            stats,
79        }
80    }
81
82    /// Get value from cache, setting reference bit on access
83    pub fn get(&self, key: &[u8]) -> Option<Bytes> {
84        self.get_entry(key, None)
85    }
86
87    pub(crate) fn get_for_record(&self, key: &[u8], record: &Arc<Record>) -> Option<Bytes> {
88        self.get_entry(key, Some(record))
89    }
90
91    fn get_entry(&self, key: &[u8], record: Option<&Arc<Record>>) -> Option<Bytes> {
92        let hash = murmur3_32(key, 0);
93        let bucket_idx = (hash as usize) % CACHE_BUCKETS;
94
95        let bucket = self.buckets[bucket_idx].read();
96
97        for entry in bucket.iter() {
98            if entry.key != key {
99                continue;
100            }
101            let generation_matches = match (record, entry.record.as_ref()) {
102                (Some(expected), Some(cached)) => {
103                    std::ptr::eq(cached.as_ptr(), Arc::as_ptr(expected))
104                }
105                (Some(_), None) => false,
106                (None, _) => true,
107            };
108            if generation_matches {
109                // Set reference bit on access (CLOCK algorithm)
110                entry.reference_bit.store(true, Ordering::Relaxed);
111                return Some(entry.value.clone());
112            }
113        }
114
115        None
116    }
117
118    /// Insert value into cache, triggering eviction if needed
119    pub fn insert(&self, key: Vec<u8>, value: Bytes) {
120        self.insert_entry(key, value, None);
121    }
122
123    pub(crate) fn insert_for_record(&self, key: Vec<u8>, value: Bytes, record: &Arc<Record>) {
124        self.insert_entry(key, value, Some(record));
125    }
126
127    fn insert_entry(&self, key: Vec<u8>, value: Bytes, record: Option<&Arc<Record>>) {
128        let size = key.len() + value.len() + std::mem::size_of::<CacheEntry>();
129
130        // Don't cache very large values
131        let high_watermark = self.high_watermark.load(Ordering::Relaxed);
132        if size > high_watermark / 4 {
133            return;
134        }
135
136        // Check if we need to evict before inserting
137        let current_usage = self.stats.cache_memory.load(Ordering::Relaxed);
138        if current_usage + size > high_watermark {
139            self.evict_entries();
140        }
141
142        let hash = murmur3_32(&key, 0);
143        let bucket_idx = (hash as usize) % CACHE_BUCKETS;
144
145        let mut bucket = self.buckets[bucket_idx].write();
146
147        // Check if key already exists and update
148        for entry in bucket.iter_mut() {
149            if entry.key == key {
150                if !can_replace_generation(entry.record.as_ref(), record) {
151                    return;
152                }
153                let old_size = entry.size;
154                entry.value = value;
155                entry.record = record.map(Arc::downgrade);
156                entry.size = size;
157                entry.reference_bit.store(true, Ordering::Relaxed);
158
159                // Update memory usage
160                if size > old_size {
161                    self.stats
162                        .cache_memory
163                        .fetch_add(size - old_size, Ordering::Relaxed);
164                } else {
165                    self.stats
166                        .cache_memory
167                        .fetch_sub(old_size - size, Ordering::Relaxed);
168                }
169                return;
170            }
171        }
172
173        // Add new entry
174        let entry = CacheEntry {
175            key,
176            value,
177            record: record.map(Arc::downgrade),
178            reference_bit: AtomicBool::new(true),
179            size,
180        };
181
182        bucket.push(entry);
183        self.stats.cache_memory.fetch_add(size, Ordering::Relaxed);
184    }
185
186    /// Remove specific key from cache
187    pub fn remove(&self, key: &[u8]) {
188        self.remove_entry(key, None);
189    }
190
191    pub(crate) fn remove_for_record(&self, key: &[u8], record: &Arc<Record>) {
192        self.remove_entry(key, Some(record));
193    }
194
195    pub(crate) fn record_entry<'a>(
196        &'a self,
197        key: &[u8],
198        record: &Arc<Record>,
199    ) -> RecordCacheEntry<'a> {
200        let hash = murmur3_32(key, 0);
201        let bucket_idx = (hash as usize) % CACHE_BUCKETS;
202        let bucket = self.buckets[bucket_idx].write();
203        let position = bucket.iter().position(|entry| {
204            entry.key == key
205                && entry
206                    .record
207                    .as_ref()
208                    .is_some_and(|cached| std::ptr::eq(cached.as_ptr(), Arc::as_ptr(record)))
209        });
210
211        RecordCacheEntry {
212            bucket,
213            position,
214            stats: self.stats.as_ref(),
215        }
216    }
217
218    fn remove_entry(&self, key: &[u8], record: Option<&Arc<Record>>) {
219        let hash = murmur3_32(key, 0);
220        let bucket_idx = (hash as usize) % CACHE_BUCKETS;
221
222        let mut bucket = self.buckets[bucket_idx].write();
223
224        if let Some(pos) = bucket.iter().position(|entry| {
225            entry.key == key
226                && record.is_none_or(|expected| {
227                    entry
228                        .record
229                        .as_ref()
230                        .is_some_and(|cached| std::ptr::eq(cached.as_ptr(), Arc::as_ptr(expected)))
231                })
232        }) {
233            let entry = bucket.remove(pos);
234            self.stats
235                .cache_memory
236                .fetch_sub(entry.size, Ordering::Relaxed);
237        }
238    }
239
240    /// CLOCK algorithm eviction - scan entries circularly, evicting those without reference bit
241    pub fn evict_entries(&self) {
242        // Try to acquire eviction lock, return if already evicting
243        let _lock = match self.eviction_lock.try_lock() {
244            Some(lock) => lock,
245            None => return,
246        };
247
248        let target_usage = self.low_watermark.load(Ordering::Relaxed);
249        let mut current_usage = self.stats.cache_memory.load(Ordering::Relaxed);
250
251        if current_usage <= target_usage {
252            return;
253        }
254
255        let mut scans = 0;
256        const MAX_SCANS: usize = 3; // Maximum passes through cache
257        let mut hand = self.clock_hand.load(Ordering::Relaxed);
258
259        while current_usage > target_usage && scans < MAX_SCANS {
260            for _ in 0..CACHE_BUCKETS {
261                let bucket_index = hand % CACHE_BUCKETS;
262                hand = hand.wrapping_add(1);
263                let mut bucket = self.buckets[bucket_index].write();
264                let mut i = 0;
265
266                while i < bucket.len() {
267                    let entry = &bucket[i];
268
269                    // Check reference bit
270                    if entry.reference_bit.load(Ordering::Relaxed) {
271                        entry.reference_bit.store(false, Ordering::Relaxed);
272                        i += 1;
273                    } else {
274                        // No reference bit - evict this entry
275                        let removed = bucket.remove(i);
276                        let previous_usage = self
277                            .stats
278                            .cache_memory
279                            .fetch_sub(removed.size, Ordering::Relaxed);
280                        self.stats.record_eviction(1);
281                        current_usage = previous_usage - removed.size;
282                        // Don't increment i since we removed an element
283                    }
284
285                    if current_usage <= target_usage {
286                        break;
287                    }
288                }
289
290                if current_usage <= target_usage {
291                    break;
292                }
293            }
294
295            scans += 1;
296        }
297        self.clock_hand.store(hand, Ordering::Relaxed);
298    }
299
300    /// Clear all cache entries
301    pub fn clear(&self) {
302        let _lock = self.eviction_lock.lock();
303        for bucket in &self.buckets {
304            let removed_size = {
305                let mut entries = bucket.write();
306                let removed_size = entries.iter().map(|entry| entry.size).sum();
307                entries.clear();
308                removed_size
309            };
310            self.stats
311                .cache_memory
312                .fetch_sub(removed_size, Ordering::Relaxed);
313            #[cfg(test)]
314            crate::test_hooks::pause_at(crate::test_hooks::AFTER_CACHE_BUCKET_CLEAR);
315        }
316
317        self.clock_hand.store(0, Ordering::Relaxed);
318    }
319
320    /// Get current cache statistics
321    pub fn stats(&self) -> CacheStats {
322        CacheStats {
323            entries: 0, // Calculate from buckets if needed
324            memory_usage: self.stats.cache_memory.load(Ordering::Relaxed),
325            high_watermark: self.high_watermark.load(Ordering::Relaxed),
326            low_watermark: self.low_watermark.load(Ordering::Relaxed),
327        }
328    }
329
330    /// Adjust cache watermarks dynamically
331    pub fn adjust_watermarks(&self, high_mb: usize, low_mb: usize) {
332        let high = high_mb * MB;
333        let low = low_mb * MB;
334
335        if high > low && high <= CACHE_MAX_SIZE {
336            // Max 1GB for cache
337            // Update watermarks atomically
338            self.high_watermark.store(high, Ordering::Relaxed);
339            self.low_watermark.store(low, Ordering::Relaxed);
340
341            // Trigger eviction if we're over the new high watermark
342            let current_usage = self.stats.cache_memory.load(Ordering::Relaxed);
343            if current_usage > high {
344                self.evict_entries();
345            }
346        }
347    }
348}
349
350fn can_replace_generation(cached: Option<&Weak<Record>>, incoming: Option<&Arc<Record>>) -> bool {
351    let Some(incoming) = incoming else {
352        return true;
353    };
354    if incoming.refcount.load(Ordering::Acquire) == 0 {
355        return false;
356    }
357    let Some(cached) = cached else {
358        return true;
359    };
360    if std::ptr::eq(cached.as_ptr(), Arc::as_ptr(incoming)) {
361        return true;
362    }
363
364    cached.upgrade().is_none_or(|cached| {
365        cached.refcount.load(Ordering::Acquire) == 0 || cached.timestamp < incoming.timestamp
366    })
367}
368
369#[derive(Debug, Clone)]
370pub struct CacheStats {
371    pub entries: u32,
372    pub memory_usage: usize,
373    pub high_watermark: usize,
374    pub low_watermark: usize,
375}