Skip to main content

pingora_cache/eviction/
simple_lru.rs

1// Copyright 2026 Cloudflare, Inc.
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7// http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! A simple LRU cache manager built on top of the `lru` crate
16
17use 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
42/// A simple LRU eviction manager
43///
44/// The implementation is not optimized. All operations require global locks.
45pub 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    /// Create a new [Manager] with the given total size limit `limit`.
57    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    /// Create a new [Manager] with optional watermark in addition to size limit `limit`.
70    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            // replacing a node, just need to update used size
95            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, // same size, update nothing, use 0 to match other arms' type
99            };
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    // evict items until the used capacity is below the size limit and watermark count
129    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); // we will at least pop 1 item
139
140        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                // lru empty
149                return to_evict;
150            }
151        }
152        to_evict
153    }
154
155    // This could use a lot of memory to buffer the serialized data in memory and could lock the LRU
156    // for too long
157    fn serialize(&self) -> Result<Vec<u8>> {
158        use rmp_serde::encode::Serializer;
159        use serde::ser::SerializeSeq;
160        use serde::ser::Serializer as _;
161        // NOTE: This could use a lot of memory to buffer the serialized data in memory
162        let mut ser = Serializer::new(vec![]);
163        // NOTE: This long for loop could lock the LRU for too long
164        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(); // write to vec, safe
170        }
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); // insert in the back
204        }
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            // create a temporary filename using a randomized u32 hash to minimize the chance of multiple writers writing to the same tmp file
296            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(); // unused value as a placeholder
396        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        // lru si full (4) now
406
407        let key4 = CacheKey::new("d", "1").to_compact();
408        let v = lru.admit(key4, 2, until);
409        // need to reduce used by at least 2, both key1 and key2 are evicted to make room for 3
410        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(); // unused value as a placeholder
432        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 is full (4) now
442        // make key1 most recently used
443        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(); // unused value as a placeholder
457        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 is full (4) now
467        // remove key1
468        lru.remove(&key1);
469
470        // key2 is the least recently used one now
471        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(); // unused value as a placeholder
481
482        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        // need to reduce used by at least 2, both key1 and key2 are evicted to make room for 3
492        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(); // unused value as a placeholder
535        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        // lru is full (4) now
545        // update key2 to reduce its size by 1
546        let v = lru.admit(key2, 1, until);
547        assert_eq!(v.len(), 0);
548
549        // lru is not full anymore
550        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        // make key4 larger
555        let v = lru.admit(key4, 2, until);
556        // need to evict now
557        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(); // unused value as a placeholder
566        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 is full (4) now
576        // make key1 most recently used
577        lru.access(&key1, 1, until);
578        assert_eq!(v.len(), 0);
579
580        // load lru2 with lru's data
581        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(); // unused value as a placeholder
596        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 is full (4) now
606        // make key1 most recently used
607        lru.access(&key1, 1, until);
608        assert_eq!(v.len(), 0);
609
610        // load lru2 with lru's data
611        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        // admit 6 items of size 1
628        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        // test items were evicted due to watermark
634        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}