use super::Result;
use crate::redis::{collector::AsRedisPairs, RedisModel, RedisRead};
use deadpool_redis::{
redis::{AsyncCommands, Expiry, ToRedisArgs},
Config, Connection, Pool,
};
use futures::future::join_all;
use std::fmt::Debug;
#[derive(Debug, Clone)]
pub struct Client {
pool: Pool,
}
impl Client {
pub async fn default() -> Result<Self> {
Self::from_url("redis://127.0.0.1:6379").await
}
pub fn from_pool(pool: Pool) -> Self {
Self { pool }
}
pub async fn from_url(url: &str) -> Result<Self> {
let config = Config::from_url(url);
Self::connect(&config).await
}
pub async fn connect(config: &Config) -> Result<Self> {
let pool = config.create_pool(Some(deadpool_redis::Runtime::Tokio1))?;
Ok(Self { pool })
}
pub async fn connection(&self) -> Result<Connection> {
Ok(self.pool.get().await?)
}
}
impl Client {
pub async fn get<V, K>(&self, key: K) -> Result<Option<V>>
where
V: RedisRead,
K: for<'a> ToRedisArgs + Send + Sync,
{
let mut connection = self.connection().await?;
Ok(connection.get(key).await?)
}
pub async fn mget<K, T, V>(&self, keys: K) -> Result<Vec<Option<V>>>
where
V: RedisRead,
K: IntoIterator<Item = T> + ToRedisArgs + Send + Sync,
T: for<'a> ToRedisArgs + Send + Sync,
{
let mut connection = self.connection().await?;
Ok(connection.mget(keys).await?)
}
pub async fn get_ex<V, K>(&self, key: K, expire_at: Expiry) -> Result<Option<V>>
where
V: RedisRead,
K: for<'a> ToRedisArgs + Send + Sync,
{
let mut connection = self.connection().await?;
Ok(connection.get_ex(key, expire_at).await?)
}
pub async fn get_del<V, K>(&self, key: K) -> Result<Option<V>>
where
V: RedisRead,
K: for<'a> ToRedisArgs + Send + Sync,
{
let mut connection = self.connection().await?;
Ok(connection.get_del(key).await?)
}
pub async fn getset<M, V>(&self, model: &M) -> Result<Option<V>>
where
M: RedisModel,
V: RedisRead,
{
let mut connection = self.connection().await?;
Ok(connection.getset(model.key()?, model.value()?).await?)
}
}
impl Client {
pub async fn set<M>(&self, model: &M) -> Result<String>
where
M: RedisModel,
{
let mut connection = self.connection().await?;
Ok(connection.set(model.key()?, model.value()?).await?)
}
pub async fn mset<M, P>(&self, pairs: P) -> Result<String>
where
M: RedisModel,
P: AsRedisPairs<M> + Send + Sync,
{
let mut connection = self.connection().await?;
let pairs = pairs.as_pairs();
Ok(connection.mset(&pairs).await?)
}
pub async fn mset_nx<M, P>(&self, pairs: P) -> Result<bool>
where
M: RedisModel,
P: AsRedisPairs<M> + Send + Sync,
{
let mut connection = self.connection().await?;
let pairs = pairs.as_pairs();
Ok(connection.mset_nx(&pairs).await?)
}
pub async fn set_nx<M>(&self, model: &M) -> Result<bool>
where
M: RedisModel,
{
let mut connection = self.connection().await?;
Ok(connection.set_nx(model.key()?, model.value()?).await?)
}
pub async fn set_ex<M>(&self, model: &M, secs: u64) -> Result<String>
where
M: RedisModel,
{
let mut connection = self.connection().await?;
Ok(connection
.set_ex(model.key()?, model.value()?, secs)
.await?)
}
}
impl Client {
pub async fn del<K>(&self, key: K) -> Result<bool>
where
K: for<'a> ToRedisArgs + Send + Sync,
{
let mut connection = self.connection().await?;
Ok(connection.del(key).await?)
}
pub async fn mdel<K, T>(&self, keys: K) -> Result<usize>
where
K: IntoIterator<Item = T>,
T: for<'a> ToRedisArgs + Send + Sync,
{
let mut futures = vec![];
for key in keys {
futures.push(self.del(key));
}
let results = join_all(futures).await;
Ok(results
.iter()
.filter(|result| matches!(result, Ok(true)))
.count())
}
}
impl Client {
pub async fn exists<K>(&self, key: K) -> Result<bool>
where
K: for<'a> ToRedisArgs + Send + Sync,
{
let mut connection = self.connection().await?;
Ok(connection.exists(key).await?)
}
pub async fn ping(&self) -> Result<String> {
let mut connection = self.connection().await?;
Ok(connection.ping().await?)
}
pub async fn rename<K1, K2>(&self, key: K1, new_key: K2) -> Result<String>
where
K1: for<'a> ToRedisArgs + Send + Sync,
K2: for<'a> ToRedisArgs + Send + Sync,
{
let mut connection = self.connection().await?;
Ok(connection.rename(key, new_key).await?)
}
pub async fn rename_nx<K1, K2>(&self, key: K1, new_key: K2) -> Result<bool>
where
K1: for<'a> ToRedisArgs + Send + Sync,
K2: for<'a> ToRedisArgs + Send + Sync,
{
let mut connection = self.connection().await?;
Ok(connection.rename_nx(key, new_key).await?)
}
}
#[cfg(test)]
mod tests {
type Result<T> = super::Result<T>;
use std::time::Duration;
use crate::redis;
use crate::redis::macros::FromRedisValue;
use crate::redis::RedisModel;
use serde::{Deserialize, Serialize};
use uuid::Uuid;
use super::*;
#[derive(Debug, Clone, Serialize, Deserialize, FromRedisValue, PartialEq)]
struct Tst {
key: String,
a: u64,
b: u64,
}
impl RedisModel for Tst {
type Key = String;
type Value = String;
fn key_ref(&self) -> &Self::Key {
&self.key
}
fn key(&self) -> redis::Result<Self::Key> {
Ok(self.key.clone())
}
fn value(&self) -> redis::Result<impl deadpool_redis::redis::ToRedisArgs + Send + Sync> {
Ok(serde_json::to_string(&self)?)
}
fn value_ref(&self) -> &Self::Value {
static PLACEHOLDER: String = String::new();
&PLACEHOLDER
}
}
impl Tst {
pub fn inc(mut self, value: u64) -> Self {
self.a += value;
self.b += value;
self
}
pub fn default(key: impl AsRef<str>) -> Self {
Self {
key: key.as_ref().to_string(),
a: 3,
b: 4,
}
}
}
async fn get_client() -> Client {
Client::default().await.unwrap()
}
#[tokio::test]
async fn test_redis_get() -> Result<()> {
let client = get_client().await;
let key = "test_redis_get";
let fx_model = Tst::default(key);
client.set(&fx_model).await?;
assert_eq!(Some(fx_model), client.get(key).await?);
client.del(key).await?;
Ok(())
}
#[tokio::test]
async fn test_redis_mget() -> Result<()> {
let client = get_client().await;
let key1 = "test_redis_mget1".to_string();
let key2 = "test_redis_mget2".to_string();
let model1 = Tst::default(&key1);
let model2 = Tst::default(&key2);
let tuple1 = (key1.clone(), serde_json::to_string(&model1)?);
let tuple2 = (key2.clone(), serde_json::to_string(&model2)?);
client.mset([&tuple1, &tuple2]).await?;
let got: Vec<Option<Tst>> = client.mget(&[&key1, &key2]).await?;
assert_eq!(vec![Some(model1), Some(model2)], got);
client.mdel([&key1, &key2]).await?;
Ok(())
}
#[tokio::test]
async fn test_redis_get_ex() -> Result<()> {
let client = get_client().await;
let key = "test_redis_get_ex";
let fx_model = Tst::default(key);
client.set(&fx_model).await?;
assert_eq!(Some(fx_model), client.get_ex(key, Expiry::EX(2)).await?);
tokio::time::sleep(Duration::from_secs(3)).await;
assert_eq!(None::<Tst>, client.get(key).await?);
Ok(())
}
#[tokio::test]
async fn test_redis_get_del() -> Result<()> {
let client = get_client().await;
let key = "test_redis_get_del";
let fx_model = Tst::default(key);
client.set(&fx_model).await?;
assert_eq!(Some(fx_model), client.get_del(key).await?);
assert_eq!(None::<Tst>, client.get(key).await?);
Ok(())
}
#[tokio::test]
async fn test_redis_getset() -> Result<()> {
let client = get_client().await;
let key = "test_redis_getset";
let to_get = Tst::default(key);
let to_set = Tst::default(key).inc(5);
client.set(&to_get).await?;
assert_eq!(Some(to_get), client.getset(&to_set).await?);
assert_eq!(Some(to_set), client.get(key).await?);
client.del(key).await?;
Ok(())
}
#[tokio::test]
async fn test_redis_set() -> Result<()> {
let client = get_client().await;
let key = "test_redis_set";
let fx_model = Tst::default(key);
assert_eq!("OK", client.set(&fx_model).await?);
assert_eq!(Some(fx_model), client.get(key).await?);
client.del(key).await?;
Ok(())
}
#[tokio::test]
async fn test_redis_mset() -> Result<()> {
let client = get_client().await;
let key1 = "test_redis_mset1".to_string();
let key2 = "test_redis_mset2".to_string();
let model1 = Tst::default(&key1);
let model2 = Tst::default(&key2);
let tuple1 = (key1.clone(), serde_json::to_string(&model1)?);
let tuple2 = (key2.clone(), serde_json::to_string(&model2)?);
assert_eq!("OK", client.mset([&tuple1, &tuple2]).await?);
assert_eq!(Some(model1), client.get(&key1).await?);
assert_eq!(Some(model2), client.get(&key2).await?);
client.mdel([&key1, &key2]).await?;
Ok(())
}
#[tokio::test]
async fn test_redis_set_ex() -> Result<()> {
let client = get_client().await;
let key = "test_redis_set_ex";
let fx_model = Tst::default(key);
assert_eq!("OK", client.set_ex(&fx_model, 2).await?);
assert_eq!(Some(fx_model), client.get(key).await?);
tokio::time::sleep(Duration::from_secs(3)).await;
assert_eq!(None::<Tst>, client.get(key).await?);
client.del(key).await?;
Ok(())
}
#[tokio::test]
async fn test_redis_set_nx() -> Result<()> {
let client = get_client().await;
let key = "test_redis_set_nx1".to_string();
let model1 = Tst::default(&key);
let model2 = Tst::default(&key).inc(5);
let _tuple1 = (key.clone(), serde_json::to_string(&model1)?);
let _tuple2 = (key.clone(), serde_json::to_string(&model2)?);
assert!(client.set_nx(&model1).await?);
assert_eq!(Some(model1.clone()), client.get(&key).await?);
assert!(!client.set_nx(&model2).await?);
assert_eq!(Some(model1), client.get(&key).await?);
client.del(&key).await?;
Ok(())
}
#[tokio::test]
async fn test_redis_mset_nx() -> Result<()> {
let client = get_client().await;
let key1 = format!("test_redis_mset_nx1_{}", Uuid::new_v4());
let key2 = format!("test_redis_mset_nx2_{}", Uuid::new_v4());
let model1_before = Tst::default(&key1);
let model2_before = Tst::default(&key2);
let model1_after = Tst::default(&key1).inc(3);
let model2_after = Tst::default(&key2).inc(3);
let tuple1_before = (key1.clone(), serde_json::to_string(&model1_before)?);
let tuple2_before = (key2.clone(), serde_json::to_string(&model2_before)?);
let tuple1_after = (key1.clone(), serde_json::to_string(&model1_after)?);
let tuple2_after = (key2.clone(), serde_json::to_string(&model2_after)?);
assert!(client.mset_nx([&tuple1_before, &tuple2_before]).await?);
assert_eq!(
vec![Some(model1_before.clone()), Some(model2_before.clone())],
client.mget(&[&key1, &key2]).await?
);
assert!(!client.mset_nx([&tuple1_after, &tuple2_after]).await?);
assert_eq!(
vec![Some(model1_before), Some(model2_before)],
client.mget(&[&key1, &key2]).await?
);
client.mdel([&key1, &key2]).await?;
Ok(())
}
#[tokio::test]
async fn test_redis_del() -> Result<()> {
let client = get_client().await;
let key = "test_redis_del".to_string();
let fx_model = Tst::default(&key);
let tuple = (key.clone(), serde_json::to_string(&fx_model)?);
client.mset([&tuple]).await?;
assert_eq!(Some(fx_model), client.get(&key).await?);
assert!(client.del(&key).await?);
assert_eq!(None::<Tst>, client.get(&key).await?);
assert!(!client.del(&key).await?);
Ok(())
}
#[tokio::test]
async fn test_redis_mdel() -> Result<()> {
let client = get_client().await;
let key1 = "test_redis_mdel1".to_string();
let key2 = "test_redis_mdel2".to_string();
let model1 = Tst::default(&key1);
let model2 = Tst::default(&key2);
let tuple1 = (key1.clone(), serde_json::to_string(&model1)?);
let tuple2 = (key2.clone(), serde_json::to_string(&model2)?);
client.mset([&tuple1, &tuple2]).await?;
assert_eq!(Some(model1), client.get(&key1).await?);
assert_eq!(Some(model2), client.get(&key2).await?);
assert_eq!(2, client.mdel([&key1, &key2]).await?);
assert_eq!(None::<Tst>, client.get(&key1).await?);
assert_eq!(None::<Tst>, client.get(&key2).await?);
assert_eq!(0, client.mdel([&key1, &key2]).await?);
Ok(())
}
#[tokio::test]
async fn test_redis_exists() -> Result<()> {
let client = get_client().await;
let key = "test_redis_exists".to_string();
assert!(!client.exists(&key).await?);
let fx_model = Tst::default(&key);
let tuple = (key.clone(), serde_json::to_string(&fx_model)?);
client.mset([&tuple]).await?;
assert!(client.exists(&key).await?);
client.del(&key).await?;
Ok(())
}
#[tokio::test]
async fn test_redis_ping() -> Result<()> {
let client = get_client().await;
assert_eq!("PONG", client.ping().await?);
Ok(())
}
#[tokio::test]
async fn test_redis_rename() -> Result<()> {
let client = get_client().await;
let key = "test_redis_rename".to_string();
let new_key = "test_redis_rename_new".to_string();
let fx_key_model = Tst::default(&key);
let tuple = (key.clone(), serde_json::to_string(&fx_key_model)?);
client.mset([&tuple]).await?;
assert_eq!("OK", client.rename(&key, &new_key).await?);
let res: Option<Tst> = client.get(&key).await?;
assert_eq!(None, res);
assert_eq!(Some(fx_key_model), client.get(&new_key).await?);
client.del(&new_key).await?;
Ok(())
}
#[tokio::test]
async fn test_redis_rename_nx() -> Result<()> {
let client = get_client().await;
let key = "test_redis_rename_nx".to_string();
let new_key = "test_redis_rename_nx_new".to_string();
let fx_key_model = Tst::default(&key);
let tuple = (key.clone(), serde_json::to_string(&fx_key_model)?);
client.mset([&tuple]).await?;
assert!(client.rename_nx(&key, &new_key).await?);
let res: Option<Tst> = client.get(&key).await?;
assert_eq!(None, res);
assert_eq!(Some(fx_key_model.clone()), client.get(&new_key).await?);
let fx_key_new_model = Tst::default(&key);
let tuple2 = (key.clone(), serde_json::to_string(&fx_key_new_model)?);
client.mset([&tuple2]).await?;
assert!(!client.rename_nx(&new_key, &key).await?);
client.mdel([&key, &new_key]).await?;
Ok(())
}
}