Skip to main content

ordinary_storage/stores/
cache.rs

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