use bytes::{BufMut, Bytes, BytesMut};
use hashbrown::{DefaultHashBuilder, HashMap, HashSet};
use ordinary_config::{
CacheLimits, StoredCache as StoredCacheConfig, StoredCache, StoredCachePolicy,
};
use parking_lot::Mutex;
use quick_cache::sync::Cache;
use quick_cache::{Lifecycle, UnitWeighter};
use saferlmdb::{
self as lmdb, Database, DatabaseOptions, Environment, ReadTransaction, WriteTransaction, put,
};
use smallvec::SmallVec;
use std::sync::Arc;
use tracing::instrument;
use uuid::Uuid;
pub struct Lookup<'a> {
pub base: Bytes,
pub checks: &'a [Bytes],
}
impl<'a> Lookup<'a> {
pub fn new(base: Bytes, checks: &'a [Bytes]) -> Self {
Lookup { base, checks }
}
pub fn set(&self, prefix: &mut BytesMut) -> HashSet<Bytes> {
let prefix_len = prefix.len();
let mut set = HashSet::new();
for check in self.checks {
prefix.put_slice(check.as_ref());
set.insert(prefix.clone().into());
prefix.truncate(prefix_len);
}
set
}
}
#[derive(Clone)]
struct EvictionHandler {
env: Arc<Environment>,
cache_db: Arc<Database<'static>>,
inventory: Arc<Mutex<(u64, usize)>>,
}
impl Lifecycle<Bytes, Option<Bytes>> for EvictionHandler {
type RequestState = ();
fn on_evict(&self, _state: &mut Self::RequestState, key: Bytes, _val: Option<Bytes>) {
if let Ok(txn) = WriteTransaction::new(self.env.clone()) {
{
let mut access = txn.access();
match access.get::<[u8], [u8]>(&self.cache_db, key.as_ref()) {
Ok(val) => {
let mut lock = self.inventory.lock();
lock.0 -= key.len() as u64;
lock.0 -= val.len() as u64;
lock.1 -= 1;
}
Err(err) => {
tracing::warn!(%err);
}
}
if let Err(err) = access.del_key::<[u8]>(&self.cache_db, key.as_ref()) {
tracing::warn!(%err);
}
}
if let Err(err) = txn.commit() {
tracing::error!(%err);
}
}
}
}
pub enum CacheKind {
Http,
Function,
}
impl CacheKind {
fn as_u8(&self) -> u8 {
match self {
CacheKind::Http => 0,
CacheKind::Function => 1,
}
}
}
#[derive(Debug, Clone, Eq, Hash, PartialEq)]
pub enum CacheDependency {
DatabaseModel([u8; 16]),
}
type DependencyMap = Arc<Mutex<HashMap<CacheDependency, HashSet<Bytes>>>>;
pub struct CacheStore {
quick_cache:
Arc<Cache<Bytes, Option<Bytes>, UnitWeighter, DefaultHashBuilder, EvictionHandler>>,
pub limits: CacheLimits,
env: Arc<Environment>,
cache_db: Arc<Database<'static>>,
log_sizes: bool,
dependency_map: DependencyMap,
inventory: Arc<Mutex<(u64, usize)>>,
}
impl CacheStore {
#[allow(clippy::too_many_lines, clippy::missing_panics_doc)]
pub fn new(
limits: CacheLimits,
env: &Arc<Environment>,
log_sizes: bool,
) -> anyhow::Result<Self> {
let inventory = Arc::new(Mutex::new((0, 0)));
let cache_db = Arc::new(Database::open(
env.clone(),
Some("cache"),
&DatabaseOptions::new(lmdb::db::Flags::CREATE),
)?);
let quick_cache = Arc::new(Cache::with(
2000,
2000,
UnitWeighter,
DefaultHashBuilder::default(),
EvictionHandler {
env: env.clone(),
cache_db: cache_db.clone(),
inventory: inventory.clone(),
},
));
let txn = WriteTransaction::new(env.clone())?;
{
let mut access = txn.access();
let mut cursor = txn.cursor(cache_db.clone())?;
let mut keys = vec![];
if let Ok((key, _val)) = cursor.first::<[u8], [u8]>(&access) {
keys.push(key.to_vec());
}
while let Ok((key, _val)) = cursor.next::<[u8], [u8]>(&access) {
keys.push(key.to_vec());
}
for key in keys {
access.del_key(&cache_db, key.as_slice())?;
}
}
txn.commit()?;
let mut dep_map: HashMap<CacheDependency, HashSet<Bytes>> = HashMap::new();
let mut inventory_lock = inventory.lock();
let txn = ReadTransaction::new(env.clone())?;
let access = txn.access();
let mut cursor = txn.cursor(cache_db.clone())?;
if let Ok((key, val)) = cursor.first::<[u8], [u8]>(&access) {
inventory_lock.0 += key.len() as u64;
inventory_lock.0 += val.len() as u64;
inventory_lock.1 += 1;
Self::seed_from_lmdb_entry(&quick_cache, &mut dep_map, key, val)?;
}
while let Ok((key, val)) = cursor.next::<[u8], [u8]>(&access) {
inventory_lock.0 += key.len() as u64;
inventory_lock.0 += val.len() as u64;
inventory_lock.1 += 1;
Self::seed_from_lmdb_entry(&quick_cache, &mut dep_map, key, val)?;
}
tracing::info!(
entries = inventory_lock.1,
size = log_sizes.then_some(display(
bytesize::ByteSize(inventory_lock.0).display().si_short()
))
);
drop(inventory_lock);
Ok(Self {
quick_cache,
limits,
env: env.clone(),
cache_db,
log_sizes,
dependency_map: Arc::new(Mutex::new(dep_map)),
inventory,
})
}
#[allow(clippy::similar_names)]
fn seed_from_lmdb_entry(
quick_cache: &Arc<
Cache<Bytes, Option<Bytes>, UnitWeighter, DefaultHashBuilder, EvictionHandler>,
>,
dependency_map: &mut HashMap<CacheDependency, HashSet<Bytes>>,
key: &[u8],
val: &[u8],
) -> anyhow::Result<()> {
let root = flexbuffers::Reader::get_root(val)?;
let root_vec = root.as_vector();
let internal = root_vec.idx(0).as_vector();
let deps_vec = internal.idx(0).as_vector();
for dep in &deps_vec {
let dep_vec = dep.as_vector();
let dep_kind = dep_vec.idx(0).as_u8();
let dep = if dep_kind == 0 {
let uuid = Uuid::from_slice(dep_vec.idx(1).as_blob().0)?;
CacheDependency::DatabaseModel(*uuid.as_bytes())
} else {
continue;
};
if let Some(inverse_set) = dependency_map.get_mut(&dep) {
inverse_set.insert(Bytes::copy_from_slice(key));
} else {
let mut inverse_set = HashSet::new();
inverse_set.insert(Bytes::copy_from_slice(key));
dependency_map.insert(dep, inverse_set);
}
}
let is_quick_cache = internal.idx(1).as_bool();
if is_quick_cache {
let keep_in_memory = internal.idx(2).as_bool();
if keep_in_memory {
quick_cache.insert(
Bytes::copy_from_slice(key),
Some(Bytes::copy_from_slice(root_vec.idx(1).as_blob().0)),
);
} else {
quick_cache.insert(Bytes::copy_from_slice(key), None);
}
}
Ok(())
}
#[allow(clippy::type_complexity)]
#[instrument(skip_all, err)]
pub fn check(
&self,
config: &StoredCacheConfig,
cache_kind: CacheKind,
lookup: Lookup,
) -> anyhow::Result<Option<Bytes>> {
let mut cache_key_mut = BytesMut::new();
cache_key_mut.put_u8(cache_kind.as_u8());
let mut rows = 0;
match config.policy {
StoredCachePolicy::Permanent => {
let txn = ReadTransaction::new(self.env.clone())?;
let access = txn.access();
let mut cursor = txn.cursor(self.cache_db.clone())?;
let checks_set = lookup.set(&mut cache_key_mut);
cache_key_mut.put_slice(lookup.base.as_ref());
if let Ok((key, val)) =
cursor.seek_range_k::<[u8], [u8]>(&access, cache_key_mut.as_ref())
{
rows += 1;
if !key.starts_with(cache_key_mut.as_ref()) {
tracing::info!(hit = false, rows);
return Ok(None);
}
if checks_set.contains(key) {
let root = flexbuffers::Reader::get_root(val)?;
let val = root.as_vector().idx(1).as_blob();
tracing::info!(hit = true, rows);
return Ok(Some(Bytes::copy_from_slice(val.0)));
}
} else {
tracing::info!(hit = false, rows);
return Ok(None);
}
while let Ok((key, val)) = cursor.next::<[u8], [u8]>(&access) {
rows += 1;
if !key.starts_with(cache_key_mut.as_ref()) {
tracing::info!(hit = false, rows);
return Ok(None);
}
if checks_set.contains(key) {
let root = flexbuffers::Reader::get_root(val)?;
let val = root.as_vector().idx(1).as_blob();
tracing::info!(hit = true, rows);
return Ok(Some(Bytes::copy_from_slice(val.0)));
}
}
tracing::info!(hit = false, rows);
Ok(None)
}
StoredCachePolicy::QuickCache => {
for check in lookup.checks {
rows += 1;
cache_key_mut.put_slice(check.as_ref());
if let Some(val) = self.quick_cache.get(cache_key_mut.as_ref()) {
if let Some(val) = val {
tracing::info!(hit = true, rows);
return Ok(Some(val));
}
let txn = ReadTransaction::new(self.env.clone())?;
let access = txn.access();
if let Ok(res) =
access.get::<[u8], [u8]>(&self.cache_db, cache_key_mut.as_ref())
{
let root = flexbuffers::Reader::get_root(res)?;
let val = root.as_vector().idx(1).as_blob();
tracing::info!(hit = true, rows);
return Ok(Some(Bytes::copy_from_slice(val.0)));
}
}
cache_key_mut.truncate(1);
}
tracing::info!(hit = false, rows);
Ok(None)
}
}
}
#[allow(
clippy::too_many_lines,
clippy::too_many_arguments,
clippy::similar_names
)]
#[instrument(skip_all, err)]
pub fn write(
&self,
config: &StoredCacheConfig,
cache_kind: CacheKind,
cache_key: &[u8],
cache_val: &[u8],
dependencies: Option<SmallVec<[CacheDependency; 13]>>,
value_in_memory: bool,
) -> anyhow::Result<()> {
let mut cache_key_mut = BytesMut::new();
cache_key_mut.put_u8(cache_kind.as_u8());
cache_key_mut.put(cache_key);
let cache_key: Bytes = cache_key_mut.into();
let mut builder = flexbuffers::Builder::new(&flexbuffers::BuilderOptions::SHARE_NONE);
let mut builder_vec = builder.start_vector();
let mut internal_vec = builder_vec.start_vector();
let mut deps_vec = internal_vec.start_vector();
if config.evict_on_dependency_change == Some(true)
&& let Some(dependencies) = dependencies
{
let mut dep_lock = self.dependency_map.lock();
for dep in dependencies {
let mut dep_vec = deps_vec.start_vector();
match dep {
CacheDependency::DatabaseModel(uuid) => {
dep_vec.push(0u8);
dep_vec.push(flexbuffers::Blob(uuid.as_ref()));
}
}
dep_vec.end_vector();
if let Some(inverse_set) = dep_lock.get_mut(&dep) {
inverse_set.insert(cache_key.clone());
} else {
let mut inverse_set = HashSet::new();
inverse_set.insert(cache_key.clone());
(*dep_lock).insert(dep, inverse_set);
}
}
}
deps_vec.end_vector();
if let StoredCachePolicy::QuickCache = config.policy {
internal_vec.push(true);
internal_vec.push(value_in_memory);
} else {
internal_vec.push(false);
}
internal_vec.end_vector();
builder_vec.push(flexbuffers::Blob(cache_val));
builder_vec.end_vector();
let val = builder.view();
let size = (val.len() + cache_key.len()) as u64;
let mut lock = self.inventory.lock();
if let Some(max_size) = config.max_size
&& size + lock.0 > max_size
{
tracing::warn!("item causes 'max_size' to be exceeded");
return Ok(());
}
if let Some(max_count) = config.max_count
&& 1 + lock.1 > max_count
{
tracing::warn!("item causes 'max_count' to be exceeded");
return Ok(());
}
lock.0 += size;
lock.1 += 1;
let storage_size = lock.0;
let entries = lock.1;
drop(lock);
let txn = WriteTransaction::new(self.env.clone())?;
{
let mut access = txn.access();
access.put::<[u8], [u8]>(
&self.cache_db,
cache_key.as_ref(),
val,
&put::Flags::empty(),
)?;
}
txn.commit()?;
if let StoredCachePolicy::QuickCache = config.policy {
self.quick_cache.insert(
cache_key.clone(),
value_in_memory.then_some(Bytes::copy_from_slice(cache_val)),
);
}
tracing::info!(
total.entries = entries,
total.size = self.log_sizes.then_some(display(
bytesize::ByteSize(storage_size).display().si_short()
)),
item.size = self
.log_sizes
.then_some(display(bytesize::ByteSize(size).display().si_short())),
"stored"
);
Ok(())
}
#[instrument(skip_all, err)]
pub async fn dependency_evict(&self, dependencies: Vec<CacheDependency>) -> anyhow::Result<()> {
{
let txn = WriteTransaction::new(self.env.clone())?;
let mut lock_dep_map = self.dependency_map.lock();
let mut lock_inventory = self.inventory.lock();
{
let mut access = txn.access();
for dependency in dependencies {
if let Some(addrs) = lock_dep_map.get(&dependency) {
for cache_key in addrs {
{
tracing::debug!("evicting for dependency");
if self.quick_cache.remove(cache_key).is_none() {
match access
.get::<[u8], [u8]>(&self.cache_db, cache_key.as_ref())
{
Ok(val) => {
lock_inventory.0 -= cache_key.len() as u64;
lock_inventory.0 -= val.len() as u64;
lock_inventory.1 -= 1;
}
Err(err) => {
tracing::warn!(%err);
}
}
if let Err(err) =
access.del_key(&self.cache_db, cache_key.as_ref())
{
tracing::warn!(%err);
}
}
}
}
}
lock_dep_map.remove(&dependency);
}
}
txn.commit()?;
}
Ok(())
}
#[instrument(skip_all, err)]
pub async fn artifact_evict(
&self,
config: &StoredCache,
kind: CacheKind,
idx: u8,
) -> anyhow::Result<()> {
let mut key = BytesMut::new();
key.put_u8(kind.as_u8());
key.put_u8(idx);
let mut lock_inventory = self.inventory.lock();
let txn = WriteTransaction::new(self.env.clone())?;
{
let mut access = txn.access();
let mut cursor = txn.cursor(self.cache_db.clone())?;
let mut keys = vec![];
if let Ok((key, val)) = cursor.seek_range_k::<[u8], [u8]>(&access, key.as_ref()) {
lock_inventory.0 -= key.len() as u64;
lock_inventory.0 -= val.len() as u64;
lock_inventory.1 -= 1;
keys.push(Bytes::copy_from_slice(key));
}
while let Ok((key, val)) = cursor.next::<[u8], [u8]>(&access) {
lock_inventory.0 -= key.len() as u64;
lock_inventory.0 -= val.len() as u64;
lock_inventory.1 -= 1;
keys.push(Bytes::copy_from_slice(key));
}
for key in keys {
match config.policy {
StoredCachePolicy::Permanent => {
if let Err(err) = access.del_key::<[u8]>(&self.cache_db, key.as_ref()) {
tracing::error!(%err);
}
}
StoredCachePolicy::QuickCache => {
self.quick_cache.remove(&key);
}
}
}
}
txn.commit()?;
Ok(())
}
}