Skip to main content

ferrox_database_redis/
lib.rs

1//! # Ferrox Database Redis (`ferrox-database-redis`)
2//!
3//! `ferrox-database-redis` provides Redis integration for caching, session storage, distributed rate limiting, and real-time Pub/Sub.
4//!
5//! ## Key Features
6//! - ⚡ **Multiplexed Connection Pool**: Efficient async Redis client backed by `bb8` or `redis-rs`.
7//! - 🔑 **Cache Helper Operations**: Strongly typed `get_json`, `set_json`, `expire`, and `del` primitives.
8//! - 📻 **Pub/Sub Subscriptions**: Asynchronous message receiver streams.
9
10use bb8::Pool;
11use bb8_redis::RedisConnectionManager;
12use redis::AsyncCommands;
13use serde::{Deserialize, Serialize};
14use ferrox_errors::AppError;
15use tracing;
16
17pub type RedisPool = Pool<RedisConnectionManager>;
18
19#[derive(Clone)]
20pub struct RedisClient {
21    pub pool: RedisPool,
22}
23
24impl RedisClient {
25    pub async fn connect(connection_string: &str) -> Result<Self, AppError> {
26        let manager = RedisConnectionManager::new(connection_string)
27            .map_err(|e| AppError::DatabaseError(format!("Redis URL Parse Error: {}", e)))?;
28        
29        // Setup an async connection pool
30        let pool = Pool::builder()
31            .max_size(15) // Enough for most web backends
32            .build(manager)
33            .await
34            .map_err(|e| AppError::DatabaseError(format!("Redis Pool Error: {}", e)))?;
35
36        // Test connection
37        let pool_clone = pool.clone();
38        let mut conn = pool_clone.get().await
39            .map_err(|e| AppError::DatabaseError(format!("Redis Ping Failed: {}", e)))?;
40            
41        let _: String = redis::cmd("PING")
42            .query_async(&mut *conn)
43            .await
44            .map_err(|e| AppError::DatabaseError(format!("Redis Ping Error: {}", e)))?;
45
46        tracing::info!("Connected to Redis via bb8 connection pool.");
47
48        Ok(Self { pool })
49    }
50
51    /// Helper to store generic serializable objects in Redis with TTL
52    pub async fn set_json<T: Serialize>(&self, key: &str, value: &T, ttl_seconds: u64) -> Result<(), AppError> {
53        let mut conn = self.pool.get().await
54            .map_err(|e| AppError::DatabaseError(e.to_string()))?;
55            
56        let json_str = serde_json::to_string(value)
57            .map_err(|e| AppError::InternalServerError(Box::new(e)))?;
58            
59        let _: () = conn.set_ex(key, json_str, ttl_seconds)
60            .await
61            .map_err(|e| AppError::DatabaseError(e.to_string()))?;
62            
63        Ok(())
64    }
65
66    /// Helper to retrieve generic objects from Redis
67    pub async fn get_json<T: for<'de> Deserialize<'de>>(&self, key: &str) -> Result<Option<T>, AppError> {
68        let mut conn = self.pool.get().await
69            .map_err(|e| AppError::DatabaseError(e.to_string()))?;
70            
71        let result: Option<String> = conn.get(key)
72            .await
73            .map_err(|e| AppError::DatabaseError(e.to_string()))?;
74            
75        match result {
76            Some(json_str) => {
77                let obj = serde_json::from_str(&json_str)
78                    .map_err(|e| AppError::InternalServerError(Box::new(e)))?;
79                Ok(Some(obj))
80            },
81            None => Ok(None)
82        }
83    }
84}