1use 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
105type DependencyMap = Arc<Mutex<HashMap<CacheDependency, HashSet<Bytes>>>>;
107
108pub struct CacheStore {
109 quick_cache:
114 Arc<Cache<Bytes, Option<Bytes>, UnitWeighter, DefaultHashBuilder, EvictionHandler>>,
115
116 pub limits: CacheLimits,
117 env: Arc<Environment>,
118
119 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 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 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 #[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 #[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 self.quick_cache.remove(&key);
629 }
630 }
631 }
632 }
633
634 txn.commit()?;
635
636 Ok(())
637 }
638}