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>,
);
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),
}
}
}
pub trait SortedSetProvider {
fn zadd(&self, key: &str, score: f64, member: &str) -> Result<u32, &'static str>;
fn zremrangebyscore(&self, key: &str, min: f64, max: f64) -> Result<u32, &'static str>;
fn zcard(&self, key: &str) -> Result<u32, &'static str>;
fn zrangebyscore_withscores(
&self,
key: &str,
min: f64,
max: f64,
limit: Option<usize>,
) -> Result<Vec<(String, f64)>, &'static str>;
fn expire(&self, key: &str, ttl_secs: u64) -> Result<(), &'static str>;
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>)
}
}
}
#[cfg(feature = "bullmq")]
pub async fn bullmq(
self,
) -> Result<
Box<dyn crate::types::bullmq::BullMqProvider + Send + Sync>,
crate::types::queues::BullMqError,
> {
match self {
ProviderEither::Redis(provider) => {
let bull = crate::providers::bullmq::BullMq::connect(provider.url.clone()).await?;
Ok(Box::new(bull))
}
}
}
}
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");
}
}
}
}
}