1use actix_web::web::ThinData;
8
9use crate::{
10 config::ServerConfig,
11 constants::{
12 SYSTEM_CLEANUP_CRON_SCHEDULE, TRANSACTION_CLEANUP_CRON_SCHEDULE,
13 WORKER_SYSTEM_CLEANUP_RETRIES, WORKER_TOKEN_SWAP_REQUEST_RETRIES,
14 WORKER_TRANSACTION_CLEANUP_RETRIES,
15 },
16 jobs::{
17 notification_handler, relayer_health_check_handler, system_cleanup_handler,
18 token_swap_cron_handler, token_swap_request_handler, transaction_cleanup_handler,
19 transaction_request_handler, transaction_status_handler, transaction_submission_handler,
20 Job, JobProducerTrait, NotificationSend, RelayerHealthCheck, SystemCleanupCronReminder,
21 TokenSwapCronReminder, TokenSwapRequest, TransactionCleanupCronReminder,
22 TransactionRequest, TransactionSend, TransactionStatusCheck,
23 },
24 models::{
25 DefaultAppState, NetworkRepoModel, NotificationRepoModel, RelayerNetworkPolicy,
26 RelayerRepoModel, SignerRepoModel, ThinDataAppState, TransactionRepoModel,
27 },
28 repositories::{
29 ApiKeyRepositoryTrait, NetworkRepository, PluginRepositoryTrait, RelayerRepository,
30 Repository, TransactionCounterTrait, TransactionRepository,
31 },
32};
33use apalis::prelude::*;
34
35use apalis::layers::retry::backoff::MakeBackoff;
36use apalis::layers::retry::{backoff::ExponentialBackoffMaker, RetryPolicy};
37use apalis::layers::ErrorHandlingLayer;
38
39pub use tower::util::rng::HasherRng;
41
42use apalis_cron::CronStream;
43use eyre::Result;
44use std::{str::FromStr, time::Duration};
45use tokio::signal::unix::SignalKind;
46use tracing::{debug, error, info};
47
48use crate::metrics::observe_queue_pickup_latency;
49
50use super::{filter_relayers_for_swap, QueueType, WorkerContext};
51use crate::queues::retry_config::{
52 RetryBackoffConfig, NOTIFICATION_BACKOFF, RELAYER_HEALTH_BACKOFF, STATUS_EVM_BACKOFF,
53 STATUS_GENERIC_BACKOFF, STATUS_STELLAR_BACKOFF, SYSTEM_CLEANUP_BACKOFF,
54 TOKEN_SWAP_CRON_BACKOFF, TOKEN_SWAP_REQUEST_BACKOFF, TX_CLEANUP_BACKOFF, TX_REQUEST_BACKOFF,
55 TX_SUBMISSION_BACKOFF,
56};
57use crate::queues::WorkerHandle;
58
59const PICKUP_LATENCY_CLOCK_SKEW_THRESHOLD_MS: i64 = 60 * 60 * 1000;
76
77fn observe_redis_pickup_latency(
88 attempt: usize,
89 available_at: Option<&String>,
90 job_timestamp: &str,
91 queue_type: &str,
92) {
93 if attempt != 1 {
94 return;
95 }
96 let baseline_epoch_secs = available_at
97 .and_then(|s| s.parse::<i64>().ok())
98 .or_else(|| job_timestamp.parse::<i64>().ok());
99 let Some(baseline_epoch_secs) = baseline_epoch_secs else {
100 tracing::warn!(
101 queue_type = queue_type,
102 available_at = ?available_at,
103 job_timestamp = %job_timestamp,
104 "skipping queue_pickup_latency: failed to parse both available_at and job_timestamp"
105 );
106 return;
107 };
108 let now_ms = chrono::Utc::now().timestamp_millis();
109 let delta_ms = now_ms - baseline_epoch_secs * 1000;
110 if available_at.is_none() && delta_ms > PICKUP_LATENCY_CLOCK_SKEW_THRESHOLD_MS {
111 tracing::warn!(
112 queue_type = queue_type,
113 latency_ms = delta_ms,
114 "queue_pickup_latency above sanity threshold for non-scheduled job; check producer/consumer clock skew"
115 );
116 }
117 let latency_secs = delta_ms.max(0) as f64 / 1000.0;
118 observe_queue_pickup_latency(queue_type, "redis", latency_secs);
119}
120
121async fn apalis_transaction_request_handler(
122 job: Job<TransactionRequest>,
123 state: Data<ThinData<DefaultAppState>>,
124 attempt: Attempt,
125 task_id: TaskId,
126) -> Result<(), apalis::prelude::Error> {
127 observe_redis_pickup_latency(
128 attempt.current(),
129 job.available_at.as_ref(),
130 &job.timestamp,
131 QueueType::TransactionRequest.queue_name(),
132 );
133 let ctx = WorkerContext::new(attempt.current(), task_id.to_string());
134 transaction_request_handler(job, (*state).clone(), ctx)
135 .await
136 .map_err(Into::into)
137}
138
139async fn apalis_transaction_submission_handler(
140 job: Job<TransactionSend>,
141 state: Data<ThinData<DefaultAppState>>,
142 attempt: Attempt,
143 task_id: TaskId,
144) -> Result<(), apalis::prelude::Error> {
145 observe_redis_pickup_latency(
146 attempt.current(),
147 job.available_at.as_ref(),
148 &job.timestamp,
149 QueueType::TransactionSubmission.queue_name(),
150 );
151 let ctx = WorkerContext::new(attempt.current(), task_id.to_string());
152 transaction_submission_handler(job, (*state).clone(), ctx)
153 .await
154 .map_err(Into::into)
155}
156
157async fn apalis_transaction_status_handler(
158 job: Job<TransactionStatusCheck>,
159 state: Data<ThinData<DefaultAppState>>,
160 attempt: Attempt,
161 task_id: TaskId,
162) -> Result<(), apalis::prelude::Error> {
163 observe_redis_pickup_latency(
164 attempt.current(),
165 job.available_at.as_ref(),
166 &job.timestamp,
167 QueueType::StatusCheck.queue_name(),
168 );
169 let ctx = WorkerContext::new(attempt.current(), task_id.to_string());
170 transaction_status_handler(job, (*state).clone(), ctx)
171 .await
172 .map_err(Into::into)
173}
174
175async fn apalis_transaction_status_evm_handler(
176 job: Job<TransactionStatusCheck>,
177 state: Data<ThinData<DefaultAppState>>,
178 attempt: Attempt,
179 task_id: TaskId,
180) -> Result<(), apalis::prelude::Error> {
181 observe_redis_pickup_latency(
182 attempt.current(),
183 job.available_at.as_ref(),
184 &job.timestamp,
185 QueueType::StatusCheckEvm.queue_name(),
186 );
187 let ctx = WorkerContext::new(attempt.current(), task_id.to_string());
188 transaction_status_handler(job, (*state).clone(), ctx)
189 .await
190 .map_err(Into::into)
191}
192
193async fn apalis_transaction_status_stellar_handler(
194 job: Job<TransactionStatusCheck>,
195 state: Data<ThinData<DefaultAppState>>,
196 attempt: Attempt,
197 task_id: TaskId,
198) -> Result<(), apalis::prelude::Error> {
199 observe_redis_pickup_latency(
200 attempt.current(),
201 job.available_at.as_ref(),
202 &job.timestamp,
203 QueueType::StatusCheckStellar.queue_name(),
204 );
205 let ctx = WorkerContext::new(attempt.current(), task_id.to_string());
206 transaction_status_handler(job, (*state).clone(), ctx)
207 .await
208 .map_err(Into::into)
209}
210
211async fn apalis_notification_handler(
212 job: Job<NotificationSend>,
213 state: Data<ThinData<DefaultAppState>>,
214 attempt: Attempt,
215 task_id: TaskId,
216) -> Result<(), apalis::prelude::Error> {
217 observe_redis_pickup_latency(
218 attempt.current(),
219 job.available_at.as_ref(),
220 &job.timestamp,
221 QueueType::Notification.queue_name(),
222 );
223 let ctx = WorkerContext::new(attempt.current(), task_id.to_string());
224 notification_handler(job, (*state).clone(), ctx)
225 .await
226 .map_err(Into::into)
227}
228
229async fn apalis_token_swap_request_handler(
230 job: Job<TokenSwapRequest>,
231 state: Data<ThinData<DefaultAppState>>,
232 attempt: Attempt,
233 task_id: TaskId,
234) -> Result<(), apalis::prelude::Error> {
235 observe_redis_pickup_latency(
236 attempt.current(),
237 job.available_at.as_ref(),
238 &job.timestamp,
239 QueueType::TokenSwapRequest.queue_name(),
240 );
241 let ctx = WorkerContext::new(attempt.current(), task_id.to_string());
242 token_swap_request_handler(job, (*state).clone(), ctx)
243 .await
244 .map_err(Into::into)
245}
246
247async fn apalis_relayer_health_check_handler(
248 job: Job<RelayerHealthCheck>,
249 state: Data<ThinData<DefaultAppState>>,
250 attempt: Attempt,
251 task_id: TaskId,
252) -> Result<(), apalis::prelude::Error> {
253 observe_redis_pickup_latency(
254 attempt.current(),
255 job.available_at.as_ref(),
256 &job.timestamp,
257 QueueType::RelayerHealthCheck.queue_name(),
258 );
259 let ctx = WorkerContext::new(attempt.current(), task_id.to_string());
260 relayer_health_check_handler(job, (*state).clone(), ctx)
261 .await
262 .map_err(Into::into)
263}
264
265async fn apalis_transaction_cleanup_handler(
266 _job: TransactionCleanupCronReminder,
267 state: Data<ThinData<DefaultAppState>>,
268 attempt: Attempt,
269 task_id: TaskId,
270) -> Result<(), apalis::prelude::Error> {
271 let ctx = WorkerContext::new(attempt.current(), task_id.to_string());
272 transaction_cleanup_handler(TransactionCleanupCronReminder(), (*state).clone(), ctx)
273 .await
274 .map_err(Into::into)
275}
276
277async fn apalis_system_cleanup_handler(
278 _job: SystemCleanupCronReminder,
279 state: Data<ThinData<DefaultAppState>>,
280 attempt: Attempt,
281 task_id: TaskId,
282) -> Result<(), apalis::prelude::Error> {
283 let ctx = WorkerContext::new(attempt.current(), task_id.to_string());
284 system_cleanup_handler(SystemCleanupCronReminder(), (*state).clone(), ctx)
285 .await
286 .map_err(Into::into)
287}
288
289async fn apalis_token_swap_cron_handler(
290 _job: TokenSwapCronReminder,
291 relayer_id: Data<String>,
292 state: Data<ThinData<DefaultAppState>>,
293 attempt: Attempt,
294 task_id: TaskId,
295) -> Result<(), apalis::prelude::Error> {
296 let ctx = WorkerContext::new(attempt.current(), task_id.to_string());
297 token_swap_cron_handler(
298 TokenSwapCronReminder(),
299 (*relayer_id).clone(),
300 (*state).clone(),
301 ctx,
302 )
303 .await
304 .map_err(Into::into)
305}
306
307const TRANSACTION_REQUEST: &str = "transaction_request";
308const TRANSACTION_SENDER: &str = "transaction_sender";
309const TRANSACTION_STATUS_CHECKER: &str = "transaction_status_checker";
311const TRANSACTION_STATUS_CHECKER_EVM: &str = "transaction_status_checker_evm";
313const TRANSACTION_STATUS_CHECKER_STELLAR: &str = "transaction_status_checker_stellar";
314const NOTIFICATION_SENDER: &str = "notification_sender";
315const TOKEN_SWAP_REQUEST: &str = "token_swap_request";
316const TRANSACTION_CLEANUP: &str = "transaction_cleanup";
317const RELAYER_HEALTH_CHECK: &str = "relayer_health_check";
318const SYSTEM_CLEANUP: &str = "system_cleanup";
319
320fn create_backoff(initial_ms: u64, max_ms: u64, jitter: f64) -> Result<ExponentialBackoffMaker> {
330 let maker = ExponentialBackoffMaker::new(
331 Duration::from_millis(initial_ms),
332 Duration::from_millis(max_ms),
333 jitter,
334 HasherRng::default(),
335 )?;
336
337 Ok(maker)
338}
339
340fn create_backoff_from_config(cfg: RetryBackoffConfig) -> Result<ExponentialBackoffMaker> {
341 create_backoff(cfg.initial_ms, cfg.max_ms, cfg.jitter)
342}
343
344pub async fn initialize_redis_workers<J, RR, TR, NR, NFR, SR, TCR, PR, AKR>(
352 app_state: ThinDataAppState<J, RR, TR, NR, NFR, SR, TCR, PR, AKR>,
353 handle: tokio::runtime::Handle,
354 mut shutdown_rx: tokio::sync::watch::Receiver<bool>,
355) -> Result<WorkerHandle>
356where
357 J: JobProducerTrait + Send + Sync + 'static,
358 RR: RelayerRepository + Repository<RelayerRepoModel, String> + Send + Sync + 'static,
359 TR: TransactionRepository + Repository<TransactionRepoModel, String> + Send + Sync + 'static,
360 NR: NetworkRepository + Repository<NetworkRepoModel, String> + Send + Sync + 'static,
361 NFR: Repository<NotificationRepoModel, String> + Send + Sync + 'static,
362 SR: Repository<SignerRepoModel, String> + Send + Sync + 'static,
363 TCR: TransactionCounterTrait + Send + Sync + 'static,
364 PR: PluginRepositoryTrait + Send + Sync + 'static,
365 AKR: ApiKeyRepositoryTrait + Send + Sync + 'static,
366{
367 let queue_backend = app_state
368 .job_producer
369 .get_queue_backend()
370 .ok_or_else(|| eyre::eyre!("Queue backend is not available"))?;
371 let queue = queue_backend
372 .queue()
373 .cloned()
374 .ok_or_else(|| eyre::eyre!("Redis queue is not available for active backend"))?;
375
376 let transaction_request_queue_worker = WorkerBuilder::new(TRANSACTION_REQUEST)
377 .layer(ErrorHandlingLayer::new())
378 .retry(
379 RetryPolicy::retries(QueueType::TransactionRequest.max_retries())
380 .with_backoff(create_backoff_from_config(TX_REQUEST_BACKOFF)?.make_backoff()),
381 )
382 .enable_tracing()
383 .catch_panic()
384 .concurrency(ServerConfig::get_worker_concurrency(
385 QueueType::TransactionRequest.concurrency_env_key(),
386 QueueType::TransactionRequest.default_concurrency(),
387 ))
388 .data(app_state.clone())
389 .backend(queue.transaction_request_queue.clone())
390 .build_fn(apalis_transaction_request_handler);
391
392 let transaction_submission_queue_worker = WorkerBuilder::new(TRANSACTION_SENDER)
393 .layer(ErrorHandlingLayer::new())
394 .enable_tracing()
395 .catch_panic()
396 .retry(
397 RetryPolicy::retries(QueueType::TransactionSubmission.max_retries())
398 .with_backoff(create_backoff_from_config(TX_SUBMISSION_BACKOFF)?.make_backoff()),
399 )
400 .concurrency(ServerConfig::get_worker_concurrency(
401 QueueType::TransactionSubmission.concurrency_env_key(),
402 QueueType::TransactionSubmission.default_concurrency(),
403 ))
404 .data(app_state.clone())
405 .backend(queue.transaction_submission_queue.clone())
406 .build_fn(apalis_transaction_submission_handler);
407
408 let transaction_status_queue_worker = WorkerBuilder::new(TRANSACTION_STATUS_CHECKER)
411 .layer(ErrorHandlingLayer::new())
412 .enable_tracing()
413 .catch_panic()
414 .retry(
415 RetryPolicy::retries(QueueType::StatusCheck.max_retries())
416 .with_backoff(create_backoff_from_config(STATUS_GENERIC_BACKOFF)?.make_backoff()),
417 )
418 .concurrency(ServerConfig::get_worker_concurrency(
419 QueueType::StatusCheck.concurrency_env_key(),
420 QueueType::StatusCheck.default_concurrency(),
421 ))
422 .data(app_state.clone())
423 .backend(queue.transaction_status_queue.clone())
424 .build_fn(apalis_transaction_status_handler);
425
426 let transaction_status_queue_worker_evm = WorkerBuilder::new(TRANSACTION_STATUS_CHECKER_EVM)
429 .layer(ErrorHandlingLayer::new())
430 .enable_tracing()
431 .catch_panic()
432 .retry(
433 RetryPolicy::retries(QueueType::StatusCheck.max_retries())
434 .with_backoff(create_backoff_from_config(STATUS_EVM_BACKOFF)?.make_backoff()),
435 )
436 .concurrency(ServerConfig::get_worker_concurrency(
437 QueueType::StatusCheckEvm.concurrency_env_key(),
438 QueueType::StatusCheckEvm.default_concurrency(),
439 ))
440 .data(app_state.clone())
441 .backend(queue.transaction_status_queue_evm.clone())
442 .build_fn(apalis_transaction_status_evm_handler);
443
444 let transaction_status_queue_worker_stellar =
447 WorkerBuilder::new(TRANSACTION_STATUS_CHECKER_STELLAR)
448 .layer(ErrorHandlingLayer::new())
449 .enable_tracing()
450 .catch_panic()
451 .retry(
452 RetryPolicy::retries(QueueType::StatusCheckStellar.max_retries()).with_backoff(
453 create_backoff_from_config(STATUS_STELLAR_BACKOFF)?.make_backoff(),
454 ),
455 )
456 .concurrency(ServerConfig::get_worker_concurrency(
457 QueueType::StatusCheckStellar.concurrency_env_key(),
458 QueueType::StatusCheckStellar.default_concurrency(),
459 ))
460 .data(app_state.clone())
461 .backend(queue.transaction_status_queue_stellar.clone())
462 .build_fn(apalis_transaction_status_stellar_handler);
463
464 let notification_queue_worker = WorkerBuilder::new(NOTIFICATION_SENDER)
465 .layer(ErrorHandlingLayer::new())
466 .enable_tracing()
467 .catch_panic()
468 .retry(
469 RetryPolicy::retries(QueueType::Notification.max_retries())
470 .with_backoff(create_backoff_from_config(NOTIFICATION_BACKOFF)?.make_backoff()),
471 )
472 .concurrency(ServerConfig::get_worker_concurrency(
473 QueueType::Notification.concurrency_env_key(),
474 QueueType::Notification.default_concurrency(),
475 ))
476 .data(app_state.clone())
477 .backend(queue.notification_queue.clone())
478 .build_fn(apalis_notification_handler);
479
480 let token_swap_request_queue_worker = WorkerBuilder::new(TOKEN_SWAP_REQUEST)
481 .layer(ErrorHandlingLayer::new())
482 .enable_tracing()
483 .catch_panic()
484 .retry(
485 RetryPolicy::retries(QueueType::TokenSwapRequest.max_retries()).with_backoff(
486 create_backoff_from_config(TOKEN_SWAP_REQUEST_BACKOFF)?.make_backoff(),
487 ),
488 )
489 .concurrency(ServerConfig::get_worker_concurrency(
490 QueueType::TokenSwapRequest.concurrency_env_key(),
491 QueueType::TokenSwapRequest.default_concurrency(),
492 ))
493 .data(app_state.clone())
494 .backend(queue.token_swap_request_queue.clone())
495 .build_fn(apalis_token_swap_request_handler);
496
497 let transaction_cleanup_queue_worker = WorkerBuilder::new(TRANSACTION_CLEANUP)
498 .layer(ErrorHandlingLayer::new())
499 .enable_tracing()
500 .catch_panic()
501 .retry(
502 RetryPolicy::retries(WORKER_TRANSACTION_CLEANUP_RETRIES)
503 .with_backoff(create_backoff_from_config(TX_CLEANUP_BACKOFF)?.make_backoff()),
504 )
505 .concurrency(ServerConfig::get_worker_concurrency(TRANSACTION_CLEANUP, 1)) .data(app_state.clone())
507 .backend(CronStream::new(
508 apalis_cron::Schedule::from_str(TRANSACTION_CLEANUP_CRON_SCHEDULE)?,
509 ))
510 .build_fn(apalis_transaction_cleanup_handler);
511
512 let system_cleanup_queue_worker = WorkerBuilder::new(SYSTEM_CLEANUP)
513 .layer(ErrorHandlingLayer::new())
514 .enable_tracing()
515 .catch_panic()
516 .retry(
517 RetryPolicy::retries(WORKER_SYSTEM_CLEANUP_RETRIES)
518 .with_backoff(create_backoff_from_config(SYSTEM_CLEANUP_BACKOFF)?.make_backoff()),
519 )
520 .concurrency(1)
521 .data(app_state.clone())
522 .backend(CronStream::new(apalis_cron::Schedule::from_str(
523 SYSTEM_CLEANUP_CRON_SCHEDULE,
524 )?))
525 .build_fn(apalis_system_cleanup_handler);
526
527 let relayer_health_check_worker = WorkerBuilder::new(RELAYER_HEALTH_CHECK)
528 .layer(ErrorHandlingLayer::new())
529 .enable_tracing()
530 .catch_panic()
531 .retry(
532 RetryPolicy::retries(QueueType::RelayerHealthCheck.max_retries())
533 .with_backoff(create_backoff_from_config(RELAYER_HEALTH_BACKOFF)?.make_backoff()),
534 )
535 .concurrency(ServerConfig::get_worker_concurrency(
536 QueueType::RelayerHealthCheck.concurrency_env_key(),
537 QueueType::RelayerHealthCheck.default_concurrency(),
538 ))
539 .data(app_state.clone())
540 .backend(queue.relayer_health_check_queue.clone())
541 .build_fn(apalis_relayer_health_check_handler);
542
543 let monitor = Monitor::new()
544 .register(transaction_request_queue_worker)
545 .register(transaction_submission_queue_worker)
546 .register(transaction_status_queue_worker)
547 .register(transaction_status_queue_worker_evm)
548 .register(transaction_status_queue_worker_stellar)
549 .register(notification_queue_worker)
550 .register(token_swap_request_queue_worker)
551 .register(transaction_cleanup_queue_worker)
552 .register(system_cleanup_queue_worker)
553 .register(relayer_health_check_worker)
554 .on_event(monitor_handle_event)
555 .shutdown_timeout(Duration::from_millis(5000));
556
557 let monitor_future = monitor.run_with_signal(async move {
558 let mut sigint = tokio::signal::unix::signal(SignalKind::interrupt())
559 .map_err(|e| std::io::Error::other(format!("Failed to create SIGINT signal: {e}")))?;
560 let mut sigterm = tokio::signal::unix::signal(SignalKind::terminate())
561 .map_err(|e| std::io::Error::other(format!("Failed to create SIGTERM signal: {e}")))?;
562
563 debug!("Workers monitor started");
564
565 tokio::select! {
566 _ = sigint.recv() => debug!("Received SIGINT."),
567 _ = sigterm.recv() => debug!("Received SIGTERM."),
568 _ = shutdown_rx.changed() => debug!("Received programmatic shutdown signal."),
569 };
570
571 debug!("Workers monitor shutting down");
572
573 Ok(())
574 });
575 let monitor_handle = handle.spawn(async move {
580 if let Err(e) = monitor_future.await {
581 error!(error = %e, "monitor error");
582 }
583 });
584 debug!("Workers monitor spawned on pipeline runtime");
585
586 Ok(WorkerHandle::Tokio(monitor_handle))
587}
588
589pub async fn initialize_redis_token_swap_workers<J, RR, TR, NR, NFR, SR, TCR, PR, AKR>(
595 app_state: ThinDataAppState<J, RR, TR, NR, NFR, SR, TCR, PR, AKR>,
596 handle: tokio::runtime::Handle,
597 mut shutdown_rx: tokio::sync::watch::Receiver<bool>,
598) -> Result<Option<WorkerHandle>>
599where
600 J: JobProducerTrait + Send + Sync + 'static,
601 RR: RelayerRepository + Repository<RelayerRepoModel, String> + Send + Sync + 'static,
602 TR: TransactionRepository + Repository<TransactionRepoModel, String> + Send + Sync + 'static,
603 NR: NetworkRepository + Repository<NetworkRepoModel, String> + Send + Sync + 'static,
604 NFR: Repository<NotificationRepoModel, String> + Send + Sync + 'static,
605 SR: Repository<SignerRepoModel, String> + Send + Sync + 'static,
606 TCR: TransactionCounterTrait + Send + Sync + 'static,
607 PR: PluginRepositoryTrait + Send + Sync + 'static,
608 AKR: ApiKeyRepositoryTrait + Send + Sync + 'static,
609{
610 let active_relayers = app_state.relayer_repository.list_active().await?;
611 let relayers_with_swap_enabled = filter_relayers_for_swap(active_relayers);
612
613 if relayers_with_swap_enabled.is_empty() {
614 debug!("No relayers with swap enabled");
615 return Ok(None);
616 }
617 info!(
618 "Found {} relayers with swap enabled",
619 relayers_with_swap_enabled.len()
620 );
621
622 let mut workers = Vec::new();
623
624 let swap_backoff = create_backoff_from_config(TOKEN_SWAP_CRON_BACKOFF)?.make_backoff();
625
626 for relayer in relayers_with_swap_enabled {
627 debug!(relayer = ?relayer, "found relayer with swap enabled");
628
629 let (cron_schedule, network_type) = match &relayer.policies {
630 RelayerNetworkPolicy::Solana(policy) => match policy.get_swap_config() {
631 Some(config) => match config.cron_schedule {
632 Some(schedule) => (schedule, "solana".to_string()),
633 None => {
634 debug!(relayer_id = %relayer.id, "No cron schedule specified for Solana relayer; skipping");
635 continue;
636 }
637 },
638 None => {
639 debug!(relayer_id = %relayer.id, "No swap configuration specified for Solana relayer; skipping");
640 continue;
641 }
642 },
643 RelayerNetworkPolicy::Stellar(policy) => match policy.get_swap_config() {
644 Some(config) => match config.cron_schedule {
645 Some(schedule) => (schedule, "stellar".to_string()),
646 None => {
647 debug!(relayer_id = %relayer.id, "No cron schedule specified for Stellar relayer; skipping");
648 continue;
649 }
650 },
651 None => {
652 debug!(relayer_id = %relayer.id, "No swap configuration specified for Stellar relayer; skipping");
653 continue;
654 }
655 },
656 RelayerNetworkPolicy::Evm(_) => {
657 debug!(relayer_id = %relayer.id, "EVM relayers do not support swap; skipping");
658 continue;
659 }
660 };
661
662 let calendar_schedule = match apalis_cron::Schedule::from_str(&cron_schedule) {
663 Ok(schedule) => schedule,
664 Err(e) => {
665 error!(relayer_id = %relayer.id, error = %e, "Failed to parse cron schedule; skipping");
666 continue;
667 }
668 };
669
670 let worker = WorkerBuilder::new(format!(
672 "{}-swap-schedule-{}",
673 network_type,
674 relayer.id.clone()
675 ))
676 .layer(ErrorHandlingLayer::new())
677 .enable_tracing()
678 .catch_panic()
679 .retry(
680 RetryPolicy::retries(WORKER_TOKEN_SWAP_REQUEST_RETRIES)
681 .with_backoff(swap_backoff.clone()),
682 )
683 .concurrency(1)
684 .data(relayer.id.clone())
685 .data(app_state.clone())
686 .backend(CronStream::new(calendar_schedule))
687 .build_fn(apalis_token_swap_cron_handler);
688
689 workers.push(worker);
690 debug!(
691 relayer_id = %relayer.id,
692 network_type = %network_type,
693 "Created worker for relayer with swap enabled"
694 );
695 }
696
697 let mut monitor = Monitor::new()
698 .on_event(monitor_handle_event)
699 .shutdown_timeout(Duration::from_millis(5000));
700
701 for worker in workers {
703 monitor = monitor.register(worker);
704 }
705
706 let monitor_future = monitor.run_with_signal(async move {
707 let mut sigint = tokio::signal::unix::signal(SignalKind::interrupt())
708 .map_err(|e| std::io::Error::other(format!("Failed to create SIGINT signal: {e}")))?;
709 let mut sigterm = tokio::signal::unix::signal(SignalKind::terminate())
710 .map_err(|e| std::io::Error::other(format!("Failed to create SIGTERM signal: {e}")))?;
711
712 debug!("Swap Monitor started");
713
714 tokio::select! {
715 _ = sigint.recv() => debug!("Received SIGINT."),
716 _ = sigterm.recv() => debug!("Received SIGTERM."),
717 _ = shutdown_rx.changed() => debug!("Received programmatic shutdown signal."),
718 };
719
720 debug!("Swap Monitor shutting down");
721
722 Ok(())
723 });
724 let monitor_handle = handle.spawn(async move {
725 if let Err(e) = monitor_future.await {
726 error!(error = %e, "monitor error");
727 }
728 });
729 Ok(Some(WorkerHandle::Tokio(monitor_handle)))
730}
731
732fn monitor_handle_event(e: Worker<Event>) {
733 let worker_id = e.id();
734 match e.inner() {
735 Event::Engage(task_id) => {
736 debug!(worker_id = %worker_id, task_id = %task_id, "worker got a job");
737 }
738 Event::Error(e) => {
739 error!(worker_id = %worker_id, error = %e, "worker encountered an error");
740 }
741 Event::Exit => {
742 debug!(worker_id = %worker_id, "worker exited");
743 }
744 Event::Idle => {
745 debug!(worker_id = %worker_id, "worker is idle");
746 }
747 Event::Start => {
748 debug!(worker_id = %worker_id, "worker started");
749 }
750 Event::Stop => {
751 debug!(worker_id = %worker_id, "worker stopped");
752 }
753 _ => {}
754 }
755}
756
757#[cfg(test)]
758mod tests {
759 use super::*;
760 use crate::queues::retry_config::{
761 NOTIFICATION_BACKOFF, RELAYER_HEALTH_BACKOFF, STATUS_EVM_BACKOFF, STATUS_GENERIC_BACKOFF,
762 STATUS_STELLAR_BACKOFF, SYSTEM_CLEANUP_BACKOFF, TOKEN_SWAP_CRON_BACKOFF,
763 TOKEN_SWAP_REQUEST_BACKOFF, TX_CLEANUP_BACKOFF, TX_REQUEST_BACKOFF, TX_SUBMISSION_BACKOFF,
764 };
765
766 #[test]
769 fn test_create_backoff_with_valid_parameters() {
770 let result = create_backoff(200, 5000, 0.99);
771 assert!(
772 result.is_ok(),
773 "Should create backoff with valid parameters"
774 );
775 }
776
777 #[test]
778 fn test_create_backoff_with_zero_initial() {
779 let result = create_backoff(0, 5000, 0.99);
780 assert!(
781 result.is_ok(),
782 "Should handle zero initial delay (edge case)"
783 );
784 }
785
786 #[test]
787 fn test_create_backoff_with_equal_initial_and_max() {
788 let result = create_backoff(1000, 1000, 0.5);
789 assert!(result.is_ok(), "Should handle equal initial and max delays");
790 }
791
792 #[test]
793 fn test_create_backoff_with_zero_jitter() {
794 let result = create_backoff(500, 5000, 0.0);
795 assert!(result.is_ok(), "Should handle zero jitter");
796 }
797
798 #[test]
799 fn test_create_backoff_with_max_jitter() {
800 let result = create_backoff(500, 5000, 1.0);
801 assert!(result.is_ok(), "Should handle maximum jitter (1.0)");
802 }
803
804 #[test]
805 fn test_create_backoff_with_small_values() {
806 let result = create_backoff(1, 10, 0.5);
807 assert!(result.is_ok(), "Should handle very small delay values");
808 }
809
810 #[test]
811 fn test_create_backoff_with_large_values() {
812 let result = create_backoff(10000, 60000, 0.99);
813 assert!(result.is_ok(), "Should handle large delay values");
814 }
815
816 #[test]
817 fn test_create_backoff_from_config_profiles() {
818 let profiles = [
819 TX_REQUEST_BACKOFF,
820 TX_SUBMISSION_BACKOFF,
821 STATUS_GENERIC_BACKOFF,
822 STATUS_EVM_BACKOFF,
823 STATUS_STELLAR_BACKOFF,
824 NOTIFICATION_BACKOFF,
825 TOKEN_SWAP_REQUEST_BACKOFF,
826 TX_CLEANUP_BACKOFF,
827 SYSTEM_CLEANUP_BACKOFF,
828 RELAYER_HEALTH_BACKOFF,
829 TOKEN_SWAP_CRON_BACKOFF,
830 ];
831
832 for cfg in profiles {
833 let result = create_backoff_from_config(cfg);
834 assert!(
835 result.is_ok(),
836 "backoff profile should be constructible: {:?}",
837 cfg
838 );
839 }
840 }
841
842 #[test]
843 fn test_create_backoff_from_config_produces_usable_backoff() {
844 let profiles = [
845 TX_REQUEST_BACKOFF,
846 TX_SUBMISSION_BACKOFF,
847 STATUS_GENERIC_BACKOFF,
848 STATUS_EVM_BACKOFF,
849 STATUS_STELLAR_BACKOFF,
850 NOTIFICATION_BACKOFF,
851 TOKEN_SWAP_REQUEST_BACKOFF,
852 TX_CLEANUP_BACKOFF,
853 SYSTEM_CLEANUP_BACKOFF,
854 RELAYER_HEALTH_BACKOFF,
855 TOKEN_SWAP_CRON_BACKOFF,
856 ];
857
858 for cfg in profiles {
859 let mut maker = create_backoff_from_config(cfg).unwrap();
860 let _backoff = maker.make_backoff();
862 }
863 }
864
865 #[test]
866 fn test_create_backoff_with_initial_greater_than_max_errors() {
867 let result = create_backoff(10000, 100, 0.5);
868 assert!(
869 result.is_err(),
870 "initial > max should be rejected by ExponentialBackoffMaker"
871 );
872 }
873
874 #[test]
877 fn test_all_backoff_configs_have_valid_initial_le_max() {
878 let profiles: &[(&str, RetryBackoffConfig)] = &[
879 ("TX_REQUEST", TX_REQUEST_BACKOFF),
880 ("TX_SUBMISSION", TX_SUBMISSION_BACKOFF),
881 ("STATUS_GENERIC", STATUS_GENERIC_BACKOFF),
882 ("STATUS_EVM", STATUS_EVM_BACKOFF),
883 ("STATUS_STELLAR", STATUS_STELLAR_BACKOFF),
884 ("NOTIFICATION", NOTIFICATION_BACKOFF),
885 ("TOKEN_SWAP_REQUEST", TOKEN_SWAP_REQUEST_BACKOFF),
886 ("TX_CLEANUP", TX_CLEANUP_BACKOFF),
887 ("SYSTEM_CLEANUP", SYSTEM_CLEANUP_BACKOFF),
888 ("RELAYER_HEALTH", RELAYER_HEALTH_BACKOFF),
889 ("TOKEN_SWAP_CRON", TOKEN_SWAP_CRON_BACKOFF),
890 ];
891
892 for (name, cfg) in profiles {
893 assert!(
894 cfg.initial_ms <= cfg.max_ms,
895 "{name}: initial_ms ({}) must be <= max_ms ({})",
896 cfg.initial_ms,
897 cfg.max_ms
898 );
899 }
900 }
901
902 #[test]
903 fn test_all_backoff_configs_have_valid_jitter_range() {
904 let profiles: &[(&str, RetryBackoffConfig)] = &[
905 ("TX_REQUEST", TX_REQUEST_BACKOFF),
906 ("TX_SUBMISSION", TX_SUBMISSION_BACKOFF),
907 ("STATUS_GENERIC", STATUS_GENERIC_BACKOFF),
908 ("STATUS_EVM", STATUS_EVM_BACKOFF),
909 ("STATUS_STELLAR", STATUS_STELLAR_BACKOFF),
910 ("NOTIFICATION", NOTIFICATION_BACKOFF),
911 ("TOKEN_SWAP_REQUEST", TOKEN_SWAP_REQUEST_BACKOFF),
912 ("TX_CLEANUP", TX_CLEANUP_BACKOFF),
913 ("SYSTEM_CLEANUP", SYSTEM_CLEANUP_BACKOFF),
914 ("RELAYER_HEALTH", RELAYER_HEALTH_BACKOFF),
915 ("TOKEN_SWAP_CRON", TOKEN_SWAP_CRON_BACKOFF),
916 ];
917
918 for (name, cfg) in profiles {
919 assert!(
920 (0.0..=1.0).contains(&cfg.jitter),
921 "{name}: jitter ({}) must be in [0.0, 1.0]",
922 cfg.jitter
923 );
924 }
925 }
926
927 #[test]
928 fn test_all_backoff_configs_have_positive_initial_ms() {
929 let profiles: &[(&str, RetryBackoffConfig)] = &[
930 ("TX_REQUEST", TX_REQUEST_BACKOFF),
931 ("TX_SUBMISSION", TX_SUBMISSION_BACKOFF),
932 ("STATUS_GENERIC", STATUS_GENERIC_BACKOFF),
933 ("STATUS_EVM", STATUS_EVM_BACKOFF),
934 ("STATUS_STELLAR", STATUS_STELLAR_BACKOFF),
935 ("NOTIFICATION", NOTIFICATION_BACKOFF),
936 ("TOKEN_SWAP_REQUEST", TOKEN_SWAP_REQUEST_BACKOFF),
937 ("TX_CLEANUP", TX_CLEANUP_BACKOFF),
938 ("SYSTEM_CLEANUP", SYSTEM_CLEANUP_BACKOFF),
939 ("RELAYER_HEALTH", RELAYER_HEALTH_BACKOFF),
940 ("TOKEN_SWAP_CRON", TOKEN_SWAP_CRON_BACKOFF),
941 ];
942
943 for (name, cfg) in profiles {
944 assert!(
945 cfg.initial_ms > 0,
946 "{name}: initial_ms must be positive, got {}",
947 cfg.initial_ms
948 );
949 }
950 }
951
952 #[test]
955 fn test_worker_name_constants_are_nonempty() {
956 let names = [
957 TRANSACTION_REQUEST,
958 TRANSACTION_SENDER,
959 TRANSACTION_STATUS_CHECKER,
960 TRANSACTION_STATUS_CHECKER_EVM,
961 TRANSACTION_STATUS_CHECKER_STELLAR,
962 NOTIFICATION_SENDER,
963 TOKEN_SWAP_REQUEST,
964 TRANSACTION_CLEANUP,
965 RELAYER_HEALTH_CHECK,
966 SYSTEM_CLEANUP,
967 ];
968
969 for name in &names {
970 assert!(!name.is_empty(), "Worker name constant must not be empty");
971 }
972 }
973
974 #[test]
975 fn test_worker_name_constants_are_unique() {
976 let names = [
977 TRANSACTION_REQUEST,
978 TRANSACTION_SENDER,
979 TRANSACTION_STATUS_CHECKER,
980 TRANSACTION_STATUS_CHECKER_EVM,
981 TRANSACTION_STATUS_CHECKER_STELLAR,
982 NOTIFICATION_SENDER,
983 TOKEN_SWAP_REQUEST,
984 TRANSACTION_CLEANUP,
985 RELAYER_HEALTH_CHECK,
986 SYSTEM_CLEANUP,
987 ];
988
989 for (i, a) in names.iter().enumerate() {
990 for (j, b) in names.iter().enumerate() {
991 if i != j {
992 assert_ne!(
993 a, b,
994 "Worker names must be unique: '{}' at index {} and {}",
995 a, i, j
996 );
997 }
998 }
999 }
1000 }
1001
1002 #[test]
1003 fn test_worker_names_match_concurrency_env_keys() {
1004 assert_eq!(
1008 TRANSACTION_REQUEST,
1009 QueueType::TransactionRequest.concurrency_env_key()
1010 );
1011 assert_eq!(
1012 TRANSACTION_SENDER,
1013 QueueType::TransactionSubmission.concurrency_env_key()
1014 );
1015 assert_eq!(
1016 TRANSACTION_STATUS_CHECKER,
1017 QueueType::StatusCheck.concurrency_env_key()
1018 );
1019 assert_eq!(
1020 TRANSACTION_STATUS_CHECKER_EVM,
1021 QueueType::StatusCheckEvm.concurrency_env_key()
1022 );
1023 assert_eq!(
1024 TRANSACTION_STATUS_CHECKER_STELLAR,
1025 QueueType::StatusCheckStellar.concurrency_env_key()
1026 );
1027 assert_eq!(
1028 NOTIFICATION_SENDER,
1029 QueueType::Notification.concurrency_env_key()
1030 );
1031 assert_eq!(
1032 TOKEN_SWAP_REQUEST,
1033 QueueType::TokenSwapRequest.concurrency_env_key()
1034 );
1035 assert_eq!(
1036 RELAYER_HEALTH_CHECK,
1037 QueueType::RelayerHealthCheck.concurrency_env_key()
1038 );
1039 }
1040
1041 fn make_worker_event(event: Event) -> Worker<Event> {
1044 let worker_id = WorkerId::from_str("test-worker").unwrap();
1045 Worker::new(worker_id, event)
1046 }
1047
1048 #[test]
1049 fn test_monitor_handle_event_start_does_not_panic() {
1050 monitor_handle_event(make_worker_event(Event::Start));
1051 }
1052
1053 #[test]
1054 fn test_monitor_handle_event_engage_does_not_panic() {
1055 let task_id = TaskId::new();
1056 monitor_handle_event(make_worker_event(Event::Engage(task_id)));
1057 }
1058
1059 #[test]
1060 fn test_monitor_handle_event_idle_does_not_panic() {
1061 monitor_handle_event(make_worker_event(Event::Idle));
1062 }
1063
1064 #[test]
1065 fn test_monitor_handle_event_error_does_not_panic() {
1066 let error: Box<dyn std::error::Error + Send + Sync> = "test error".to_string().into();
1067 monitor_handle_event(make_worker_event(Event::Error(error)));
1068 }
1069
1070 #[test]
1071 fn test_monitor_handle_event_stop_does_not_panic() {
1072 monitor_handle_event(make_worker_event(Event::Stop));
1073 }
1074
1075 #[test]
1076 fn test_monitor_handle_event_exit_does_not_panic() {
1077 monitor_handle_event(make_worker_event(Event::Exit));
1078 }
1079
1080 #[test]
1081 fn test_monitor_handle_event_custom_does_not_panic() {
1082 monitor_handle_event(make_worker_event(Event::Custom("test-custom".to_string())));
1083 }
1084
1085 fn pickup_sample_count(queue_type: &str) -> u64 {
1088 crate::metrics::QUEUE_PICKUP_LATENCY
1089 .with_label_values(&[queue_type, "redis"])
1090 .get_sample_count()
1091 }
1092
1093 #[test]
1094 fn test_observe_redis_pickup_latency_records_on_first_attempt() {
1095 let queue = "test-pickup-first-attempt";
1096 let before = pickup_sample_count(queue);
1097
1098 let ts = chrono::Utc::now().timestamp().to_string();
1099 observe_redis_pickup_latency(1, None, &ts, queue);
1100
1101 assert_eq!(pickup_sample_count(queue), before + 1);
1102 }
1103
1104 #[test]
1105 fn test_observe_redis_pickup_latency_skips_retry_attempts() {
1106 let queue = "test-pickup-skip-retry";
1107 let before = pickup_sample_count(queue);
1108
1109 let ts = chrono::Utc::now().timestamp().to_string();
1110 observe_redis_pickup_latency(2, None, &ts, queue);
1111 observe_redis_pickup_latency(99, None, &ts, queue);
1112
1113 assert_eq!(pickup_sample_count(queue), before);
1114 }
1115
1116 #[test]
1117 fn test_observe_redis_pickup_latency_prefers_available_at_over_timestamp() {
1118 let queue = "test-pickup-prefers-available-at";
1122 let histogram = crate::metrics::QUEUE_PICKUP_LATENCY.with_label_values(&[queue, "redis"]);
1123 let sum_before = histogram.get_sample_sum();
1124
1125 let now = chrono::Utc::now().timestamp();
1126 let available_at = now.to_string();
1127 let stale_timestamp = (now - 3600).to_string();
1128
1129 observe_redis_pickup_latency(1, Some(&available_at), &stale_timestamp, queue);
1130
1131 let delta = histogram.get_sample_sum() - sum_before;
1132 assert!(
1133 delta < 5.0,
1134 "expected near-zero latency when available_at is now, got {delta}"
1135 );
1136 }
1137
1138 #[test]
1139 fn test_observe_redis_pickup_latency_falls_back_to_timestamp_when_available_at_absent() {
1140 let queue = "test-pickup-fallback-timestamp";
1141 let before = pickup_sample_count(queue);
1142
1143 let ts = chrono::Utc::now().timestamp().to_string();
1144 observe_redis_pickup_latency(1, None, &ts, queue);
1145
1146 assert_eq!(pickup_sample_count(queue), before + 1);
1147 }
1148
1149 #[test]
1150 fn test_observe_redis_pickup_latency_clamps_negative_skew_to_zero() {
1151 let queue = "test-pickup-clamps-negative";
1154 let histogram = crate::metrics::QUEUE_PICKUP_LATENCY.with_label_values(&[queue, "redis"]);
1155 let sum_before = histogram.get_sample_sum();
1156 let count_before = histogram.get_sample_count();
1157
1158 let future_ts = (chrono::Utc::now().timestamp() + 3600).to_string();
1159 observe_redis_pickup_latency(1, None, &future_ts, queue);
1160
1161 assert_eq!(histogram.get_sample_count(), count_before + 1);
1162 let delta_sum = histogram.get_sample_sum() - sum_before;
1163 assert!(
1164 delta_sum.abs() < f64::EPSILON,
1165 "expected 0 latency for future baseline, got {delta_sum}"
1166 );
1167 }
1168
1169 #[test]
1170 fn test_observe_redis_pickup_latency_skips_when_baseline_unparsable() {
1171 let queue = "test-pickup-unparsable";
1172 let before = pickup_sample_count(queue);
1173
1174 observe_redis_pickup_latency(1, None, "not-a-number", queue);
1175
1176 assert_eq!(pickup_sample_count(queue), before);
1177 }
1178}