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_submission, transaction_status (critical path)
231        // - Low-frequency: request, 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                // Stellar fast resubmit schedules 5s*attempt delayed submit jobs here.
245                // A 20s sweep quantizes those retries and lets the status checker double-fire.
246                high_frequency.clone(),
247            ),
248            transaction_status_queue: Self::storage(
249                &format!("{redis_key_prefix}transaction_status_queue"),
250                conn_status,
251                high_frequency.clone(),
252            ),
253            transaction_status_queue_evm: Self::storage(
254                &format!("{redis_key_prefix}transaction_status_queue_evm"),
255                conn_status_evm,
256                high_frequency.clone(),
257            ),
258            transaction_status_queue_stellar: Self::storage(
259                &format!("{redis_key_prefix}transaction_status_queue_stellar"),
260                conn_status_stellar,
261                high_frequency.clone(),
262            ),
263            // Lower-frequency queues
264            notification_queue: Self::storage(
265                &format!("{redis_key_prefix}notification_queue"),
266                conn_notification,
267                low_frequency.clone(), // scheduling not used
268            ),
269            token_swap_request_queue: Self::storage(
270                &format!("{redis_key_prefix}token_swap_request_queue"),
271                conn_swap,
272                low_frequency.clone(), // scheduling not used
273            ),
274            relayer_health_check_queue: Self::storage(
275                &format!("{redis_key_prefix}relayer_health_check_queue"),
276                conn_health,
277                low_frequency.clone(), // scheduling not used
278            ),
279            redis_connections,
280        })
281    }
282
283    /// Returns the Redis connection pools.
284    ///
285    /// This provides access to both primary and reader pools for handlers
286    /// that need Redis pool-based access (e.g., for metadata storage, distributed locking).
287    pub fn redis_connections(&self) -> Arc<RedisConnections> {
288        self.redis_connections.clone()
289    }
290}
291
292#[cfg(test)]
293mod tests {
294    use super::*;
295
296    #[test]
297    fn test_queue_storage_configuration() {
298        // Test the config creation logic without actual Redis connections
299        let namespace = "test_namespace";
300        let config = Config::default().set_namespace(namespace);
301
302        assert_eq!(config.get_namespace(), namespace);
303    }
304
305    // Mock version of Queue for testing
306    #[derive(Clone, Debug)]
307    struct MockQueue {
308        pub namespace_transaction_request: String,
309        pub namespace_transaction_submission: String,
310        pub namespace_transaction_status: String,
311        pub namespace_transaction_status_evm: String,
312        pub namespace_transaction_status_stellar: String,
313        pub namespace_notification: String,
314        pub namespace_token_swap_request_queue: String,
315        pub namespace_relayer_health_check_queue: String,
316    }
317
318    impl MockQueue {
319        fn new() -> Self {
320            Self {
321                namespace_transaction_request: "transaction_request_queue".to_string(),
322                namespace_transaction_submission: "transaction_submission_queue".to_string(),
323                namespace_transaction_status: "transaction_status_queue".to_string(),
324                namespace_transaction_status_evm: "transaction_status_queue_evm".to_string(),
325                namespace_transaction_status_stellar: "transaction_status_queue_stellar"
326                    .to_string(),
327                namespace_notification: "notification_queue".to_string(),
328                namespace_token_swap_request_queue: "token_swap_request_queue".to_string(),
329                namespace_relayer_health_check_queue: "relayer_health_check_queue".to_string(),
330            }
331        }
332    }
333
334    #[test]
335    fn test_queue_namespaces() {
336        let mock_queue = MockQueue::new();
337
338        assert_eq!(
339            mock_queue.namespace_transaction_request,
340            "transaction_request_queue"
341        );
342        assert_eq!(
343            mock_queue.namespace_transaction_submission,
344            "transaction_submission_queue"
345        );
346        assert_eq!(
347            mock_queue.namespace_transaction_status,
348            "transaction_status_queue"
349        );
350        assert_eq!(
351            mock_queue.namespace_transaction_status_evm,
352            "transaction_status_queue_evm"
353        );
354        assert_eq!(
355            mock_queue.namespace_transaction_status_stellar,
356            "transaction_status_queue_stellar"
357        );
358        assert_eq!(mock_queue.namespace_notification, "notification_queue");
359        assert_eq!(
360            mock_queue.namespace_token_swap_request_queue,
361            "token_swap_request_queue"
362        );
363        assert_eq!(
364            mock_queue.namespace_relayer_health_check_queue,
365            "relayer_health_check_queue"
366        );
367    }
368
369    #[test]
370    fn test_queue_config_with_prefix() {
371        // Test that namespace includes prefix when set
372        let prefix = "myprefix:queue:";
373        let queue_name = "transaction_request_queue";
374        let full_namespace = format!("{prefix}{queue_name}");
375
376        let config = Config::default().set_namespace(&full_namespace);
377        assert_eq!(
378            config.get_namespace(),
379            "myprefix:queue:transaction_request_queue"
380        );
381    }
382
383    #[test]
384    fn test_queue_config_without_prefix() {
385        // Test that namespace works without prefix
386        let queue_name = "transaction_request_queue";
387
388        let config = Config::default().set_namespace(queue_name);
389        assert_eq!(config.get_namespace(), "transaction_request_queue");
390    }
391
392    #[test]
393    fn test_queue_config_default() {
394        let config = QueueConfig::default();
395
396        assert_eq!(config.enqueue_scheduled, Duration::from_secs(30));
397    }
398
399    #[test]
400    fn test_queue_config_high_throughput() {
401        let config = QueueConfig::high_frequency();
402
403        // High frequency should have faster enqueue_scheduled
404        assert_eq!(config.enqueue_scheduled, Duration::from_secs(2));
405
406        // Verify it's faster than default
407        let default = QueueConfig::default();
408        assert!(config.enqueue_scheduled < default.enqueue_scheduled);
409    }
410
411    #[test]
412    fn test_queue_config_low_frequency() {
413        let config = QueueConfig::low_frequency();
414
415        assert_eq!(config.enqueue_scheduled, Duration::from_secs(20));
416
417        // Low frequency should have longer enqueue_scheduled than high frequency
418        let high = QueueConfig::high_frequency();
419        assert!(config.enqueue_scheduled > high.enqueue_scheduled);
420    }
421}