openzeppelin_relayer/queues/redis/
worker.rs

1//! Redis/Apalis worker initialization.
2//!
3//! This module contains all Apalis-specific worker creation logic for the Redis
4//! queue backend, including WorkerBuilder configurations, Monitor setup,
5//! backoff strategies, and token swap cron workers.
6
7use 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
39/// Re-exports from [`tower::util`]
40pub 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
59// ---------------------------------------------------------------------------
60// Apalis adapter functions
61//
62// These thin adapters are the ONLY place where Apalis-specific handler types
63// (Data, Attempt, Worker<Context>, TaskId, RedisContext) appear. They convert
64// Apalis types → WorkerContext and HandlerError → apalis::prelude::Error,
65// keeping all handler business logic backend-neutral.
66// ---------------------------------------------------------------------------
67
68/// Sanity threshold (ms) for non-scheduled latency observations.
69///
70/// Latencies above this for a job whose only baseline was `Job.timestamp`
71/// almost certainly indicate clock skew between producer and consumer rather
72/// than a real backlog of that duration. We still observe the value (so the
73/// `+Inf` bucket reflects reality), but log a warning so operators can detect
74/// bad data instead of alerting on it.
75const PICKUP_LATENCY_CLOCK_SKEW_THRESHOLD_MS: i64 = 60 * 60 * 1000;
76
77/// Observe queue pickup latency for Redis/Apalis workers.
78///
79/// Uses `available_at` (the intended availability time) when present to exclude
80/// intentional scheduling delay. Falls back to `timestamp` (job creation time)
81/// for immediate jobs. Only records on the initial attempt — apalis
82/// `Attempt::current()` is 1-indexed, so `attempt == 1` is the first delivery;
83/// subsequent attempts would inflate the metric with retry backoff time.
84///
85/// If the chosen baseline fails to parse, the alternative is tried so a single
86/// corrupted field cannot silently drop the observation.
87fn 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";
309// Generic transaction status checker
310const TRANSACTION_STATUS_CHECKER: &str = "transaction_status_checker";
311// Network specific status checkers
312const 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
320/// Creates an exponential backoff with configurable parameters
321///
322/// # Arguments
323/// * `initial_ms` - Initial delay in milliseconds (e.g., 200)
324/// * `max_ms` - Maximum delay in milliseconds (e.g., 5000)
325/// * `jitter` - Jitter factor 0.0-1.0 (e.g., 0.99 for high jitter)
326///
327/// # Returns
328/// A configured backoff instance ready for use with RetryPolicy
329fn 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
344/// Initializes Redis/Apalis workers and starts the lifecycle monitor.
345///
346/// # Arguments
347/// * `app_state` - Application state containing the job producer and configuration
348/// * `shutdown_rx` - Backend-level shutdown signal (mirrors the SQS/PubSub `watch`
349///   pattern). Selected alongside SIGINT/SIGTERM so `QueueBackend::shutdown()` can
350///   stop this monitor on a programmatic shutdown, not just on an OS signal.
351pub 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    // Generic status checker
409    // Uses medium settings that work reasonably for most chains
410    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    // EVM status checker - slower retries to avoid premature resubmission
427    // EVM has longer block times (~12s) and needs time for resubmission logic
428    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    // Stellar status checker - fast retries for fast finality
445    // Stellar has sub-second finality, needs more frequent status checks
446    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)) // Default to 1 to avoid DB conflicts
506        .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    // Re-home the apalis Monitor onto the multi-thread pipeline runtime so its
576    // queue workers distribute across worker threads instead of pinning to the
577    // actix System arbiter's single thread. The JoinHandle is returned so the
578    // Monitor can be joined on graceful shutdown (drain in-flight work).
579    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
589/// Initializes swap workers for Solana and Stellar relayers.
590/// This function creates and registers workers for relayers that have swap enabled and cron schedule set.
591///
592/// `shutdown_rx` mirrors the SQS/PubSub `watch` shutdown pattern (see
593/// [`initialize_redis_workers`]).
594pub 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        // Create worker and add to the workers vector
671        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    // Register all workers with the monitor
702    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    // ── create_backoff tests ───────────────────────────────────────────
767
768    #[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            // Calling make_backoff() should not panic
861            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    // ── Backoff config invariant tests ─────────────────────────────────
875
876    #[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    // ── Worker name constant tests ─────────────────────────────────────
953
954    #[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        // The WorkerBuilder name for each queue-type-backed worker should match
1005        // the concurrency_env_key used with ServerConfig::get_worker_concurrency,
1006        // so that concurrency configuration picks up the correct env var.
1007        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    // ── monitor_handle_event tests ─────────────────────────────────────
1042
1043    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    // ── observe_redis_pickup_latency tests ─────────────────────────────
1086
1087    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        // available_at is "now" so latency should be ~0; job_timestamp is far in the
1119        // past, so if the fallback were used we'd see a large value. We verify the
1120        // preference indirectly via the histogram sum delta.
1121        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        // baseline in the future (producer clock ahead of consumer) → delta negative
1152        // → clamped to 0, but the observation still records.
1153        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}