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
10pub struct ClockCache {
13 buckets: Vec<RwLock<Vec<CacheEntry>>>,
15
16 clock_hand: AtomicUsize,
18
19 high_watermark: AtomicUsize,
21
22 low_watermark: AtomicUsize,
24
25 eviction_lock: Mutex<()>,
27
28 stats: Arc<Statistics>,
30}
31
32struct CacheEntry {
33 key: Vec<u8>,
34 value: Bytes,
35 record: Option<Weak<Record>>,
36
37 reference_bit: AtomicBool,
39
40 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 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 entry.reference_bit.store(true, Ordering::Relaxed);
111 return Some(entry.value.clone());
112 }
113 }
114
115 None
116 }
117
118 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 let high_watermark = self.high_watermark.load(Ordering::Relaxed);
132 if size > high_watermark / 4 {
133 return;
134 }
135
136 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 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 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 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 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 pub fn evict_entries(&self) {
242 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; 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 if entry.reference_bit.load(Ordering::Relaxed) {
271 entry.reference_bit.store(false, Ordering::Relaxed);
272 i += 1;
273 } else {
274 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 }
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 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 pub fn stats(&self) -> CacheStats {
322 CacheStats {
323 entries: 0, 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 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 self.high_watermark.store(high, Ordering::Relaxed);
339 self.low_watermark.store(low, Ordering::Relaxed);
340
341 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}