use crate::error::{OxCacheError, OxCacheResult};
#[cfg(feature = "telemetry")]
use crate::i18n::messages::MSG_LOG_BRIDGE_SHUTDOWN_SKIPPED;
use crate::i18n::messages::{
MSG_DETAIL_ASYNC_ATOMIC_CAS, MSG_DETAIL_ASYNC_ATOMIC_INCREMENT,
MSG_DETAIL_ASYNC_ATOMIC_SET_IF_ABSENT, MSG_DETAIL_SYNC_ATOMIC_CAS,
MSG_DETAIL_SYNC_ATOMIC_INCREMENT, MSG_DETAIL_SYNC_ATOMIC_SET_IF_ABSENT,
MSG_DETAIL_SYNC_REQUIRES_MULTI_THREAD, MSG_DETAIL_SYNC_REQUIRES_RUNTIME,
MSG_PANIC_BRIDGE_TEMP_RUNTIME, t,
};
use async_trait::async_trait;
use std::sync::Arc;
use std::time::Duration;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum BackendKind {
Moka,
DashMap,
Redis,
Valkey,
Dragonfly,
Aerospike,
Chain,
Mock,
Disk,
Unknown,
}
impl BackendKind {
pub fn is_memory(&self) -> bool {
matches!(
self,
BackendKind::Moka | BackendKind::DashMap | BackendKind::Mock
)
}
pub fn is_distributed(&self) -> bool {
matches!(
self,
BackendKind::Redis
| BackendKind::Valkey
| BackendKind::Dragonfly
| BackendKind::Aerospike
)
}
pub fn is_composite(&self) -> bool {
matches!(self, BackendKind::Chain)
}
pub fn name(&self) -> &'static str {
match self {
BackendKind::Moka => "moka",
BackendKind::DashMap => "dashmap",
BackendKind::Redis => "redis",
BackendKind::Valkey => "valkey",
BackendKind::Dragonfly => "dragonfly",
BackendKind::Aerospike => "aerospike",
BackendKind::Chain => "chain",
BackendKind::Mock => "mock",
BackendKind::Disk => "disk",
BackendKind::Unknown => "unknown",
}
}
}
pub fn glob_match(pattern: &str, text: &str) -> bool {
let p: Vec<char> = pattern.chars().collect();
let t: Vec<char> = text.chars().collect();
let (m, n) = (p.len(), t.len());
let mut pi = 0;
let mut ti = 0;
let mut star_pi: Option<usize> = None;
let mut star_ti: usize = 0;
while ti < n {
if pi < m && p[pi] == '?' {
pi += 1;
ti += 1;
} else if pi < m && p[pi] == '*' {
star_pi = Some(pi);
star_ti = ti;
pi += 1;
} else if pi < m && p[pi] == t[ti] {
pi += 1;
ti += 1;
} else if let Some(sp) = star_pi {
pi = sp + 1;
star_ti += 1;
ti = star_ti;
} else {
return false;
}
}
while pi < m && p[pi] == '*' {
pi += 1;
}
pi == m
}
#[async_trait]
pub trait CacheReader: Send + Sync + 'static {
async fn get(&self, key: &str) -> OxCacheResult<Option<Vec<u8>>>;
async fn exists(&self, key: &str) -> OxCacheResult<bool>;
async fn ttl(&self, key: &str) -> OxCacheResult<Option<Duration>>;
async fn len(&self) -> OxCacheResult<u64>;
async fn is_empty(&self) -> OxCacheResult<bool> {
Ok(self.len().await?.eq(&0))
}
async fn capacity(&self) -> OxCacheResult<u64>;
async fn stats(&self) -> OxCacheResult<std::collections::HashMap<String, String>>;
async fn get_many(&self, keys: &[String]) -> OxCacheResult<Vec<Option<Vec<u8>>>> {
let mut results = Vec::with_capacity(keys.len());
for key in keys {
results.push(self.get(key).await?);
}
Ok(results)
}
async fn keys(&self, pattern: &str) -> OxCacheResult<Vec<String>> {
let _ = pattern;
Ok(vec![])
}
}
pub type CacheSetItem = (Arc<str>, Arc<Vec<u8>>, Option<Duration>);
#[async_trait]
pub trait CacheWriter: Send + Sync + 'static {
async fn set(
&self,
key: Arc<str>,
value: Arc<Vec<u8>>,
ttl: Option<Duration>,
) -> OxCacheResult<()>;
async fn delete(&self, key: &str) -> OxCacheResult<()>;
async fn clear(&self) -> OxCacheResult<()>;
async fn expire(&self, key: &str, ttl: Duration) -> OxCacheResult<bool>;
async fn set_many(&self, items: &[CacheSetItem]) -> OxCacheResult<()> {
for (key, value, ttl) in items {
self.set(key.clone(), value.clone(), *ttl).await?;
}
Ok(())
}
async fn delete_many(&self, keys: &[String]) -> OxCacheResult<()> {
for key in keys {
self.delete(key).await?;
}
Ok(())
}
}
#[async_trait]
pub trait CacheConnector: Send + Sync + 'static {
async fn health_check(&self) -> OxCacheResult<()>;
async fn shutdown(&self);
fn backend_kind(&self) -> BackendKind;
#[cfg(feature = "lua")]
fn as_lua_executor(&self) -> Option<&dyn LuaExecutor> {
None
}
fn as_atomic_writer(&self) -> Option<&dyn AtomicCacheWriter> {
None
}
}
#[cfg(feature = "lua")]
#[async_trait]
pub trait LuaExecutor: Send + Sync {
async fn eval_lua(
&self,
script: &str,
keys: &[&str],
args: &[&str],
) -> OxCacheResult<redis::Value>;
async fn eval_sha(
&self,
sha: &str,
keys: &[&str],
args: &[&str],
) -> OxCacheResult<redis::Value>;
async fn script_load(&self, script: &str) -> OxCacheResult<String>;
}
#[async_trait]
pub trait AtomicCacheWriter: Send + Sync + 'static {
async fn incr(&self, key: &str, delta: i64, ttl: Option<Duration>) -> OxCacheResult<i64>;
async fn compare_and_swap(
&self,
key: &str,
expected: Option<&[u8]>,
new: Vec<u8>,
ttl: Option<Duration>,
) -> OxCacheResult<bool>;
async fn set_if_absent(
&self,
key: &str,
value: Vec<u8>,
ttl: Option<Duration>,
) -> OxCacheResult<bool>;
}
#[async_trait]
pub trait CacheBackend: CacheReader + CacheWriter + CacheConnector + 'static {}
#[async_trait]
impl<T: CacheReader + CacheWriter + CacheConnector + 'static> CacheBackend for T {}
pub trait SyncCacheReader: Send + Sync + 'static {
fn get(&self, key: &str) -> OxCacheResult<Option<Vec<u8>>>;
fn exists(&self, key: &str) -> OxCacheResult<bool>;
fn ttl(&self, key: &str) -> OxCacheResult<Option<Duration>>;
fn len(&self) -> OxCacheResult<u64>;
fn is_empty(&self) -> OxCacheResult<bool> {
Ok(self.len()? == 0)
}
fn capacity(&self) -> OxCacheResult<u64>;
fn stats(&self) -> OxCacheResult<std::collections::HashMap<String, String>>;
fn get_many(&self, keys: &[String]) -> OxCacheResult<Vec<Option<Vec<u8>>>> {
let mut results = Vec::with_capacity(keys.len());
for key in keys {
results.push(self.get(key)?);
}
Ok(results)
}
fn keys(&self, pattern: &str) -> OxCacheResult<Vec<String>> {
let _ = pattern;
Ok(vec![])
}
}
pub trait SyncCacheWriter: Send + Sync + 'static {
fn set(&self, key: Arc<str>, value: Arc<Vec<u8>>, ttl: Option<Duration>) -> OxCacheResult<()>;
fn delete(&self, key: &str) -> OxCacheResult<()>;
fn clear(&self) -> OxCacheResult<()>;
fn expire(&self, key: &str, ttl: Duration) -> OxCacheResult<bool>;
fn set_many(&self, items: &[CacheSetItem]) -> OxCacheResult<()> {
for (key, value, ttl) in items {
self.set(key.clone(), value.clone(), *ttl)?;
}
Ok(())
}
fn delete_many(&self, keys: &[String]) -> OxCacheResult<()> {
for key in keys {
self.delete(key)?;
}
Ok(())
}
}
pub trait SyncCacheConnector: Send + Sync + 'static {
fn health_check(&self) -> OxCacheResult<()>;
fn shutdown(&self);
fn backend_kind(&self) -> BackendKind;
fn as_sync_atomic_writer(&self) -> Option<&dyn SyncAtomicCacheWriter> {
None
}
}
pub trait SyncCacheBackend:
SyncCacheReader + SyncCacheWriter + SyncCacheConnector + 'static
{
}
impl<T: SyncCacheReader + SyncCacheWriter + SyncCacheConnector + 'static> SyncCacheBackend for T {}
pub trait SyncAtomicCacheWriter: Send + Sync + 'static {
fn incr(&self, key: &str, delta: i64, ttl: Option<Duration>) -> OxCacheResult<i64>;
fn compare_and_swap(
&self,
key: &str,
expected: Option<&[u8]>,
new: Vec<u8>,
ttl: Option<Duration>,
) -> OxCacheResult<bool>;
fn set_if_absent(
&self,
key: &str,
value: Vec<u8>,
ttl: Option<Duration>,
) -> OxCacheResult<bool>;
}
pub struct SyncBackendAdapter {
inner: Arc<dyn SyncCacheBackend>,
}
impl SyncBackendAdapter {
pub fn new(inner: Arc<dyn SyncCacheBackend>) -> Self {
Self { inner }
}
}
impl SyncCacheReader for SyncBackendAdapter {
fn get(&self, key: &str) -> OxCacheResult<Option<Vec<u8>>> {
self.inner.get(key)
}
fn exists(&self, key: &str) -> OxCacheResult<bool> {
self.inner.exists(key)
}
fn ttl(&self, key: &str) -> OxCacheResult<Option<Duration>> {
self.inner.ttl(key)
}
fn len(&self) -> OxCacheResult<u64> {
self.inner.len()
}
fn capacity(&self) -> OxCacheResult<u64> {
self.inner.capacity()
}
fn stats(&self) -> OxCacheResult<std::collections::HashMap<String, String>> {
self.inner.stats()
}
fn get_many(&self, keys: &[String]) -> OxCacheResult<Vec<Option<Vec<u8>>>> {
self.inner.get_many(keys)
}
fn keys(&self, pattern: &str) -> OxCacheResult<Vec<String>> {
self.inner.keys(pattern)
}
}
impl SyncCacheWriter for SyncBackendAdapter {
fn set(&self, key: Arc<str>, value: Arc<Vec<u8>>, ttl: Option<Duration>) -> OxCacheResult<()> {
self.inner.set(key, value, ttl)
}
fn delete(&self, key: &str) -> OxCacheResult<()> {
self.inner.delete(key)
}
fn clear(&self) -> OxCacheResult<()> {
self.inner.clear()
}
fn expire(&self, key: &str, ttl: Duration) -> OxCacheResult<bool> {
self.inner.expire(key, ttl)
}
fn set_many(&self, items: &[CacheSetItem]) -> OxCacheResult<()> {
self.inner.set_many(items)
}
fn delete_many(&self, keys: &[String]) -> OxCacheResult<()> {
self.inner.delete_many(keys)
}
}
impl SyncCacheConnector for SyncBackendAdapter {
fn health_check(&self) -> OxCacheResult<()> {
self.inner.health_check()
}
fn shutdown(&self) {
self.inner.shutdown()
}
fn backend_kind(&self) -> BackendKind {
self.inner.backend_kind()
}
}
#[async_trait]
impl CacheReader for SyncBackendAdapter {
async fn get(&self, key: &str) -> OxCacheResult<Option<Vec<u8>>> {
self.inner.get(key)
}
async fn exists(&self, key: &str) -> OxCacheResult<bool> {
self.inner.exists(key)
}
async fn ttl(&self, key: &str) -> OxCacheResult<Option<Duration>> {
self.inner.ttl(key)
}
async fn len(&self) -> OxCacheResult<u64> {
self.inner.len()
}
async fn capacity(&self) -> OxCacheResult<u64> {
self.inner.capacity()
}
async fn stats(&self) -> OxCacheResult<std::collections::HashMap<String, String>> {
self.inner.stats()
}
async fn get_many(&self, keys: &[String]) -> OxCacheResult<Vec<Option<Vec<u8>>>> {
self.inner.get_many(keys)
}
async fn keys(&self, pattern: &str) -> OxCacheResult<Vec<String>> {
self.inner.keys(pattern)
}
}
#[async_trait]
impl CacheWriter for SyncBackendAdapter {
async fn set(
&self,
key: Arc<str>,
value: Arc<Vec<u8>>,
ttl: Option<Duration>,
) -> OxCacheResult<()> {
self.inner.set(key, value, ttl)
}
async fn delete(&self, key: &str) -> OxCacheResult<()> {
self.inner.delete(key)
}
async fn clear(&self) -> OxCacheResult<()> {
self.inner.clear()
}
async fn expire(&self, key: &str, ttl: Duration) -> OxCacheResult<bool> {
self.inner.expire(key, ttl)
}
async fn set_many(&self, items: &[CacheSetItem]) -> OxCacheResult<()> {
self.inner.set_many(items)
}
async fn delete_many(&self, keys: &[String]) -> OxCacheResult<()> {
self.inner.delete_many(keys)
}
}
#[async_trait]
impl CacheConnector for SyncBackendAdapter {
async fn health_check(&self) -> OxCacheResult<()> {
self.inner.health_check()
}
async fn shutdown(&self) {
self.inner.shutdown()
}
fn backend_kind(&self) -> BackendKind {
self.inner.backend_kind()
}
fn as_atomic_writer(&self) -> Option<&dyn AtomicCacheWriter> {
match self.inner.as_sync_atomic_writer() {
Some(_) => Some(self),
None => None,
}
}
}
#[async_trait]
impl AtomicCacheWriter for SyncBackendAdapter {
async fn incr(&self, key: &str, delta: i64, ttl: Option<Duration>) -> OxCacheResult<i64> {
match self.inner.as_sync_atomic_writer() {
Some(w) => w.incr(key, delta, ttl),
None => Err(OxCacheError::NotSupported(t(
MSG_DETAIL_SYNC_ATOMIC_INCREMENT,
&[],
))),
}
}
async fn compare_and_swap(
&self,
key: &str,
expected: Option<&[u8]>,
new: Vec<u8>,
ttl: Option<Duration>,
) -> OxCacheResult<bool> {
match self.inner.as_sync_atomic_writer() {
Some(w) => w.compare_and_swap(key, expected, new, ttl),
None => Err(OxCacheError::NotSupported(t(
MSG_DETAIL_SYNC_ATOMIC_CAS,
&[],
))),
}
}
async fn set_if_absent(
&self,
key: &str,
value: Vec<u8>,
ttl: Option<Duration>,
) -> OxCacheResult<bool> {
match self.inner.as_sync_atomic_writer() {
Some(w) => w.set_if_absent(key, value, ttl),
None => Err(OxCacheError::NotSupported(t(
MSG_DETAIL_SYNC_ATOMIC_SET_IF_ABSENT,
&[],
))),
}
}
}
pub(crate) fn multi_thread_bridge_handle() -> OxCacheResult<tokio::runtime::Handle> {
let handle = tokio::runtime::Handle::try_current().map_err(|e| {
OxCacheError::NotSupported(t(
MSG_DETAIL_SYNC_REQUIRES_RUNTIME,
&[("err", e.to_string())],
))
})?;
if handle.runtime_flavor() == tokio::runtime::RuntimeFlavor::CurrentThread {
return Err(OxCacheError::NotSupported(t(
MSG_DETAIL_SYNC_REQUIRES_MULTI_THREAD,
&[],
)));
}
Ok(handle)
}
pub struct AsyncToSyncBridge {
inner: Arc<dyn CacheBackend>,
}
impl AsyncToSyncBridge {
pub fn new(inner: Arc<dyn CacheBackend>) -> Self {
Self { inner }
}
}
impl SyncCacheReader for AsyncToSyncBridge {
fn get(&self, key: &str) -> OxCacheResult<Option<Vec<u8>>> {
let handle = multi_thread_bridge_handle()?;
tokio::task::block_in_place(|| handle.block_on(CacheReader::get(&*self.inner, key)))
}
fn exists(&self, key: &str) -> OxCacheResult<bool> {
let handle = multi_thread_bridge_handle()?;
tokio::task::block_in_place(|| handle.block_on(CacheReader::exists(&*self.inner, key)))
}
fn ttl(&self, key: &str) -> OxCacheResult<Option<Duration>> {
let handle = multi_thread_bridge_handle()?;
tokio::task::block_in_place(|| handle.block_on(CacheReader::ttl(&*self.inner, key)))
}
fn len(&self) -> OxCacheResult<u64> {
let handle = multi_thread_bridge_handle()?;
tokio::task::block_in_place(|| handle.block_on(CacheReader::len(&*self.inner)))
}
fn capacity(&self) -> OxCacheResult<u64> {
let handle = multi_thread_bridge_handle()?;
tokio::task::block_in_place(|| handle.block_on(CacheReader::capacity(&*self.inner)))
}
fn stats(&self) -> OxCacheResult<std::collections::HashMap<String, String>> {
let handle = multi_thread_bridge_handle()?;
tokio::task::block_in_place(|| handle.block_on(CacheReader::stats(&*self.inner)))
}
fn get_many(&self, keys: &[String]) -> OxCacheResult<Vec<Option<Vec<u8>>>> {
let handle = multi_thread_bridge_handle()?;
tokio::task::block_in_place(|| handle.block_on(CacheReader::get_many(&*self.inner, keys)))
}
fn keys(&self, pattern: &str) -> OxCacheResult<Vec<String>> {
let handle = multi_thread_bridge_handle()?;
tokio::task::block_in_place(|| handle.block_on(CacheReader::keys(&*self.inner, pattern)))
}
}
impl SyncCacheWriter for AsyncToSyncBridge {
fn set(&self, key: Arc<str>, value: Arc<Vec<u8>>, ttl: Option<Duration>) -> OxCacheResult<()> {
let handle = multi_thread_bridge_handle()?;
tokio::task::block_in_place(|| {
handle.block_on(CacheWriter::set(&*self.inner, key, value, ttl))
})
}
fn delete(&self, key: &str) -> OxCacheResult<()> {
let handle = multi_thread_bridge_handle()?;
tokio::task::block_in_place(|| handle.block_on(CacheWriter::delete(&*self.inner, key)))
}
fn clear(&self) -> OxCacheResult<()> {
let handle = multi_thread_bridge_handle()?;
tokio::task::block_in_place(|| handle.block_on(CacheWriter::clear(&*self.inner)))
}
fn expire(&self, key: &str, ttl: Duration) -> OxCacheResult<bool> {
let handle = multi_thread_bridge_handle()?;
tokio::task::block_in_place(|| handle.block_on(CacheWriter::expire(&*self.inner, key, ttl)))
}
fn set_many(&self, items: &[CacheSetItem]) -> OxCacheResult<()> {
let handle = multi_thread_bridge_handle()?;
tokio::task::block_in_place(|| handle.block_on(CacheWriter::set_many(&*self.inner, items)))
}
fn delete_many(&self, keys: &[String]) -> OxCacheResult<()> {
let handle = multi_thread_bridge_handle()?;
tokio::task::block_in_place(|| {
handle.block_on(CacheWriter::delete_many(&*self.inner, keys))
})
}
}
impl SyncCacheConnector for AsyncToSyncBridge {
fn health_check(&self) -> OxCacheResult<()> {
let handle = multi_thread_bridge_handle()?;
tokio::task::block_in_place(|| handle.block_on(CacheConnector::health_check(&*self.inner)))
}
fn shutdown(&self) {
match tokio::runtime::Handle::try_current() {
Ok(handle) if handle.runtime_flavor() == tokio::runtime::RuntimeFlavor::MultiThread => {
tokio::task::block_in_place(|| {
handle.block_on(CacheConnector::shutdown(&*self.inner))
});
}
Ok(_) => {
record_bridge_shutdown_rejected();
warn_bridge_shutdown_rejected();
}
Err(_) => {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap_or_else(|e| panic!("{}: {e:?}", t(MSG_PANIC_BRIDGE_TEMP_RUNTIME, &[])));
rt.block_on(CacheConnector::shutdown(&*self.inner));
}
}
}
fn backend_kind(&self) -> BackendKind {
CacheConnector::backend_kind(&*self.inner)
}
fn as_sync_atomic_writer(&self) -> Option<&dyn SyncAtomicCacheWriter> {
if self.inner.as_atomic_writer().is_some() {
Some(self)
} else {
None
}
}
}
impl SyncAtomicCacheWriter for AsyncToSyncBridge {
fn incr(&self, key: &str, delta: i64, ttl: Option<Duration>) -> OxCacheResult<i64> {
let handle = multi_thread_bridge_handle()?;
tokio::task::block_in_place(|| {
handle.block_on(match self.inner.as_atomic_writer() {
Some(w) => w.incr(key, delta, ttl),
None => {
return Err(OxCacheError::NotSupported(t(
MSG_DETAIL_ASYNC_ATOMIC_INCREMENT,
&[],
)));
}
})
})
}
fn compare_and_swap(
&self,
key: &str,
expected: Option<&[u8]>,
new: Vec<u8>,
ttl: Option<Duration>,
) -> OxCacheResult<bool> {
let handle = multi_thread_bridge_handle()?;
tokio::task::block_in_place(|| {
handle.block_on(match self.inner.as_atomic_writer() {
Some(w) => w.compare_and_swap(key, expected, new, ttl),
None => {
return Err(OxCacheError::NotSupported(t(
MSG_DETAIL_ASYNC_ATOMIC_CAS,
&[],
)));
}
})
})
}
fn set_if_absent(
&self,
key: &str,
value: Vec<u8>,
ttl: Option<Duration>,
) -> OxCacheResult<bool> {
let handle = multi_thread_bridge_handle()?;
tokio::task::block_in_place(|| {
handle.block_on(match self.inner.as_atomic_writer() {
Some(w) => w.set_if_absent(key, value, ttl),
None => {
return Err(OxCacheError::NotSupported(t(
MSG_DETAIL_ASYNC_ATOMIC_SET_IF_ABSENT,
&[],
)));
}
})
})
}
}
#[cfg(feature = "metrics")]
#[inline]
fn record_bridge_shutdown_rejected() {
crate::infra::metrics::unified::GLOBAL_UNIFIED_METRICS
.increment_counter("oxcache_bridge_shutdown_rejected_total", 1);
}
#[cfg(not(feature = "metrics"))]
#[inline]
fn record_bridge_shutdown_rejected() {}
#[cfg(feature = "telemetry")]
#[inline]
fn warn_bridge_shutdown_rejected() {
tracing::warn!(
target = "oxcache::backend",
"{}",
t(MSG_LOG_BRIDGE_SHUTDOWN_SKIPPED, &[])
);
}
#[cfg(not(feature = "telemetry"))]
#[inline]
fn warn_bridge_shutdown_rejected() {}
#[cfg(test)]
mod tests;