openzeppelin_relayer/queues/redis/
queue.rs

1//! Queue management module for job processing.
2//!
3//! This module provides Redis-backed queue implementation for handling different types of jobs:
4//! - Transaction requests
5//! - Transaction submissions
6//! - Transaction status checks
7//! - Notifications
8//! - Solana swap requests
9//! - Relayer health checks
10use std::{env, sync::Arc};
11
12use apalis_redis::{Config, RedisStorage};
13use color_eyre::{eyre, Result};
14use redis::aio::{ConnectionManager, ConnectionManagerConfig};
15use serde::{Deserialize, Serialize};
16use tokio::time::Duration;
17use tracing::info;
18
19use crate::queues::redis::refreshing_connection::RefreshingConnection;
20use crate::{config::ServerConfig, utils::RedisConnections};
21
22use crate::jobs::{
23    Job, NotificationSend, RelayerHealthCheck, TokenSwapRequest, TransactionRequest,
24    TransactionSend, TransactionStatusCheck,
25};
26
27/// Storage type for all queues.
28///
29/// Uses [`RefreshingConnection`] instead of a bare `ConnectionManager` so that
30/// connections are dropped/reopened on a bounded lifetime, letting them follow
31/// the endpoint's DNS whenever it changes (e.g. after an ElastiCache failover
32/// repoints the endpoint to a new node). Connections are ALSO rebuilt
33/// reactively when a reply indicates the connection is pinned to a read-only
34/// node, which accelerates recovery; the bounded lifetime remains the
35/// guaranteed backstop.
36pub type QueueStorage<T> = RedisStorage<T, RefreshingConnection<ConnectionManager>>;
37
38#[derive(Clone)]
39pub struct Queue {
40    pub transaction_request_queue: QueueStorage<Job<TransactionRequest>>,
41    pub transaction_submission_queue: QueueStorage<Job<TransactionSend>>,
42    /// Default/fallback status queue for backward compatibility, Solana, and future networks
43    pub transaction_status_queue: QueueStorage<Job<TransactionStatusCheck>>,
44    /// EVM-specific status queue with slower retries
45    pub transaction_status_queue_evm: QueueStorage<Job<TransactionStatusCheck>>,
46    /// Stellar-specific status queue with fast retries
47    pub transaction_status_queue_stellar: QueueStorage<Job<TransactionStatusCheck>>,
48    pub notification_queue: QueueStorage<Job<NotificationSend>>,
49    pub token_swap_request_queue: QueueStorage<Job<TokenSwapRequest>>,
50    pub relayer_health_check_queue: QueueStorage<Job<RelayerHealthCheck>>,
51    /// Redis connection pools for handlers that need pool-based access.
52    /// Provides both primary (write) and reader (read) pools for:
53    /// - Distributed locking
54    /// - Status check metadata (failure counters)
55    /// - Any handler-specific Redis operations
56    redis_connections: Arc<RedisConnections>,
57}
58
59impl std::fmt::Debug for Queue {
60    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
61        f.debug_struct("Queue")
62            .field("transaction_request_queue", &"RedisStorage<...>")
63            .field("transaction_submission_queue", &"RedisStorage<...>")
64            .field("transaction_status_queue", &"RedisStorage<...>")
65            .field("transaction_status_queue_evm", &"RedisStorage<...>")
66            .field("transaction_status_queue_stellar", &"RedisStorage<...>")
67            .field("notification_queue", &"RedisStorage<...>")
68            .field("token_swap_request_queue", &"RedisStorage<...>")
69            .field("relayer_health_check_queue", &"RedisStorage<...>")
70            .field("redis_connections", &"RedisConnections")
71            .finish()
72    }
73}
74
75/// Configuration for queue storage tuning.
76#[derive(Clone, Debug)]
77struct QueueConfig {
78    /// How often to move scheduled jobs to active queue (default: 30s)
79    enqueue_scheduled: Duration,
80}
81
82impl Default for QueueConfig {
83    fn default() -> Self {
84        Self {
85            enqueue_scheduled: Duration::from_secs(30),
86        }
87    }
88}
89
90impl QueueConfig {
91    /// - Faster poll interval for lower latency
92    fn high_frequency() -> Self {
93        Self {
94            enqueue_scheduled: Duration::from_secs(2),
95        }
96    }
97
98    /// Configuration for lower-frequency queues.
99    /// - Smaller buffer (less memory pressure)
100    /// - Slower scheduled job polling (reduces Redis load)
101    fn low_frequency() -> Self {
102        Self {
103            enqueue_scheduled: Duration::from_secs(20),
104        }
105    }
106}
107
108impl Queue {
109    /// Creates a RedisStorage for a specific job type using a ConnectionManager.
110    ///
111    /// # Arguments
112    /// * `namespace` - Redis key namespace for this queue
113    /// * `conn` - ConnectionManager with auto-reconnect
114    /// * `queue_config` - Tuning parameters for this queue
115    ///
116    /// ConnectionManager provides automatic reconnection on connection failures,
117    /// ensuring queue processing continues even if the Redis connection drops temporarily.
118    fn storage<T: Serialize + for<'de> Deserialize<'de>>(
119        namespace: &str,
120        conn: RefreshingConnection<ConnectionManager>,
121        queue_config: QueueConfig,
122    ) -> QueueStorage<T> {
123        let config = Config::default()
124            .set_namespace(namespace)
125            .set_enqueue_scheduled(queue_config.enqueue_scheduled);
126
127        RedisStorage::new_with_config(conn, config)
128    }
129
130    /// Creates a `RefreshingConnection` wrapping a `ConnectionManager` with the
131    /// standard queue configuration.
132    ///
133    /// Each ConnectionManager represents a single Redis connection with auto-reconnect.
134    /// Creating separate managers for different queue types enables parallel Redis operations.
135    ///
136    /// The connection is wrapped in a [`RefreshingConnection`] so it is
137    /// recycled on age and on read-only replies (see [`QueueStorage`]).
138    async fn create_connection_manager(
139        client: &redis::Client,
140        queue_timeout: Duration,
141        max_age_ms: u64,
142    ) -> Result<RefreshingConnection<ConnectionManager>> {
143        let conn_config = ConnectionManagerConfig::new()
144            .set_connection_timeout(queue_timeout)
145            .set_response_timeout(queue_timeout)
146            .set_number_of_retries(2)
147            .set_max_delay(1000);
148
149        let conn = ConnectionManager::new_with_config(client.clone(), conn_config.clone())
150            .await
151            .map_err(|e| eyre::eyre!("Failed to create Redis connection manager: {}", e))?;
152
153        Ok(RefreshingConnection::new(
154            client.clone(),
155            conn_config,
156            conn,
157            max_age_ms,
158        ))
159    }
160
161    /// Sets up all job queues with properly configured Redis connections.
162    ///
163    /// # Architecture
164    /// - **Queue storages**: Each queue gets its own `ConnectionManager` for maximum parallelism
165    /// - **Handler operations**: Use `redis_connections` pool for metadata, locking, counters
166    ///
167    /// # Connection Strategy
168    /// Each queue has a dedicated Redis connection to prevent contention under high throughput.
169    /// This allows 8 parallel Redis operations (one per queue type).
170    ///
171    /// # Arguments
172    /// * `redis_connections` - Redis connection pools for handler operations.
173    ///
174    /// # Connection Configuration
175    /// - `connection_timeout`: Max time to establish TCP connection to Redis
176    /// - `response_timeout`: Max time to wait for Redis command responses
177    /// - Auto-reconnect: ConnectionManager automatically reconnects on failures
178    pub async fn setup(redis_connections: Arc<RedisConnections>) -> Result<Self> {
179        let server_config = ServerConfig::from_env();
180        let redis_url = &server_config.redis_url;
181
182        // Create Redis client
183        let client = redis::Client::open(redis_url.as_str())
184            .map_err(|e| eyre::eyre!("Failed to create Redis client for queue: {}", e))?;
185
186        // Configure timeout for all ConnectionManagers
187        // Worst case calculation: 3 attempts × 5s timeout + ~0.3s backoff = ~15.3s
188        let queue_timeout = Duration::from_secs(5);
189
190        // Bounded connection lifetime (see `QueueStorage`). `0` disables only
191        // the age backstop, not reactive read-only healing.
192        let max_age_ms = server_config.redis_connection_max_age_ms;
193
194        // Create one ConnectionManager per queue to prevent connection contention.
195        // Each ConnectionManager is a single Redis connection with auto-reconnect.
196        let conn_tx_request =
197            Self::create_connection_manager(&client, queue_timeout, max_age_ms).await?;
198        let conn_tx_submit =
199            Self::create_connection_manager(&client, queue_timeout, max_age_ms).await?;
200        let conn_status =
201            Self::create_connection_manager(&client, queue_timeout, max_age_ms).await?;
202        let conn_status_evm =
203            Self::create_connection_manager(&client, queue_timeout, max_age_ms).await?;
204        let conn_status_stellar =
205            Self::create_connection_manager(&client, queue_timeout, max_age_ms).await?;
206        let conn_notification =
207            Self::create_connection_manager(&client, queue_timeout, max_age_ms).await?;
208        let conn_swap = Self::create_connection_manager(&client, queue_timeout, max_age_ms).await?;
209        let conn_health =
210            Self::create_connection_manager(&client, queue_timeout, max_age_ms).await?;
211
212        info!(
213            redis_url = %redis_url,
214            connection_timeout_ms = 5000,
215            response_timeout_ms = 5000,
216            retries = 2,
217            max_backoff_ms = 1000,
218            connection_count = 8,
219            "Queue setup: created dedicated ConnectionManager per queue"
220        );
221
222        // use REDIS_KEY_PREFIX only if set, otherwise do not use it
223        let redis_key_prefix = env::var("REDIS_KEY_PREFIX")
224            .ok()
225            .filter(|v| !v.is_empty())
226            .map(|value| format!("{value}:queue:"))
227            .unwrap_or_default();
228
229        // Queue configurations:
230        // - High-frequency: transaction_status (critical path)
231        // - Low-frequency: request, submission, notifications, health checks, swaps
232        let high_frequency = QueueConfig::high_frequency();
233        let low_frequency = QueueConfig::low_frequency();
234
235        Ok(Self {
236            transaction_request_queue: Self::storage(
237                &format!("{redis_key_prefix}transaction_request_queue"),
238                conn_tx_request,
239                low_frequency.clone(), // scheduling not used
240            ),
241            transaction_submission_queue: Self::storage(
242                &format!("{redis_key_prefix}transaction_submission_queue"),
243                conn_tx_submit,
244                low_frequency.clone(), // scheduling not used
245            ),
246            transaction_status_queue: Self::storage(
247                &format!("{redis_key_prefix}transaction_status_queue"),
248                conn_status,
249                high_frequency.clone(),
250            ),
251            transaction_status_queue_evm: Self::storage(
252                &format!("{redis_key_prefix}transaction_status_queue_evm"),
253                conn_status_evm,
254                high_frequency.clone(),
255            ),
256            transaction_status_queue_stellar: Self::storage(
257                &format!("{redis_key_prefix}transaction_status_queue_stellar"),
258                conn_status_stellar,
259                high_frequency.clone(),
260            ),
261            // Lower-frequency queues
262            notification_queue: Self::storage(
263                &format!("{redis_key_prefix}notification_queue"),
264                conn_notification,
265                low_frequency.clone(), // scheduling not used
266            ),
267            token_swap_request_queue: Self::storage(
268                &format!("{redis_key_prefix}token_swap_request_queue"),
269                conn_swap,
270                low_frequency.clone(), // scheduling not used
271            ),
272            relayer_health_check_queue: Self::storage(
273                &format!("{redis_key_prefix}relayer_health_check_queue"),
274                conn_health,
275                low_frequency.clone(), // scheduling not used
276            ),
277            redis_connections,
278        })
279    }
280
281    /// Returns the Redis connection pools.
282    ///
283    /// This provides access to both primary and reader pools for handlers
284    /// that need Redis pool-based access (e.g., for metadata storage, distributed locking).
285    pub fn redis_connections(&self) -> Arc<RedisConnections> {
286        self.redis_connections.clone()
287    }
288}
289
290#[cfg(test)]
291mod tests {
292    use super::*;
293
294    #[test]
295    fn test_queue_storage_configuration() {
296        // Test the config creation logic without actual Redis connections
297        let namespace = "test_namespace";
298        let config = Config::default().set_namespace(namespace);
299
300        assert_eq!(config.get_namespace(), namespace);
301    }
302
303    // Mock version of Queue for testing
304    #[derive(Clone, Debug)]
305    struct MockQueue {
306        pub namespace_transaction_request: String,
307        pub namespace_transaction_submission: String,
308        pub namespace_transaction_status: String,
309        pub namespace_transaction_status_evm: String,
310        pub namespace_transaction_status_stellar: String,
311        pub namespace_notification: String,
312        pub namespace_token_swap_request_queue: String,
313        pub namespace_relayer_health_check_queue: String,
314    }
315
316    impl MockQueue {
317        fn new() -> Self {
318            Self {
319                namespace_transaction_request: "transaction_request_queue".to_string(),
320                namespace_transaction_submission: "transaction_submission_queue".to_string(),
321                namespace_transaction_status: "transaction_status_queue".to_string(),
322                namespace_transaction_status_evm: "transaction_status_queue_evm".to_string(),
323                namespace_transaction_status_stellar: "transaction_status_queue_stellar"
324                    .to_string(),
325                namespace_notification: "notification_queue".to_string(),
326                namespace_token_swap_request_queue: "token_swap_request_queue".to_string(),
327                namespace_relayer_health_check_queue: "relayer_health_check_queue".to_string(),
328            }
329        }
330    }
331
332    #[test]
333    fn test_queue_namespaces() {
334        let mock_queue = MockQueue::new();
335
336        assert_eq!(
337            mock_queue.namespace_transaction_request,
338            "transaction_request_queue"
339        );
340        assert_eq!(
341            mock_queue.namespace_transaction_submission,
342            "transaction_submission_queue"
343        );
344        assert_eq!(
345            mock_queue.namespace_transaction_status,
346            "transaction_status_queue"
347        );
348        assert_eq!(
349            mock_queue.namespace_transaction_status_evm,
350            "transaction_status_queue_evm"
351        );
352        assert_eq!(
353            mock_queue.namespace_transaction_status_stellar,
354            "transaction_status_queue_stellar"
355        );
356        assert_eq!(mock_queue.namespace_notification, "notification_queue");
357        assert_eq!(
358            mock_queue.namespace_token_swap_request_queue,
359            "token_swap_request_queue"
360        );
361        assert_eq!(
362            mock_queue.namespace_relayer_health_check_queue,
363            "relayer_health_check_queue"
364        );
365    }
366
367    #[test]
368    fn test_queue_config_with_prefix() {
369        // Test that namespace includes prefix when set
370        let prefix = "myprefix:queue:";
371        let queue_name = "transaction_request_queue";
372        let full_namespace = format!("{prefix}{queue_name}");
373
374        let config = Config::default().set_namespace(&full_namespace);
375        assert_eq!(
376            config.get_namespace(),
377            "myprefix:queue:transaction_request_queue"
378        );
379    }
380
381    #[test]
382    fn test_queue_config_without_prefix() {
383        // Test that namespace works without prefix
384        let queue_name = "transaction_request_queue";
385
386        let config = Config::default().set_namespace(queue_name);
387        assert_eq!(config.get_namespace(), "transaction_request_queue");
388    }
389
390    #[test]
391    fn test_queue_config_default() {
392        let config = QueueConfig::default();
393
394        assert_eq!(config.enqueue_scheduled, Duration::from_secs(30));
395    }
396
397    #[test]
398    fn test_queue_config_high_throughput() {
399        let config = QueueConfig::high_frequency();
400
401        // High frequency should have faster enqueue_scheduled
402        assert_eq!(config.enqueue_scheduled, Duration::from_secs(2));
403
404        // Verify it's faster than default
405        let default = QueueConfig::default();
406        assert!(config.enqueue_scheduled < default.enqueue_scheduled);
407    }
408
409    #[test]
410    fn test_queue_config_low_frequency() {
411        let config = QueueConfig::low_frequency();
412
413        assert_eq!(config.enqueue_scheduled, Duration::from_secs(20));
414
415        // Low frequency should have longer enqueue_scheduled than high frequency
416        let high = QueueConfig::high_frequency();
417        assert!(config.enqueue_scheduled > high.enqueue_scheduled);
418    }
419}