1use super::{CacheEntryKey, CacheEntryKeyRef, EvictionManager};
18#[cfg(test)]
19use crate::key::CompactCacheKey;
20
21use async_trait::async_trait;
22use lru::LruCache;
23use parking_lot::RwLock;
24use pingora_error::{BError, ErrorType::*, OrErr, Result};
25use rand::Rng;
26use serde::de::SeqAccess;
27use serde::{Deserialize, Serialize};
28use std::collections::hash_map::DefaultHasher;
29use std::fs::File;
30use std::hash::{Hash, Hasher};
31use std::io::prelude::*;
32use std::path::Path;
33use std::sync::atomic::{AtomicUsize, Ordering};
34use std::time::SystemTime;
35
36#[derive(Debug, Deserialize, Serialize)]
37struct Node {
38 key: CacheEntryKey,
39 size: usize,
40}
41
42pub struct Manager {
46 lru: RwLock<LruCache<u64, Node>>,
47 limit: usize,
48 items_watermark: Option<usize>,
49 used: AtomicUsize,
50 items: AtomicUsize,
51 evicted_size: AtomicUsize,
52 evicted_items: AtomicUsize,
53}
54
55impl Manager {
56 pub fn new(limit: usize) -> Self {
58 Manager {
59 lru: RwLock::new(LruCache::unbounded()),
60 limit,
61 items_watermark: None,
62 used: AtomicUsize::new(0),
63 items: AtomicUsize::new(0),
64 evicted_size: AtomicUsize::new(0),
65 evicted_items: AtomicUsize::new(0),
66 }
67 }
68
69 pub fn new_with_watermark(limit: usize, items_watermark: Option<usize>) -> Self {
71 Manager {
72 lru: RwLock::new(LruCache::unbounded()),
73 limit,
74 items_watermark,
75 used: AtomicUsize::new(0),
76 items: AtomicUsize::new(0),
77 evicted_size: AtomicUsize::new(0),
78 evicted_items: AtomicUsize::new(0),
79 }
80 }
81
82 fn insert(&self, hash_key: u64, node: CacheEntryKey, size: usize, reverse: bool) {
83 use std::cmp::Ordering::*;
84 let node = Node { key: node, size };
85 let old = {
86 let mut lru = self.lru.write();
87 let old = lru.push(hash_key, node);
88 if reverse && old.is_none() {
89 lru.demote(&hash_key);
90 }
91 old
92 };
93 if let Some(old) = old {
94 match size.cmp(&old.1.size) {
96 Greater => self.used.fetch_add(size - old.1.size, Ordering::Relaxed),
97 Less => self.used.fetch_sub(old.1.size - size, Ordering::Relaxed),
98 Equal => 0, };
100 } else {
101 self.used.fetch_add(size, Ordering::Relaxed);
102 self.items.fetch_add(1, Ordering::Relaxed);
103 }
104 }
105
106 fn increase_weight(&self, key: u64, delta: usize, max_weight: Option<usize>) -> bool {
107 let mut lru = self.lru.write();
108 let Some(node) = lru.get_key_value_mut(&key) else {
109 return false;
110 };
111 let incremented = node.1.size.saturating_add(delta);
112 let new_size = max_weight.map_or(incremented, |m| incremented.min(m).max(node.1.size));
113 debug_assert!(new_size >= node.1.size);
114 self.used
115 .fetch_add(new_size - node.1.size, Ordering::Relaxed);
116 node.1.size = new_size;
117 true
118 }
119
120 #[inline]
121 fn over_limits(&self) -> bool {
122 self.used.load(Ordering::Relaxed) > self.limit
123 || self
124 .items_watermark
125 .is_some_and(|w| self.items.load(Ordering::Relaxed) > w)
126 }
127
128 fn evict(&self) -> Vec<CacheEntryKey> {
130 if self.used.load(Ordering::Relaxed) <= self.limit
131 && self
132 .items_watermark
133 .is_none_or(|w| self.items.load(Ordering::Relaxed) <= w)
134 {
135 return vec![];
136 }
137
138 let mut to_evict = Vec::with_capacity(1); while self.over_limits() {
141 if let Some((_, node)) = self.lru.write().pop_lru() {
142 self.used.fetch_sub(node.size, Ordering::Relaxed);
143 self.items.fetch_sub(1, Ordering::Relaxed);
144 self.evicted_size.fetch_add(node.size, Ordering::Relaxed);
145 self.evicted_items.fetch_add(1, Ordering::Relaxed);
146 to_evict.push(node.key);
147 } else {
148 return to_evict;
150 }
151 }
152 to_evict
153 }
154
155 fn serialize(&self) -> Result<Vec<u8>> {
158 use rmp_serde::encode::Serializer;
159 use serde::ser::SerializeSeq;
160 use serde::ser::Serializer as _;
161 let mut ser = Serializer::new(vec![]);
163 let lru = self.lru.read();
165 let mut seq = ser
166 .serialize_seq(Some(lru.len()))
167 .or_err(InternalError, "fail to serialize node")?;
168 for item in lru.iter() {
169 seq.serialize_element(item.1).unwrap(); }
171 seq.end().or_err(InternalError, "when serializing LRU")?;
172 Ok(ser.into_inner())
173 }
174
175 fn deserialize(&self, buf: &[u8]) -> Result<()> {
176 use rmp_serde::decode::Deserializer;
177 use serde::de::Deserializer as _;
178 let mut de = Deserializer::new(buf);
179 let visitor = InsertToManager { lru: self };
180 de.deserialize_seq(visitor)
181 .or_err(InternalError, "when deserializing LRU")?;
182 Ok(())
183 }
184}
185
186struct InsertToManager<'a> {
187 lru: &'a Manager,
188}
189
190impl<'de> serde::de::Visitor<'de> for InsertToManager<'_> {
191 type Value = ();
192
193 fn expecting(&self, formatter: &mut std::fmt::Formatter) -> std::fmt::Result {
194 formatter.write_str("array of lru nodes")
195 }
196
197 fn visit_seq<A>(self, mut seq: A) -> Result<Self::Value, A::Error>
198 where
199 A: SeqAccess<'de>,
200 {
201 while let Some(node) = seq.next_element::<Node>()? {
202 let key = u64key(&node.key);
203 self.lru.insert(key, node.key, node.size, true); }
205 Ok(())
206 }
207}
208
209#[inline]
210fn u64key(key: &impl Hash) -> u64 {
211 let mut hasher = DefaultHasher::new();
212 key.hash(&mut hasher);
213 hasher.finish()
214}
215
216const FILE_NAME: &str = "simple_lru.data";
217
218#[async_trait]
219impl EvictionManager for Manager {
220 fn total_size(&self) -> usize {
221 self.used.load(Ordering::Relaxed)
222 }
223 fn total_items(&self) -> usize {
224 self.items.load(Ordering::Relaxed)
225 }
226 fn evicted_size(&self) -> usize {
227 self.evicted_size.load(Ordering::Relaxed)
228 }
229 fn evicted_items(&self) -> usize {
230 self.evicted_items.load(Ordering::Relaxed)
231 }
232
233 fn admit(
234 &self,
235 item: CacheEntryKey,
236 size: usize,
237 _fresh_until: SystemTime,
238 ) -> Vec<CacheEntryKey> {
239 let key = u64key(&item);
240 self.insert(key, item, size, false);
241 self.evict()
242 }
243
244 fn increment_weight(
245 &self,
246 item: &CacheEntryKey,
247 delta: usize,
248 max_weight: Option<usize>,
249 ) -> Vec<CacheEntryKey> {
250 let key = u64key(item);
251 if !self.increase_weight(key, delta, max_weight) {
252 self.insert(
253 key,
254 item.clone(),
255 max_weight.map_or(delta, |m| delta.min(m)).max(1),
256 false,
257 );
258 }
259 self.evict()
260 }
261
262 fn remove(&self, item: CacheEntryKeyRef<'_>) {
263 let key = u64key(&item);
264 let node = self.lru.write().pop(&key);
265 if let Some(n) = node {
266 self.used.fetch_sub(n.size, Ordering::Relaxed);
267 self.items.fetch_sub(1, Ordering::Relaxed);
268 }
269 }
270
271 fn access(&self, item: &CacheEntryKey, size: usize, _fresh_until: SystemTime) -> bool {
272 let key = u64key(item);
273 if self.lru.write().get(&key).is_none() {
274 self.insert(key, item.clone(), size, false);
275 false
276 } else {
277 true
278 }
279 }
280
281 fn peek(&self, item: &CacheEntryKey) -> bool {
282 let key = u64key(item);
283 self.lru.read().peek(&key).is_some()
284 }
285
286 async fn save(&self, dir_path: &str) -> Result<()> {
287 let data = self.serialize()?;
288 let dir_str = dir_path.to_owned();
289 tokio::task::spawn_blocking(move || {
290 let dir_path = Path::new(&dir_str);
291 std::fs::create_dir_all(dir_path)
292 .or_err_with(InternalError, || format!("fail to create {dir_str}"))?;
293
294 let final_file_path = dir_path.join(FILE_NAME);
295 let random_suffix: u32 = rand::thread_rng().gen();
297 let temp_file_path = dir_path.join(format!("{}.{:08x}.tmp", FILE_NAME, random_suffix));
298 let mut file = File::create(&temp_file_path).or_err_with(InternalError, || {
299 format!("fail to create temporary file {}", temp_file_path.display())
300 })?;
301 file.write_all(&data).or_err_with(InternalError, || {
302 format!("fail to write to {}", temp_file_path.display())
303 })?;
304 file.flush().or_err_with(InternalError, || {
305 format!("fail to flush temp file {}", temp_file_path.display())
306 })?;
307 std::fs::rename(&temp_file_path, &final_file_path).or_err_with(InternalError, || {
308 format!(
309 "fail to rename temporary file {} to {}",
310 temp_file_path.display(),
311 final_file_path.display()
312 )
313 })
314 })
315 .await
316 .or_err(InternalError, "async blocking IO failure")?
317 }
318
319 async fn load(&self, dir_path: &str) -> Result<()> {
320 let dir_path = dir_path.to_owned();
321 let data = tokio::task::spawn_blocking(move || {
322 let file_path = Path::new(&dir_path).join(FILE_NAME);
323 let mut file = File::open(file_path.clone()).or_err_with(InternalError, || {
324 format!("fail to open {}", file_path.display())
325 })?;
326 let mut buffer = Vec::with_capacity(8192);
327 file.read_to_end(&mut buffer)
328 .or_err(InternalError, "fail to read from {file_path}")?;
329 Ok::<Vec<u8>, BError>(buffer)
330 })
331 .await
332 .or_err(InternalError, "async blocking IO failure")??;
333 self.deserialize(&data)
334 }
335}
336
337#[cfg(test)]
338impl Manager {
339 fn admit(
340 &self,
341 item: CompactCacheKey,
342 size: usize,
343 fresh_until: SystemTime,
344 ) -> Vec<CompactCacheKey> {
345 EvictionManager::admit(self, CacheEntryKey::key_only(item), size, fresh_until)
346 .into_iter()
347 .map(CacheEntryKey::into_key)
348 .collect()
349 }
350
351 fn increment_weight(
352 &self,
353 item: &CompactCacheKey,
354 delta: usize,
355 max_weight: Option<usize>,
356 ) -> Vec<CompactCacheKey> {
357 EvictionManager::increment_weight(
358 self,
359 &CacheEntryKey::key_only(item.clone()),
360 delta,
361 max_weight,
362 )
363 .into_iter()
364 .map(CacheEntryKey::into_key)
365 .collect()
366 }
367
368 fn remove(&self, item: &CompactCacheKey) {
369 EvictionManager::remove(self, CacheEntryKeyRef::from_entry_id(item, None));
370 }
371
372 fn access(&self, item: &CompactCacheKey, size: usize, fresh_until: SystemTime) -> bool {
373 EvictionManager::access(
374 self,
375 &CacheEntryKey::key_only(item.clone()),
376 size,
377 fresh_until,
378 )
379 }
380
381 fn peek(&self, item: &CompactCacheKey) -> bool {
382 EvictionManager::peek(self, &CacheEntryKey::key_only(item.clone()))
383 }
384}
385
386#[cfg(test)]
387mod test {
388 use super::*;
389 use crate::CacheKey;
390
391 #[test]
392 fn test_admission() {
393 let lru = Manager::new(4);
394 let key1 = CacheKey::new("a", "1").to_compact();
395 let until = SystemTime::now(); let v = lru.admit(key1.clone(), 1, until);
397 assert_eq!(v.len(), 0);
398 let key2 = CacheKey::new("b", "1").to_compact();
399 let v = lru.admit(key2.clone(), 2, until);
400 assert_eq!(v.len(), 0);
401 let key3 = CacheKey::new("c", "1").to_compact();
402 let v = lru.admit(key3, 1, until);
403 assert_eq!(v.len(), 0);
404
405 let key4 = CacheKey::new("d", "1").to_compact();
408 let v = lru.admit(key4, 2, until);
409 assert_eq!(v.len(), 2);
411 assert_eq!(v[0], key1);
412 assert_eq!(v[1], key2);
413 }
414
415 #[test]
416 fn test_identified_entries_are_distinct() {
417 let lru = Manager::new(1);
418 let key = CacheKey::new("a", "1").to_compact();
419 let first = CacheEntryKey::identified(key.clone(), crate::CacheEntryId::new(1));
420 let second = CacheEntryKey::identified(key, crate::CacheEntryId::new(2));
421 let until = SystemTime::now();
422
423 assert!(EvictionManager::admit(&lru, first.clone(), 1, until).is_empty());
424 assert_eq!(EvictionManager::admit(&lru, second, 1, until), vec![first]);
425 }
426
427 #[test]
428 fn test_access() {
429 let lru = Manager::new(4);
430 let key1 = CacheKey::new("a", "1").to_compact();
431 let until = SystemTime::now(); let v = lru.admit(key1.clone(), 1, until);
433 assert_eq!(v.len(), 0);
434 let key2 = CacheKey::new("b", "1").to_compact();
435 let v = lru.admit(key2.clone(), 2, until);
436 assert_eq!(v.len(), 0);
437 let key3 = CacheKey::new("c", "1").to_compact();
438 let v = lru.admit(key3, 1, until);
439 assert_eq!(v.len(), 0);
440
441 lru.access(&key1, 1, until);
444 assert_eq!(v.len(), 0);
445
446 let key4 = CacheKey::new("d", "1").to_compact();
447 let v = lru.admit(key4, 2, until);
448 assert_eq!(v.len(), 1);
449 assert_eq!(v[0], key2);
450 }
451
452 #[test]
453 fn test_remove() {
454 let lru = Manager::new(4);
455 let key1 = CacheKey::new("a", "1").to_compact();
456 let until = SystemTime::now(); let v = lru.admit(key1.clone(), 1, until);
458 assert_eq!(v.len(), 0);
459 let key2 = CacheKey::new("b", "1").to_compact();
460 let v = lru.admit(key2.clone(), 2, until);
461 assert_eq!(v.len(), 0);
462 let key3 = CacheKey::new("c", "1").to_compact();
463 let v = lru.admit(key3, 1, until);
464 assert_eq!(v.len(), 0);
465
466 lru.remove(&key1);
469
470 let key4 = CacheKey::new("d", "1").to_compact();
472 let v = lru.admit(key4, 2, until);
473 assert_eq!(v.len(), 1);
474 assert_eq!(v[0], key2);
475 }
476
477 #[test]
478 fn test_access_add() {
479 let lru = Manager::new(4);
480 let until = SystemTime::now(); let key1 = CacheKey::new("a", "1").to_compact();
483 lru.access(&key1, 1, until);
484 let key2 = CacheKey::new("b", "1").to_compact();
485 lru.access(&key2, 2, until);
486 let key3 = CacheKey::new("c", "1").to_compact();
487 lru.access(&key3, 2, until);
488
489 let key4 = CacheKey::new("d", "1").to_compact();
490 let v = lru.admit(key4, 2, until);
491 assert_eq!(v.len(), 2);
493 assert_eq!(v[0], key1);
494 assert_eq!(v[1], key2);
495 }
496
497 #[test]
498 fn test_increment_weight_adds_missing_item() {
499 let lru = Manager::new(4);
500 let key1 = CacheKey::new("a", "1").to_compact();
501 assert!(lru.increment_weight(&key1, 2, None).is_empty());
502 assert!(lru.peek(&key1));
503 assert_eq!(lru.total_size(), 2);
504 assert_eq!(lru.total_items(), 1);
505
506 let key2 = CacheKey::new("b", "1").to_compact();
507 let evicted = lru.increment_weight(&key2, 100, Some(3));
508 assert_eq!(evicted, vec![key1]);
509 assert!(lru.peek(&key2));
510 assert_eq!(lru.total_size(), 3);
511 assert_eq!(lru.total_items(), 1);
512 }
513
514 #[test]
515 fn test_increment_weight_admits_zero_and_does_not_shrink() {
516 let lru = Manager::new(10);
517
518 let key1 = CacheKey::new("a", "1").to_compact();
519 assert!(lru.increment_weight(&key1, 0, None).is_empty());
520 assert!(lru.peek(&key1));
521 assert_eq!(lru.total_size(), 1);
522
523 let key2 = CacheKey::new("b", "1").to_compact();
524 assert!(lru.increment_weight(&key2, 3, None).is_empty());
525 assert!(lru.increment_weight(&key2, 100, Some(2)).is_empty());
526 assert_eq!(lru.total_size(), 4);
527 assert_eq!(lru.total_items(), 2);
528 }
529
530 #[test]
531 fn test_admit_update() {
532 let lru = Manager::new(4);
533 let key1 = CacheKey::new("a", "1").to_compact();
534 let until = SystemTime::now(); let v = lru.admit(key1.clone(), 1, until);
536 assert_eq!(v.len(), 0);
537 let key2 = CacheKey::new("b", "1").to_compact();
538 let v = lru.admit(key2.clone(), 2, until);
539 assert_eq!(v.len(), 0);
540 let key3 = CacheKey::new("c", "1").to_compact();
541 let v = lru.admit(key3, 1, until);
542 assert_eq!(v.len(), 0);
543
544 let v = lru.admit(key2, 1, until);
547 assert_eq!(v.len(), 0);
548
549 let key4 = CacheKey::new("d", "1").to_compact();
551 let v = lru.admit(key4.clone(), 1, until);
552 assert_eq!(v.len(), 0);
553
554 let v = lru.admit(key4, 2, until);
556 assert_eq!(v.len(), 1);
558 assert_eq!(v[0], key1);
559 }
560
561 #[test]
562 fn test_serde() {
563 let lru = Manager::new(4);
564 let key1 = CacheKey::new("a", "1").to_compact();
565 let until = SystemTime::now(); let v = lru.admit(key1.clone(), 1, until);
567 assert_eq!(v.len(), 0);
568 let key2 = CacheKey::new("b", "1").to_compact();
569 let v = lru.admit(key2.clone(), 2, until);
570 assert_eq!(v.len(), 0);
571 let key3 = CacheKey::new("c", "1").to_compact();
572 let v = lru.admit(key3, 1, until);
573 assert_eq!(v.len(), 0);
574
575 lru.access(&key1, 1, until);
578 assert_eq!(v.len(), 0);
579
580 let ser = lru.serialize().unwrap();
582 let lru2 = Manager::new(4);
583 lru2.deserialize(&ser).unwrap();
584
585 let key4 = CacheKey::new("d", "1").to_compact();
586 let v = lru2.admit(key4, 2, until);
587 assert_eq!(v.len(), 1);
588 assert_eq!(v[0], key2);
589 }
590
591 #[tokio::test]
592 async fn test_save_to_disk() {
593 let lru = Manager::new(4);
594 let key1 = CacheKey::new("a", "1").to_compact();
595 let until = SystemTime::now(); let v = lru.admit(key1.clone(), 1, until);
597 assert_eq!(v.len(), 0);
598 let key2 = CacheKey::new("b", "1").to_compact();
599 let v = lru.admit(key2.clone(), 2, until);
600 assert_eq!(v.len(), 0);
601 let key3 = CacheKey::new("c", "1").to_compact();
602 let v = lru.admit(key3, 1, until);
603 assert_eq!(v.len(), 0);
604
605 lru.access(&key1, 1, until);
608 assert_eq!(v.len(), 0);
609
610 lru.save("/tmp/test_simple_lru_save").await.unwrap();
612 let lru2 = Manager::new(4);
613 lru2.load("/tmp/test_simple_lru_save").await.unwrap();
614
615 let key4 = CacheKey::new("d", "1").to_compact();
616 let v = lru2.admit(key4, 2, until);
617 assert_eq!(v.len(), 1);
618 assert_eq!(v[0], key2);
619 }
620
621 #[test]
622 fn test_watermark_eviction() {
623 const SIZE_LIMIT: usize = usize::MAX / 2;
624 let lru = Manager::new_with_watermark(SIZE_LIMIT, Some(4));
625 let until = SystemTime::now();
626
627 for name in ["a", "b", "c", "d", "e", "f"] {
629 let key = CacheKey::new(name, "1").to_compact();
630 let _ = lru.admit(key, 1, until);
631 }
632
633 assert_eq!(lru.total_items(), 4);
635 assert_eq!(lru.evicted_items(), 2);
636 assert_eq!(lru.evicted_size(), 2);
637 assert!(lru.total_size() <= SIZE_LIMIT);
638 }
639}