flippico-cache 0.4.1

Flippico cache adapter
Documentation
use redis::Client;

use crate::{
    providers::redis::Redis,
    types::channels::{
        CacheSpace, CacheValue, ChannelMessage, ListChannel, ListMessage, MessageMeta,
        SubscriptionChannel,
    },
};
use dotenv::dotenv;
use std::env;
use std::future::Future;
use std::pin::Pin;

pub trait ProviderTrait {
    fn new(url: String) -> Self
    where
        Self: Sized;
}

pub trait Connectable {
    fn set_client(&mut self);
    fn get_connection(&self) -> &Client;
    fn get_mut_connection(&mut self) -> &mut Client;
}

pub type PubSubProviderError = &'static str;

pub type AsyncCallback<T> =
    Box<dyn Fn(T) -> Pin<Box<dyn Future<Output = ()> + Send + Sync>> + Send + Sync>;

pub trait PubSubProvider {
    type Channels;
    fn subscribe(
        &self,
        callback: AsyncCallback<ChannelMessage>,
        channel: Self::Channels,
    ) -> Pin<Box<dyn Future<Output = Result<(), &'static str>> + Send + '_>>;
    fn publish(&mut self, channel: Self::Channels, message: ChannelMessage);
}

pub trait FifoProvider {
    fn push(&self, channel_name: ListChannel, job_key: &str, job_payload: serde_json::Value);
    fn pop(&self, channel_name: ListChannel, job_key: &str) -> Result<ListMessage, &'static str>;

    fn get_list_name(&self, channel_name: &str, job_key: &str) -> String {
        format!("{}:{}", channel_name, job_key)
    }
}

pub trait CacheProvider {
    fn get(&self, cache_space: CacheSpace, key: &str, app_name: Option<String>) -> String;
    fn delete(&self, cache_space: CacheSpace, key: &str, app_name: Option<String>);
    fn set(
        &self,
        cache_space: CacheSpace,
        key: &str,
        value: serde_json::Value,
        app_name: Option<String>,
    );
    /// Set a key with an explicit TTL in seconds. Defaults to `set` (no expiry) if not overridden.
    fn set_with_ttl(
        &self,
        cache_space: CacheSpace,
        key: &str,
        value: serde_json::Value,
        app_name: Option<String>,
        ttl_secs: u64,
    ) {
        self.set(cache_space, key, value, app_name);
        let _ = ttl_secs;
    }

    fn get_key_name(&self, cache_space: CacheSpace, key: &str, app_name: Option<String>) -> String {
        if app_name.is_none() {
            return format!("{}:{}", cache_space.get_space(), key);
        }

        format!("{}:{}:{}", cache_space.get_space(), app_name.unwrap(), key)
    }

    fn get_value(
        &self,
        app_name: Option<String>,
        value: serde_json::Value,
    ) -> CacheValue<serde_json::Value> {
        if app_name.is_none() {
            return CacheValue {
                meta: None,
                value: Some(value),
            };
        }

        CacheValue {
            meta: Some(MessageMeta {
                app_name: app_name.unwrap(),
            }),
            value: Some(value),
        }
    }
}

/// Sorted set operations for rate limiting, leaderboards, etc.
pub trait SortedSetProvider {
    /// Add a member with the given score. Returns the number of new elements added.
    fn zadd(&self, key: &str, score: f64, member: &str) -> Result<u32, &'static str>;

    /// Remove all members with scores between min and max (inclusive).
    /// Returns the number of removed members.
    fn zremrangebyscore(&self, key: &str, min: f64, max: f64) -> Result<u32, &'static str>;

    /// Return the number of members in the sorted set.
    fn zcard(&self, key: &str) -> Result<u32, &'static str>;

    /// Return members with scores in the given range, ordered by score ascending.
    /// Includes scores in the result as `(member, score)` pairs.
    fn zrangebyscore_withscores(
        &self,
        key: &str,
        min: f64,
        max: f64,
        limit: Option<usize>,
    ) -> Result<Vec<(String, f64)>, &'static str>;

    /// Set a TTL (in seconds) on the key. Useful for auto-cleanup of sliding windows.
    fn expire(&self, key: &str, ttl_secs: u64) -> Result<(), &'static str>;

    /// Execute a Lua script with the given keys and args. Returns a vector of integers.
    fn eval_script(
        &self,
        script: &str,
        keys: &[&str],
        args: &[&str],
    ) -> Result<Vec<i64>, &'static str>;
}

pub trait ConnectableProvider: ProviderTrait + Connectable {}
impl<T> ConnectableProvider for T where T: ProviderTrait + Connectable {}

pub trait ConnectablePubSubProvider: ConnectableProvider + PubSubProvider {}
impl<T> ConnectablePubSubProvider for T where T: ConnectableProvider + PubSubProvider {}

pub trait ConnectableFifoProvider: ConnectableProvider + FifoProvider {}
impl<T> ConnectableFifoProvider for T where T: ConnectableProvider + FifoProvider {}

pub trait ConnectableCacheProvider: ConnectableProvider + CacheProvider {}
impl<T> ConnectableCacheProvider for T where T: ConnectableProvider + CacheProvider {}

pub trait ConnectableSortedSetProvider: ConnectableProvider + SortedSetProvider {}
impl<T> ConnectableSortedSetProvider for T where T: ConnectableProvider + SortedSetProvider {}

pub enum Provider {
    Redis,
}

pub enum ProviderEither {
    Redis(Box<Redis>),
}

impl ProviderEither {
    pub fn pub_sub(
        self,
    ) -> Result<
        Box<dyn ConnectablePubSubProvider<Channels = SubscriptionChannel> + Send + Sync>,
        &'static str,
    > {
        match self {
            ProviderEither::Redis(provider) => Ok(provider
                as Box<
                    dyn ConnectablePubSubProvider<Channels = SubscriptionChannel> + Send + Sync,
                >),
        }
    }

    pub fn connectable(self) -> Result<Box<dyn ConnectableProvider + Send + Sync>, &'static str> {
        match self {
            ProviderEither::Redis(provider) => {
                Ok(provider as Box<dyn ConnectableProvider + Send + Sync>)
            }
        }
    }

    pub fn fifo(self) -> Result<Box<dyn ConnectableFifoProvider + Send + Sync>, &'static str> {
        match self {
            ProviderEither::Redis(provider) => {
                Ok(provider as Box<dyn ConnectableFifoProvider + Send + Sync>)
            }
        }
    }

    pub fn cache(self) -> Result<Box<dyn ConnectableCacheProvider + Send + Sync>, &'static str> {
        match self {
            ProviderEither::Redis(provider) => {
                Ok(provider as Box<dyn ConnectableCacheProvider + Send + Sync>)
            }
        }
    }

    pub fn sorted_set(
        self,
    ) -> Result<Box<dyn ConnectableSortedSetProvider + Send + Sync>, &'static str> {
        match self {
            ProviderEither::Redis(provider) => {
                Ok(provider as Box<dyn ConnectableSortedSetProvider + Send + Sync>)
            }
        }
    }
}

impl Provider {
    pub(crate) fn create_provider(&self) -> ProviderEither {
        match self {
            Provider::Redis => {
                dotenv().ok();
                let url = env::var("FLIPPICO_CACHE_REDIS_URL");
                if let Ok(url) = url {
                    let mut provider = Redis::new(url);
                    provider.set_client();
                    ProviderEither::Redis(Box::new(provider))
                } else {
                    log::error!("FLIPPICO_CACHE_REDIS_URL is not set");
                    panic!("FLIPPICO_CACHE_REDIS_URL is not set");
                }
            }
        }
    }
}