openzeppelin_relayer/queues/redis/
backend.rs1use 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#[derive(Clone)]
32pub struct RedisBackend {
33 queue: Queue,
34 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 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 pub fn queue(&self) -> &Queue {
76 &self.queue
77 }
78}
79
80fn 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), 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 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 Ok(handles)
356 }
357
358 async fn health_check(&self) -> Result<Vec<QueueHealth>, QueueBackendError> {
359 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 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 let job = Job::new(
411 JobType::TransactionRequest,
412 TransactionRequest::new("tx1", "relayer1"),
413 );
414
415 assert!(!job.message_id.is_empty());
417 assert_eq!(job.message_id.len(), 36); 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 let now = std::time::SystemTime::now()
430 .duration_since(std::time::UNIX_EPOCH)
431 .unwrap()
432 .as_secs() as i64;
433
434 let immediate: Option<i64> = None;
436 assert_eq!(immediate, None);
437
438 let scheduled: Option<i64> = Some(now + 60);
440 assert!(scheduled.is_some());
441 assert!(scheduled.unwrap() > now);
442
443 let past: Option<i64> = Some(now - 10);
445 assert!(past.is_some());
446 }
447
448 #[test]
449 fn test_health_check_structure() {
450 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 assert_eq!(expected_queue_types.len(), 8);
466
467 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}