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 low_frequency.clone(), ),
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 notification_queue: Self::storage(
263 &format!("{redis_key_prefix}notification_queue"),
264 conn_notification,
265 low_frequency.clone(), ),
267 token_swap_request_queue: Self::storage(
268 &format!("{redis_key_prefix}token_swap_request_queue"),
269 conn_swap,
270 low_frequency.clone(), ),
272 relayer_health_check_queue: Self::storage(
273 &format!("{redis_key_prefix}relayer_health_check_queue"),
274 conn_health,
275 low_frequency.clone(), ),
277 redis_connections,
278 })
279 }
280
281 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 let namespace = "test_namespace";
298 let config = Config::default().set_namespace(namespace);
299
300 assert_eq!(config.get_namespace(), namespace);
301 }
302
303 #[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 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 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 assert_eq!(config.enqueue_scheduled, Duration::from_secs(2));
403
404 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 let high = QueueConfig::high_frequency();
417 assert!(config.enqueue_scheduled > high.enqueue_scheduled);
418 }
419}