openzeppelin_relayer/queues/redis/
queue.rs1use 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
27pub 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 pub transaction_status_queue: QueueStorage<Job<TransactionStatusCheck>>,
44 pub transaction_status_queue_evm: QueueStorage<Job<TransactionStatusCheck>>,
46 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_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#[derive(Clone, Debug)]
77struct QueueConfig {
78 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 fn high_frequency() -> Self {
93 Self {
94 enqueue_scheduled: Duration::from_secs(2),
95 }
96 }
97
98 fn low_frequency() -> Self {
102 Self {
103 enqueue_scheduled: Duration::from_secs(20),
104 }
105 }
106}
107
108impl Queue {
109 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 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 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 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 let queue_timeout = Duration::from_secs(5);
189
190 let max_age_ms = server_config.redis_connection_max_age_ms;
193
194 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 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 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(), ),
241 transaction_submission_queue: Self::storage(
242 &format!("{redis_key_prefix}transaction_submission_queue"),
243 conn_tx_submit,
244 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 notification_queue: Self::storage(
265 &format!("{redis_key_prefix}notification_queue"),
266 conn_notification,
267 low_frequency.clone(), ),
269 token_swap_request_queue: Self::storage(
270 &format!("{redis_key_prefix}token_swap_request_queue"),
271 conn_swap,
272 low_frequency.clone(), ),
274 relayer_health_check_queue: Self::storage(
275 &format!("{redis_key_prefix}relayer_health_check_queue"),
276 conn_health,
277 low_frequency.clone(), ),
279 redis_connections,
280 })
281 }
282
283 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 let namespace = "test_namespace";
300 let config = Config::default().set_namespace(namespace);
301
302 assert_eq!(config.get_namespace(), namespace);
303 }
304
305 #[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 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 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 assert_eq!(config.enqueue_scheduled, Duration::from_secs(2));
405
406 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 let high = QueueConfig::high_frequency();
419 assert!(config.enqueue_scheduled > high.enqueue_scheduled);
420 }
421}