Skip to main content

ordinary_storage/stores/
cache.rs

1// Copyright (C) 2026 The Ordinary Authors.
2//
3// SPDX-License-Identifier: BSD-3-Clause
4
5use bytes::{BufMut, Bytes, BytesMut};
6use hashbrown::{DefaultHashBuilder, HashMap, HashSet};
7use ordinary_config::{
8    CacheLimits, StoredCache as StoredCacheConfig, StoredCache, StoredCachePolicy,
9};
10use parking_lot::Mutex;
11use quick_cache::sync::Cache;
12use quick_cache::{Lifecycle, UnitWeighter};
13use saferlmdb::{
14    self as lmdb, Database, DatabaseOptions, Environment, ReadTransaction, WriteTransaction, put,
15};
16use smallvec::SmallVec;
17use std::sync::Arc;
18use tracing::instrument;
19use uuid::Uuid;
20
21pub struct Lookup<'a> {
22    pub base: Bytes,
23    pub checks: &'a [Bytes],
24}
25
26impl<'a> Lookup<'a> {
27    pub fn new(base: Bytes, checks: &'a [Bytes]) -> Self {
28        Lookup { base, checks }
29    }
30
31    pub fn set(&self, prefix: &mut BytesMut) -> HashSet<Bytes> {
32        let prefix_len = prefix.len();
33
34        let mut set = HashSet::new();
35
36        for check in self.checks {
37            prefix.put_slice(check.as_ref());
38            set.insert(prefix.clone().into());
39            prefix.truncate(prefix_len);
40        }
41
42        set
43    }
44}
45
46#[derive(Clone)]
47struct EvictionHandler {
48    env: Arc<Environment>,
49    cache_db: Arc<Database<'static>>,
50    inventory: Arc<Mutex<(u64, usize)>>,
51}
52
53impl Lifecycle<Bytes, Option<Bytes>> for EvictionHandler {
54    type RequestState = ();
55
56    fn on_evict(&self, _state: &mut Self::RequestState, key: Bytes, _val: Option<Bytes>) {
57        if let Ok(txn) = WriteTransaction::new(self.env.clone()) {
58            {
59                let mut access = txn.access();
60
61                match access.get::<[u8], [u8]>(&self.cache_db, key.as_ref()) {
62                    Ok(val) => {
63                        let mut lock = self.inventory.lock();
64
65                        lock.0 -= key.len() as u64;
66                        lock.0 -= val.len() as u64;
67                        lock.1 -= 1;
68                    }
69                    Err(err) => {
70                        tracing::warn!(%err);
71                    }
72                }
73
74                if let Err(err) = access.del_key::<[u8]>(&self.cache_db, key.as_ref()) {
75                    tracing::warn!(%err);
76                }
77            }
78
79            if let Err(err) = txn.commit() {
80                tracing::error!(%err);
81            }
82        }
83    }
84}
85
86pub enum CacheKind {
87    Http,
88    Function,
89}
90
91impl CacheKind {
92    fn as_u8(&self) -> u8 {
93        match self {
94            CacheKind::Http => 0,
95            CacheKind::Function => 1,
96        }
97    }
98}
99
100#[derive(Debug, Clone, Eq, Hash, PartialEq)]
101pub enum CacheDependency {
102    DatabaseModel([u8; 16]),
103}
104
105/// `(cache_kind, service_idx, cache_key)`
106type DependencyMap = Arc<Mutex<HashMap<CacheDependency, HashSet<Bytes>>>>;
107
108pub struct CacheStore {
109    /// `(service_kind, service_idx, key_bytes) -> ()`
110    ///
111    /// in the case of templates the `key_bytes` will likely
112    /// be `host|method|matched_path|path_params?|query_params?|incoming_body_hash?|(etag OR content_type|compression) -> Option<last_modified>`.
113    quick_cache:
114        Arc<Cache<Bytes, Option<Bytes>, UnitWeighter, DefaultHashBuilder, EvictionHandler>>,
115
116    pub limits: CacheLimits,
117    env: Arc<Environment>,
118
119    /// `key_bytes -> val_bytes`
120    ///
121    /// in the case of templates the `val_bytes` will likely be a flexbuffer
122    /// with an `etag`, `last-modified` and `response-body`.
123    ///
124    /// and the `key_bytes` will be `[host, segments, params, compression]`.
125    cache_db: Arc<Database<'static>>,
126
127    log_sizes: bool,
128
129    dependency_map: DependencyMap,
130
131    inventory: Arc<Mutex<(u64, usize)>>,
132}
133
134impl CacheStore {
135    #[allow(clippy::too_many_lines, clippy::missing_panics_doc)]
136    pub fn new(
137        limits: CacheLimits,
138        env: &Arc<Environment>,
139        log_sizes: bool,
140    ) -> anyhow::Result<Self> {
141        let inventory = Arc::new(Mutex::new((0, 0)));
142
143        let cache_db = Arc::new(Database::open(
144            env.clone(),
145            Some("cache"),
146            &DatabaseOptions::new(lmdb::db::Flags::CREATE),
147        )?);
148
149        let quick_cache = Arc::new(Cache::with(
150            2000,
151            2000,
152            UnitWeighter,
153            DefaultHashBuilder::default(),
154            EvictionHandler {
155                env: env.clone(),
156                cache_db: cache_db.clone(),
157                inventory: inventory.clone(),
158            },
159        ));
160
161        // !! [start] clean out the whole lmdb cache
162        let txn = WriteTransaction::new(env.clone())?;
163
164        {
165            let mut access = txn.access();
166            let mut cursor = txn.cursor(cache_db.clone())?;
167
168            let mut keys = vec![];
169
170            if let Ok((key, _val)) = cursor.first::<[u8], [u8]>(&access) {
171                keys.push(key.to_vec());
172            }
173
174            while let Ok((key, _val)) = cursor.next::<[u8], [u8]>(&access) {
175                keys.push(key.to_vec());
176            }
177
178            for key in keys {
179                access.del_key(&cache_db, key.as_slice())?;
180            }
181        }
182
183        txn.commit()?;
184        // !! [end] clean out the whole lmdb cache
185
186        let mut dep_map: HashMap<CacheDependency, HashSet<Bytes>> = HashMap::new();
187
188        let mut inventory_lock = inventory.lock();
189
190        let txn = ReadTransaction::new(env.clone())?;
191
192        let access = txn.access();
193        let mut cursor = txn.cursor(cache_db.clone())?;
194
195        if let Ok((key, val)) = cursor.first::<[u8], [u8]>(&access) {
196            inventory_lock.0 += key.len() as u64;
197            inventory_lock.0 += val.len() as u64;
198            inventory_lock.1 += 1;
199
200            Self::seed_from_lmdb_entry(&quick_cache, &mut dep_map, key, val)?;
201        }
202
203        while let Ok((key, val)) = cursor.next::<[u8], [u8]>(&access) {
204            inventory_lock.0 += key.len() as u64;
205            inventory_lock.0 += val.len() as u64;
206            inventory_lock.1 += 1;
207
208            Self::seed_from_lmdb_entry(&quick_cache, &mut dep_map, key, val)?;
209        }
210
211        tracing::info!(
212            entries = inventory_lock.1,
213            size = log_sizes.then_some(display(
214                bytesize::ByteSize(inventory_lock.0).display().si_short()
215            ))
216        );
217
218        drop(inventory_lock);
219
220        Ok(Self {
221            quick_cache,
222            limits,
223            env: env.clone(),
224            cache_db,
225            log_sizes,
226            dependency_map: Arc::new(Mutex::new(dep_map)),
227            inventory,
228        })
229    }
230
231    #[allow(clippy::similar_names)]
232    fn seed_from_lmdb_entry(
233        quick_cache: &Arc<
234            Cache<Bytes, Option<Bytes>, UnitWeighter, DefaultHashBuilder, EvictionHandler>,
235        >,
236        dependency_map: &mut HashMap<CacheDependency, HashSet<Bytes>>,
237        key: &[u8],
238        val: &[u8],
239    ) -> anyhow::Result<()> {
240        let root = flexbuffers::Reader::get_root(val)?;
241        let root_vec = root.as_vector();
242
243        let internal = root_vec.idx(0).as_vector();
244
245        let deps_vec = internal.idx(0).as_vector();
246
247        for dep in &deps_vec {
248            let dep_vec = dep.as_vector();
249
250            let dep_kind = dep_vec.idx(0).as_u8();
251
252            let dep = if dep_kind == 0 {
253                let uuid = Uuid::from_slice(dep_vec.idx(1).as_blob().0)?;
254                CacheDependency::DatabaseModel(*uuid.as_bytes())
255            } else {
256                continue;
257            };
258
259            if let Some(inverse_set) = dependency_map.get_mut(&dep) {
260                inverse_set.insert(Bytes::copy_from_slice(key));
261            } else {
262                let mut inverse_set = HashSet::new();
263                inverse_set.insert(Bytes::copy_from_slice(key));
264
265                dependency_map.insert(dep, inverse_set);
266            }
267        }
268
269        let is_quick_cache = internal.idx(1).as_bool();
270
271        if is_quick_cache {
272            let keep_in_memory = internal.idx(2).as_bool();
273
274            if keep_in_memory {
275                quick_cache.insert(
276                    Bytes::copy_from_slice(key),
277                    Some(Bytes::copy_from_slice(root_vec.idx(1).as_blob().0)),
278                );
279            } else {
280                quick_cache.insert(Bytes::copy_from_slice(key), None);
281            }
282        }
283        Ok(())
284    }
285
286    /// Check the `cache_db` for a hit; returns `Ok(None)` if not.
287    #[allow(clippy::type_complexity)]
288    #[instrument(skip_all, err)]
289    pub fn check(
290        &self,
291        config: &StoredCacheConfig,
292        cache_kind: CacheKind,
293        lookup: Lookup,
294    ) -> anyhow::Result<Option<Bytes>> {
295        let mut cache_key_mut = BytesMut::new();
296        cache_key_mut.put_u8(cache_kind.as_u8());
297
298        let mut rows = 0;
299
300        match config.policy {
301            StoredCachePolicy::Permanent => {
302                let txn = ReadTransaction::new(self.env.clone())?;
303                let access = txn.access();
304
305                let mut cursor = txn.cursor(self.cache_db.clone())?;
306                let checks_set = lookup.set(&mut cache_key_mut);
307                cache_key_mut.put_slice(lookup.base.as_ref());
308
309                if let Ok((key, val)) =
310                    cursor.seek_range_k::<[u8], [u8]>(&access, cache_key_mut.as_ref())
311                {
312                    rows += 1;
313
314                    if !key.starts_with(cache_key_mut.as_ref()) {
315                        tracing::info!(hit = false, rows);
316                        return Ok(None);
317                    }
318
319                    if checks_set.contains(key) {
320                        let root = flexbuffers::Reader::get_root(val)?;
321                        let val = root.as_vector().idx(1).as_blob();
322
323                        tracing::info!(hit = true, rows);
324                        return Ok(Some(Bytes::copy_from_slice(val.0)));
325                    }
326                } else {
327                    tracing::info!(hit = false, rows);
328                    return Ok(None);
329                }
330
331                while let Ok((key, val)) = cursor.next::<[u8], [u8]>(&access) {
332                    rows += 1;
333
334                    if !key.starts_with(cache_key_mut.as_ref()) {
335                        tracing::info!(hit = false, rows);
336                        return Ok(None);
337                    }
338
339                    if checks_set.contains(key) {
340                        let root = flexbuffers::Reader::get_root(val)?;
341                        let val = root.as_vector().idx(1).as_blob();
342
343                        tracing::info!(hit = true, rows);
344                        return Ok(Some(Bytes::copy_from_slice(val.0)));
345                    }
346                }
347
348                tracing::info!(hit = false, rows);
349                Ok(None)
350            }
351            StoredCachePolicy::QuickCache => {
352                for check in lookup.checks {
353                    rows += 1;
354
355                    cache_key_mut.put_slice(check.as_ref());
356
357                    if let Some(val) = self.quick_cache.get(cache_key_mut.as_ref()) {
358                        if let Some(val) = val {
359                            tracing::info!(hit = true, rows);
360                            return Ok(Some(val));
361                        }
362
363                        let txn = ReadTransaction::new(self.env.clone())?;
364                        let access = txn.access();
365
366                        if let Ok(res) =
367                            access.get::<[u8], [u8]>(&self.cache_db, cache_key_mut.as_ref())
368                        {
369                            let root = flexbuffers::Reader::get_root(res)?;
370                            let val = root.as_vector().idx(1).as_blob();
371
372                            tracing::info!(hit = true, rows);
373                            return Ok(Some(Bytes::copy_from_slice(val.0)));
374                        }
375                    }
376
377                    cache_key_mut.truncate(1);
378                }
379
380                tracing::info!(hit = false, rows);
381                Ok(None)
382            }
383        }
384    }
385
386    /// Caches item for kind and index for specified cache key.
387    ///
388    /// LMDB format (`cache_db`)
389    /// ```null
390    /// uuid_v7 ->
391    /// - internal 0 (vector)
392    ///   - deps 0 (vector)
393    ///   - quick_cache 1 (bool)
394    ///   - value_in_memory 1 (bool)
395    /// - value 1 (blob)
396    /// ```
397    #[allow(
398        clippy::too_many_lines,
399        clippy::too_many_arguments,
400        clippy::similar_names
401    )]
402    #[instrument(skip_all, err)]
403    pub fn write(
404        &self,
405        config: &StoredCacheConfig,
406        cache_kind: CacheKind,
407        cache_key: &[u8],
408        cache_val: &[u8],
409        dependencies: Option<SmallVec<[CacheDependency; 13]>>,
410        value_in_memory: bool,
411    ) -> anyhow::Result<()> {
412        let mut cache_key_mut = BytesMut::new();
413
414        cache_key_mut.put_u8(cache_kind.as_u8());
415        cache_key_mut.put(cache_key);
416
417        let cache_key: Bytes = cache_key_mut.into();
418
419        let mut builder = flexbuffers::Builder::new(&flexbuffers::BuilderOptions::SHARE_NONE);
420        let mut builder_vec = builder.start_vector();
421
422        let mut internal_vec = builder_vec.start_vector();
423
424        let mut deps_vec = internal_vec.start_vector();
425
426        if config.evict_on_dependency_change == Some(true)
427            && let Some(dependencies) = dependencies
428        {
429            let mut dep_lock = self.dependency_map.lock();
430
431            for dep in dependencies {
432                let mut dep_vec = deps_vec.start_vector();
433
434                match dep {
435                    CacheDependency::DatabaseModel(uuid) => {
436                        dep_vec.push(0u8);
437                        dep_vec.push(flexbuffers::Blob(uuid.as_ref()));
438                    }
439                }
440
441                dep_vec.end_vector();
442
443                if let Some(inverse_set) = dep_lock.get_mut(&dep) {
444                    inverse_set.insert(cache_key.clone());
445                } else {
446                    let mut inverse_set = HashSet::new();
447                    inverse_set.insert(cache_key.clone());
448
449                    (*dep_lock).insert(dep, inverse_set);
450                }
451            }
452        }
453
454        deps_vec.end_vector();
455
456        if let StoredCachePolicy::QuickCache = config.policy {
457            internal_vec.push(true);
458            internal_vec.push(value_in_memory);
459        } else {
460            internal_vec.push(false);
461        }
462
463        internal_vec.end_vector();
464
465        builder_vec.push(flexbuffers::Blob(cache_val));
466        builder_vec.end_vector();
467
468        let val = builder.view();
469        let size = (val.len() + cache_key.len()) as u64;
470
471        let mut lock = self.inventory.lock();
472
473        if let Some(max_size) = config.max_size
474            && size + lock.0 > max_size
475        {
476            tracing::warn!("item causes 'max_size' to be exceeded");
477            return Ok(());
478        }
479
480        if let Some(max_count) = config.max_count
481            && 1 + lock.1 > max_count
482        {
483            tracing::warn!("item causes 'max_count' to be exceeded");
484            return Ok(());
485        }
486
487        lock.0 += size;
488        lock.1 += 1;
489
490        let storage_size = lock.0;
491        let entries = lock.1;
492
493        drop(lock);
494
495        let txn = WriteTransaction::new(self.env.clone())?;
496
497        {
498            let mut access = txn.access();
499            access.put::<[u8], [u8]>(
500                &self.cache_db,
501                cache_key.as_ref(),
502                val,
503                &put::Flags::empty(),
504            )?;
505        }
506
507        txn.commit()?;
508
509        if let StoredCachePolicy::QuickCache = config.policy {
510            self.quick_cache.insert(
511                cache_key.clone(),
512                value_in_memory.then_some(Bytes::copy_from_slice(cache_val)),
513            );
514        }
515
516        tracing::info!(
517            total.entries = entries,
518            total.size = self.log_sizes.then_some(display(
519                bytesize::ByteSize(storage_size).display().si_short()
520            )),
521            item.size = self
522                .log_sizes
523                .then_some(display(bytesize::ByteSize(size).display().si_short())),
524            "stored"
525        );
526
527        Ok(())
528    }
529
530    #[instrument(skip_all, err)]
531    pub async fn dependency_evict(&self, dependencies: Vec<CacheDependency>) -> anyhow::Result<()> {
532        {
533            let txn = WriteTransaction::new(self.env.clone())?;
534
535            let mut lock_dep_map = self.dependency_map.lock();
536            let mut lock_inventory = self.inventory.lock();
537
538            {
539                let mut access = txn.access();
540
541                for dependency in dependencies {
542                    if let Some(addrs) = lock_dep_map.get(&dependency) {
543                        for cache_key in addrs {
544                            {
545                                tracing::debug!("evicting for dependency");
546
547                                if self.quick_cache.remove(cache_key).is_none() {
548                                    match access
549                                        .get::<[u8], [u8]>(&self.cache_db, cache_key.as_ref())
550                                    {
551                                        Ok(val) => {
552                                            lock_inventory.0 -= cache_key.len() as u64;
553                                            lock_inventory.0 -= val.len() as u64;
554                                            lock_inventory.1 -= 1;
555                                        }
556                                        Err(err) => {
557                                            tracing::warn!(%err);
558                                        }
559                                    }
560
561                                    if let Err(err) =
562                                        access.del_key(&self.cache_db, cache_key.as_ref())
563                                    {
564                                        tracing::warn!(%err);
565                                    }
566                                }
567                            }
568                        }
569                    }
570
571                    lock_dep_map.remove(&dependency);
572                }
573            }
574
575            txn.commit()?;
576        }
577
578        Ok(())
579    }
580
581    #[instrument(skip_all, err)]
582    pub async fn artifact_evict(
583        &self,
584        config: &StoredCache,
585        kind: CacheKind,
586        idx: u8,
587    ) -> anyhow::Result<()> {
588        let mut key = BytesMut::new();
589
590        key.put_u8(kind.as_u8());
591        key.put_u8(idx);
592
593        let mut lock_inventory = self.inventory.lock();
594
595        let txn = WriteTransaction::new(self.env.clone())?;
596
597        {
598            let mut access = txn.access();
599            let mut cursor = txn.cursor(self.cache_db.clone())?;
600
601            let mut keys = vec![];
602
603            if let Ok((key, val)) = cursor.seek_range_k::<[u8], [u8]>(&access, key.as_ref()) {
604                lock_inventory.0 -= key.len() as u64;
605                lock_inventory.0 -= val.len() as u64;
606                lock_inventory.1 -= 1;
607
608                keys.push(Bytes::copy_from_slice(key));
609            }
610
611            while let Ok((key, val)) = cursor.next::<[u8], [u8]>(&access) {
612                lock_inventory.0 -= key.len() as u64;
613                lock_inventory.0 -= val.len() as u64;
614                lock_inventory.1 -= 1;
615
616                keys.push(Bytes::copy_from_slice(key));
617            }
618
619            for key in keys {
620                match config.policy {
621                    StoredCachePolicy::Permanent => {
622                        if let Err(err) = access.del_key::<[u8]>(&self.cache_db, key.as_ref()) {
623                            tracing::error!(%err);
624                        }
625                    }
626                    StoredCachePolicy::QuickCache => {
627                        // lmdb cleanup happens in the eviction handler
628                        self.quick_cache.remove(&key);
629                    }
630                }
631            }
632        }
633
634        txn.commit()?;
635
636        Ok(())
637    }
638}