use crate::db::get_redis_conn;
use crate::types::DynError;
use deadpool_redis::redis::AsyncCommands;
pub async fn put(
prefix: &str,
key: &str,
values: &[&str],
expiration: Option<i64>,
) -> Result<(), DynError> {
if values.is_empty() {
return Ok(());
}
let index_key = format!("{prefix}:{key}");
let mut redis_conn = get_redis_conn().await?;
let mut pipe = redis::pipe();
pipe.sadd(&index_key, values);
if let Some(ttl) = expiration {
pipe.expire(&index_key, ttl);
}
let _: () = pipe.query_async(&mut redis_conn).await?;
Ok(())
}
pub async fn get_range(
prefix: &str,
key: &str,
skip: Option<usize>,
limit: Option<usize>,
) -> Result<Option<Vec<String>>, DynError> {
let mut redis_conn = get_redis_conn().await?;
let index_key = format!("{prefix}:{key}");
let mut cursor = "0".to_string();
let mut collected: Vec<String> = Vec::new();
let skip = skip.unwrap_or(0);
let limit = limit.unwrap_or(5);
let mut skipped = 0;
while collected.len() < limit {
let result: (String, Vec<String>) = redis::cmd("SSCAN")
.arg(&index_key)
.arg(&cursor)
.arg("COUNT")
.arg(limit)
.query_async(&mut redis_conn)
.await?;
let (new_cursor, items) = result;
for item in items {
if skipped < skip {
skipped += 1;
continue;
}
collected.push(item);
if collected.len() >= limit {
break;
}
}
cursor = new_cursor;
if cursor == "0" {
break; }
}
if collected.is_empty() {
Ok(None)
} else {
Ok(Some(collected))
}
}
pub async fn check_member(prefix: &str, key: &str, member: &str) -> Result<(bool, bool), DynError> {
let mut redis_conn = get_redis_conn().await?;
let index_key = format!("{prefix}:{key}");
let set_exists: bool = redis_conn.exists(&index_key).await?;
if set_exists {
let is_member: bool = redis_conn.sismember(&index_key, member).await?;
Ok((true, is_member))
} else {
Ok((false, false))
}
}
pub async fn get_size(prefix: &str, key: &str) -> Result<Option<usize>, DynError> {
let mut redis_conn = get_redis_conn().await?;
let index_key = format!("{prefix}:{key}");
let set_exists: bool = redis_conn.exists(&index_key).await?;
if !set_exists {
return Ok(None);
}
let set_size: usize = redis_conn.scard(&index_key).await?;
Ok(Some(set_size))
}
pub async fn get_multiple_sets(
prefix: &str,
keys: &[&str],
member: Option<&str>,
limit: Option<usize>,
) -> Result<Vec<Option<(Vec<String>, usize, bool)>>, DynError> {
let mut redis_conn = get_redis_conn().await?;
let mut pipe = redis::pipe();
for key in keys {
let index_key = format!("{prefix}:{key}");
pipe.smembers(index_key);
}
let results: Vec<Vec<String>> = pipe.query_async(&mut redis_conn).await?;
let taggers_list = results
.into_iter()
.map(|set| {
if set.is_empty() {
None
} else {
let set_length = set.len();
let is_member = member
.map(|member_to_search| set.iter().any(|s| s == member_to_search))
.unwrap_or(false);
match limit {
Some(set_limit) if set_limit < set_length => {
let limited_set = set.into_iter().take(set_limit).collect();
Some((limited_set, set_length, is_member))
}
_ => Some((set, set_length, is_member)),
}
}
})
.collect();
Ok(taggers_list)
}
pub async fn put_multiple_sets(
prefix: &str,
common_key: &[&str],
index: &[&str],
collections: &[&[&str]],
expiration: Option<i64>,
) -> Result<(), DynError> {
let mut redis_conn = get_redis_conn().await?;
let mut pipe = redis::pipe();
for (i, key) in index.iter().enumerate() {
let full_index = format!("{}:{}:{}", &prefix, common_key.join(":"), key);
if !collections[i].is_empty() {
pipe.sadd(&full_index, collections[i]); if let Some(ttl) = expiration {
pipe.expire(&full_index, ttl);
}
}
}
let _: () = pipe.query_async(&mut redis_conn).await?;
Ok(())
}
pub async fn del(prefix: &str, key: &str, values: &[&str]) -> Result<(), DynError> {
if values.is_empty() {
return Ok(());
}
let index_key = format!("{prefix}:{key}");
let mut redis_conn = get_redis_conn().await?;
let _: () = redis_conn.srem(index_key, values).await?;
Ok(())
}
pub async fn get_random_members(
prefix: &str,
key: &str,
count: isize,
) -> Result<Option<Vec<String>>, DynError> {
let mut redis_conn = get_redis_conn().await?;
let index_key = format!("{prefix}:{key}");
let set_exists: bool = redis_conn.exists(&index_key).await?;
if !set_exists {
return Ok(None);
}
let random_members: Vec<String> = redis::cmd("SRANDMEMBER")
.arg(&index_key)
.arg(count)
.query_async(&mut redis_conn)
.await?;
Ok(Some(random_members))
}