openzeppelin_relayer/queues/redis/
backend.rs

1//! Redis backend implementation using Apalis.
2//!
3//! This module provides a Redis/Apalis-backed implementation of the QueueBackend trait.
4//! It wraps the existing Queue structure and delegates to Apalis for job processing.
5
6use std::sync::Arc;
7
8use actix_web::web::ThinData;
9use apalis::prelude::Storage;
10use async_trait::async_trait;
11use tokio::sync::watch;
12use tracing::info;
13
14use crate::{
15    jobs::{
16        Job, NotificationSend, RelayerHealthCheck, TokenSwapRequest, TransactionRequest,
17        TransactionSend, TransactionStatusCheck,
18    },
19    models::{DefaultAppState, NetworkType},
20    queues::{Queue, QueueBackendType},
21    utils::RedisConnections,
22};
23
24use super::{QueueBackend, QueueBackendError, QueueHealth, QueueType, WorkerHandle};
25
26/// Redis backend using Apalis for job queue operations.
27///
28/// This is a wrapper around the existing Queue implementation that provides
29/// the QueueBackend trait interface. It delegates all operations to the
30/// existing Apalis/Redis infrastructure.
31#[derive(Clone)]
32pub struct RedisBackend {
33    queue: Queue,
34    /// Shutdown signal sender — mirrors the SQS/PubSub `watch` pattern. Sending
35    /// `true` is select!'d into the Apalis Monitor's `run_with_signal` future so
36    /// `QueueBackend::shutdown()` can stop the Redis monitors on a programmatic
37    /// shutdown, not just on a raw SIGINT/SIGTERM.
38    shutdown_tx: Arc<watch::Sender<bool>>,
39}
40
41impl std::fmt::Debug for RedisBackend {
42    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
43        f.debug_struct("RedisBackend")
44            .field("backend_type", &"redis")
45            .finish()
46    }
47}
48
49impl RedisBackend {
50    /// Creates a new Redis backend.
51    ///
52    /// This initializes all Redis-backed queues using the existing Queue::setup method.
53    ///
54    /// # Arguments
55    /// * `redis_connections` - Redis connection pools for queue operations
56    ///
57    /// # Errors
58    /// Returns QueueBackendError if queue setup fails
59    pub async fn new(redis_connections: Arc<RedisConnections>) -> Result<Self, QueueBackendError> {
60        info!("Initializing Redis queue backend");
61
62        let queue = Queue::setup(redis_connections)
63            .await
64            .map_err(|e| QueueBackendError::RedisError(e.to_string()))?;
65
66        let (shutdown_tx, _) = watch::channel(false);
67
68        Ok(Self {
69            queue,
70            shutdown_tx: Arc::new(shutdown_tx),
71        })
72    }
73
74    /// Returns a reference to the underlying Queue for compatibility with existing code.
75    pub fn queue(&self) -> &Queue {
76        &self.queue
77    }
78}
79
80/// Select status-check queue by network type.
81///
82/// EVM and Stellar use dedicated queues. Solana and unknown network types
83/// use the default status-check queue.
84fn status_check_queue_type(network_type: Option<&NetworkType>) -> QueueType {
85    match network_type {
86        Some(NetworkType::Evm) => QueueType::StatusCheckEvm,
87        Some(NetworkType::Stellar) => QueueType::StatusCheckStellar,
88        _ => QueueType::StatusCheck,
89    }
90}
91
92fn static_redis_health_statuses() -> Vec<QueueHealth> {
93    vec![
94        QueueHealth {
95            queue_type: QueueType::TransactionRequest,
96            messages_visible: Some(0), // Would need Redis LLEN query
97            messages_in_flight: 0,
98            messages_dlq: 0,
99            backend: "redis".to_string(),
100            is_healthy: true,
101        },
102        QueueHealth {
103            queue_type: QueueType::TransactionSubmission,
104            messages_visible: Some(0),
105            messages_in_flight: 0,
106            messages_dlq: 0,
107            backend: "redis".to_string(),
108            is_healthy: true,
109        },
110        QueueHealth {
111            queue_type: QueueType::StatusCheck,
112            messages_visible: Some(0),
113            messages_in_flight: 0,
114            messages_dlq: 0,
115            backend: "redis".to_string(),
116            is_healthy: true,
117        },
118        QueueHealth {
119            queue_type: QueueType::StatusCheckEvm,
120            messages_visible: Some(0),
121            messages_in_flight: 0,
122            messages_dlq: 0,
123            backend: "redis".to_string(),
124            is_healthy: true,
125        },
126        QueueHealth {
127            queue_type: QueueType::StatusCheckStellar,
128            messages_visible: Some(0),
129            messages_in_flight: 0,
130            messages_dlq: 0,
131            backend: "redis".to_string(),
132            is_healthy: true,
133        },
134        QueueHealth {
135            queue_type: QueueType::Notification,
136            messages_visible: Some(0),
137            messages_in_flight: 0,
138            messages_dlq: 0,
139            backend: "redis".to_string(),
140            is_healthy: true,
141        },
142        QueueHealth {
143            queue_type: QueueType::TokenSwapRequest,
144            messages_visible: Some(0),
145            messages_in_flight: 0,
146            messages_dlq: 0,
147            backend: "redis".to_string(),
148            is_healthy: true,
149        },
150        QueueHealth {
151            queue_type: QueueType::RelayerHealthCheck,
152            messages_visible: Some(0),
153            messages_in_flight: 0,
154            messages_dlq: 0,
155            backend: "redis".to_string(),
156            is_healthy: true,
157        },
158    ]
159}
160
161#[async_trait]
162impl QueueBackend for RedisBackend {
163    async fn produce_transaction_request(
164        &self,
165        job: Job<TransactionRequest>,
166        scheduled_on: Option<i64>,
167    ) -> Result<String, QueueBackendError> {
168        let mut storage = self.queue.transaction_request_queue.clone();
169        let job_id = job.message_id.clone();
170
171        match scheduled_on {
172            Some(on) => {
173                storage
174                    .schedule(job, on)
175                    .await
176                    .map_err(|e| QueueBackendError::RedisError(e.to_string()))?;
177            }
178            None => {
179                storage
180                    .push(job)
181                    .await
182                    .map_err(|e| QueueBackendError::RedisError(e.to_string()))?;
183            }
184        }
185
186        Ok(job_id)
187    }
188
189    async fn produce_transaction_submission(
190        &self,
191        job: Job<TransactionSend>,
192        scheduled_on: Option<i64>,
193    ) -> Result<String, QueueBackendError> {
194        let mut storage = self.queue.transaction_submission_queue.clone();
195        let job_id = job.message_id.clone();
196
197        match scheduled_on {
198            Some(on) => {
199                storage
200                    .schedule(job, on)
201                    .await
202                    .map_err(|e| QueueBackendError::RedisError(e.to_string()))?;
203            }
204            None => {
205                storage
206                    .push(job)
207                    .await
208                    .map_err(|e| QueueBackendError::RedisError(e.to_string()))?;
209            }
210        }
211
212        Ok(job_id)
213    }
214
215    async fn produce_transaction_status_check(
216        &self,
217        job: Job<TransactionStatusCheck>,
218        scheduled_on: Option<i64>,
219    ) -> Result<String, QueueBackendError> {
220        // Route by network_type to preserve existing Redis queue behavior.
221        let mut storage = match status_check_queue_type(job.data.network_type.as_ref()) {
222            QueueType::StatusCheckEvm => self.queue.transaction_status_queue_evm.clone(),
223            QueueType::StatusCheckStellar => self.queue.transaction_status_queue_stellar.clone(),
224            _ => self.queue.transaction_status_queue.clone(),
225        };
226        let job_id = job.message_id.clone();
227
228        match scheduled_on {
229            Some(on) => {
230                storage
231                    .schedule(job, on)
232                    .await
233                    .map_err(|e| QueueBackendError::RedisError(e.to_string()))?;
234            }
235            None => {
236                storage
237                    .push(job)
238                    .await
239                    .map_err(|e| QueueBackendError::RedisError(e.to_string()))?;
240            }
241        }
242
243        Ok(job_id)
244    }
245
246    async fn produce_notification(
247        &self,
248        job: Job<NotificationSend>,
249        scheduled_on: Option<i64>,
250    ) -> Result<String, QueueBackendError> {
251        let mut storage = self.queue.notification_queue.clone();
252        let job_id = job.message_id.clone();
253
254        match scheduled_on {
255            Some(on) => {
256                storage
257                    .schedule(job, on)
258                    .await
259                    .map_err(|e| QueueBackendError::RedisError(e.to_string()))?;
260            }
261            None => {
262                storage
263                    .push(job)
264                    .await
265                    .map_err(|e| QueueBackendError::RedisError(e.to_string()))?;
266            }
267        }
268
269        Ok(job_id)
270    }
271
272    async fn produce_token_swap_request(
273        &self,
274        job: Job<TokenSwapRequest>,
275        scheduled_on: Option<i64>,
276    ) -> Result<String, QueueBackendError> {
277        let mut storage = self.queue.token_swap_request_queue.clone();
278        let job_id = job.message_id.clone();
279
280        match scheduled_on {
281            Some(on) => {
282                storage
283                    .schedule(job, on)
284                    .await
285                    .map_err(|e| QueueBackendError::RedisError(e.to_string()))?;
286            }
287            None => {
288                storage
289                    .push(job)
290                    .await
291                    .map_err(|e| QueueBackendError::RedisError(e.to_string()))?;
292            }
293        }
294
295        Ok(job_id)
296    }
297
298    async fn produce_relayer_health_check(
299        &self,
300        job: Job<RelayerHealthCheck>,
301        scheduled_on: Option<i64>,
302    ) -> Result<String, QueueBackendError> {
303        let mut storage = self.queue.relayer_health_check_queue.clone();
304        let job_id = job.message_id.clone();
305
306        match scheduled_on {
307            Some(on) => {
308                storage
309                    .schedule(job, on)
310                    .await
311                    .map_err(|e| QueueBackendError::RedisError(e.to_string()))?;
312            }
313            None => {
314                storage
315                    .push(job)
316                    .await
317                    .map_err(|e| QueueBackendError::RedisError(e.to_string()))?;
318            }
319        }
320
321        Ok(job_id)
322    }
323
324    async fn initialize_workers(
325        &self,
326        app_state: Arc<ThinData<DefaultAppState>>,
327        handle: tokio::runtime::Handle,
328    ) -> Result<Vec<WorkerHandle>, QueueBackendError> {
329        info!("Initializing Redis backend workers");
330
331        let mut handles = Vec::new();
332
333        let monitor_handle = super::redis_worker::initialize_redis_workers(
334            (*app_state).clone(),
335            handle.clone(),
336            self.shutdown_tx.subscribe(),
337        )
338        .await
339        .map_err(|e| QueueBackendError::WorkerInitError(e.to_string()))?;
340        handles.push(monitor_handle);
341
342        if let Some(swap_handle) = super::redis_worker::initialize_redis_token_swap_workers(
343            (*app_state).clone(),
344            handle,
345            self.shutdown_tx.subscribe(),
346        )
347        .await
348        .map_err(|e| QueueBackendError::WorkerInitError(e.to_string()))?
349        {
350            handles.push(swap_handle);
351        }
352
353        // The returned handles are the apalis Monitor futures (re-homed onto the
354        // pipeline runtime); joining them on shutdown lets in-flight work drain.
355        Ok(handles)
356    }
357
358    async fn health_check(&self) -> Result<Vec<QueueHealth>, QueueBackendError> {
359        // Intentionally avoid per-request Redis queue depth calls here to keep
360        // health checks lightweight and avoid adding pressure to Redis.
361        // Return static backend health metadata only.
362        Ok(static_redis_health_statuses())
363    }
364
365    fn backend_type(&self) -> QueueBackendType {
366        QueueBackendType::Redis
367    }
368
369    fn shutdown(&self) {
370        info!("Redis backend: broadcasting shutdown signal to Apalis monitors");
371        let _ = self.shutdown_tx.send(true);
372    }
373}
374
375#[cfg(test)]
376mod tests {
377    use super::*;
378    use crate::jobs::{Job, JobType, TransactionRequest};
379    use crate::models::NetworkType;
380    use crate::queues::QueueType;
381
382    #[test]
383    fn test_backend_type_logic() {
384        // Test that backend_type returns Redis without requiring a Queue instance
385        // This tests the logic, not the actual implementation
386        assert_eq!(QueueBackendType::Redis.as_str(), "redis");
387        assert_eq!(QueueBackendType::Redis.to_string(), "redis");
388    }
389
390    #[test]
391    fn test_produce_transaction_status_check_routing_logic() {
392        assert_eq!(
393            status_check_queue_type(Some(&NetworkType::Evm)),
394            QueueType::StatusCheckEvm
395        );
396        assert_eq!(
397            status_check_queue_type(Some(&NetworkType::Stellar)),
398            QueueType::StatusCheckStellar
399        );
400        assert_eq!(
401            status_check_queue_type(Some(&NetworkType::Solana)),
402            QueueType::StatusCheck
403        );
404        assert_eq!(status_check_queue_type(None), QueueType::StatusCheck);
405    }
406
407    #[test]
408    fn test_job_id_extraction() {
409        // Test that job IDs are correctly extracted from jobs
410        let job = Job::new(
411            JobType::TransactionRequest,
412            TransactionRequest::new("tx1", "relayer1"),
413        );
414
415        // Verify job has a message_id
416        assert!(!job.message_id.is_empty());
417        assert_eq!(job.message_id.len(), 36); // UUID v4 format
418
419        // Verify job_id can be cloned (used in produce methods)
420        let job_id = job.message_id.clone();
421        assert_eq!(job_id, job.message_id);
422    }
423
424    #[test]
425    fn test_scheduled_on_handling() {
426        // Test that scheduled_on timestamps are handled correctly
427        // None means immediate execution, Some(timestamp) means scheduled
428
429        let now = std::time::SystemTime::now()
430            .duration_since(std::time::UNIX_EPOCH)
431            .unwrap()
432            .as_secs() as i64;
433
434        // Immediate execution
435        let immediate: Option<i64> = None;
436        assert_eq!(immediate, None);
437
438        // Scheduled execution
439        let scheduled: Option<i64> = Some(now + 60);
440        assert!(scheduled.is_some());
441        assert!(scheduled.unwrap() > now);
442
443        // Past timestamp should still be Some (handled by queue backend)
444        let past: Option<i64> = Some(now - 10);
445        assert!(past.is_some());
446    }
447
448    #[test]
449    fn test_health_check_structure() {
450        // Test the structure of health check responses without requiring Redis
451        // This verifies the expected format and fields
452
453        let expected_queue_types = vec![
454            QueueType::TransactionRequest,
455            QueueType::TransactionSubmission,
456            QueueType::StatusCheck,
457            QueueType::StatusCheckEvm,
458            QueueType::StatusCheckStellar,
459            QueueType::Notification,
460            QueueType::TokenSwapRequest,
461            QueueType::RelayerHealthCheck,
462        ];
463
464        // Verify all expected queue types exist
465        assert_eq!(expected_queue_types.len(), 8);
466
467        // Verify QueueType implements required traits
468        for queue_type in &expected_queue_types {
469            assert!(!queue_type.queue_name().is_empty());
470            assert!(!queue_type.redis_namespace().is_empty());
471        }
472
473        let statuses = static_redis_health_statuses();
474        assert_eq!(statuses.len(), expected_queue_types.len());
475        for queue_type in expected_queue_types {
476            assert!(statuses.iter().any(|h| h.queue_type == queue_type));
477        }
478        assert!(statuses.iter().all(|h| h.backend == "redis"));
479        assert!(statuses.iter().all(|h| h.is_healthy));
480    }
481
482    #[test]
483    fn test_static_redis_health_statuses_have_zero_counts() {
484        let statuses = static_redis_health_statuses();
485        assert!(!statuses.is_empty());
486        for status in statuses {
487            assert_eq!(status.messages_visible, Some(0));
488            assert_eq!(status.messages_in_flight, 0);
489            assert_eq!(status.messages_dlq, 0);
490        }
491    }
492}