use crate::backend::interface::{
AtomicCacheWriter, BackendKind, CacheConnector, CacheReader, CacheWriter,
};
use crate::backend::memory::redis::RedisBackend;
use crate::backend::score::{BackendScore, Scores};
use crate::error::OxCacheResult;
use async_trait::async_trait;
use std::collections::HashSet;
use std::sync::Arc;
use std::time::Duration;
pub struct DragonflyBackend {
inner: RedisBackend,
restrictions: DragonflyRestrictions,
}
#[derive(Debug, Clone)]
pub struct DragonflyRestrictions {
disabled_commands: HashSet<String>,
cluster_disabled: bool,
}
impl Default for DragonflyRestrictions {
fn default() -> Self {
Self {
disabled_commands: ["FLUSHALL", "FLUSHDB", "DEBUG", "MONITOR"]
.into_iter()
.map(String::from)
.collect(),
cluster_disabled: true,
}
}
}
impl DragonflyRestrictions {
pub fn with_disabled_commands(mut self, commands: Vec<String>) -> Self {
self.disabled_commands = commands.into_iter().collect();
self
}
pub fn with_cluster_disabled(mut self, disabled: bool) -> Self {
self.cluster_disabled = disabled;
self
}
pub fn is_command_disabled(&self, command: &str) -> bool {
self.disabled_commands.contains(command)
}
pub fn cluster_disabled(&self) -> bool {
self.cluster_disabled
}
}
impl DragonflyBackend {
pub async fn new(url: &str, pool_size: usize) -> OxCacheResult<Self> {
let inner = RedisBackend::builder()
.connection_string(url)
.pool_size(pool_size)
.build()
.await?;
Ok(Self {
inner,
restrictions: DragonflyRestrictions::default(),
})
}
pub fn with_restrictions(mut self, restrictions: DragonflyRestrictions) -> Self {
self.restrictions = restrictions;
self
}
pub fn inner(&self) -> &RedisBackend {
&self.inner
}
pub fn restrictions(&self) -> &DragonflyRestrictions {
&self.restrictions
}
}
#[async_trait]
impl CacheReader for DragonflyBackend {
async fn get(&self, key: &str) -> OxCacheResult<Option<Vec<u8>>> {
self.inner.get(key).await
}
async fn exists(&self, key: &str) -> OxCacheResult<bool> {
self.inner.exists(key).await
}
async fn ttl(&self, key: &str) -> OxCacheResult<Option<Duration>> {
self.inner.ttl(key).await
}
async fn len(&self) -> OxCacheResult<u64> {
self.inner.len().await
}
async fn capacity(&self) -> OxCacheResult<u64> {
self.inner.capacity().await
}
async fn stats(&self) -> OxCacheResult<std::collections::HashMap<String, String>> {
self.inner.stats().await
}
async fn keys(&self, pattern: &str) -> OxCacheResult<Vec<String>> {
self.inner.keys(pattern).await
}
}
#[async_trait]
impl CacheWriter for DragonflyBackend {
async fn set(
&self,
key: Arc<str>,
value: Arc<Vec<u8>>,
ttl: Option<Duration>,
) -> OxCacheResult<()> {
self.inner.set(key, value, ttl).await
}
async fn delete(&self, key: &str) -> OxCacheResult<()> {
self.inner.delete(key).await
}
async fn clear(&self) -> OxCacheResult<()> {
self.inner.clear().await
}
async fn expire(&self, key: &str, ttl: Duration) -> OxCacheResult<bool> {
self.inner.expire(key, ttl).await
}
async fn set_many(
&self,
items: &[(Arc<str>, Arc<Vec<u8>>, Option<Duration>)],
) -> OxCacheResult<()> {
self.inner.set_many(items).await
}
async fn delete_many(&self, keys: &[String]) -> OxCacheResult<()> {
self.inner.delete_many(keys).await
}
}
impl BackendScore for DragonflyBackend {
fn score(&self) -> u8 {
Scores::REDIS
}
fn is_persistent(&self) -> bool {
true
}
fn backend_name(&self) -> &'static str {
"dragonfly"
}
}
#[async_trait]
impl CacheConnector for DragonflyBackend {
async fn health_check(&self) -> OxCacheResult<()> {
self.inner.health_check().await
}
async fn shutdown(&self) {
self.inner.shutdown().await
}
fn backend_kind(&self) -> BackendKind {
BackendKind::Dragonfly
}
fn as_atomic_writer(&self) -> Option<&dyn AtomicCacheWriter> {
None
}
}
#[cfg(test)]
mod tests;