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 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
82type DependencyMap = Arc<Mutex<HashMap<CacheDependency, HashSet<Bytes>>>>;
84
85pub struct CacheStore {
86 quick_cache: Arc<Cache<Bytes, (), UnitWeighter, DefaultHashBuilder, EvictionHandler>>,
91
92 pub limits: CacheLimits,
93 env: Arc<Environment>,
94
95 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 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 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 #[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 #[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 self.quick_cache.remove(&key);
542 }
543 }
544 }
545 }
546
547 txn.commit()?;
548
549 Ok(())
550 }
551}