openzeppelin_relayer/queues/sqs/
backend.rs

1//! AWS SQS backend implementation.
2//!
3//! This module provides an AWS SQS-backed implementation of the QueueBackend trait.
4//! Supports both Standard and FIFO queues. By default (`SQS_QUEUE_TYPE=auto`),
5//! the queue type is auto-detected at startup by probing a reference queue.
6//! Can also be set explicitly to `standard` or `fifo`.
7
8use async_trait::async_trait;
9use aws_sdk_sqs::types::MessageAttributeValue;
10use std::collections::HashMap;
11use std::sync::Arc;
12use std::time::SystemTime;
13use tokio::sync::watch;
14use tracing::{debug, error, info, warn};
15
16use crate::{
17    config::ServerConfig,
18    jobs::{
19        Job, NotificationSend, RelayerHealthCheck, TokenSwapRequest, TransactionRequest,
20        TransactionSend, TransactionStatusCheck,
21    },
22    models::{DefaultAppState, NetworkType},
23    queues::QueueBackendType,
24    utils::{aws_error::DisplayErrorContext, classify_sdk_error},
25};
26use actix_web::web::ThinData;
27
28use super::{QueueBackend, QueueBackendError, QueueHealth, QueueType, WorkerHandle};
29
30/// SQS maximum message body size (256 KB).
31const SQS_MAX_MESSAGE_SIZE_BYTES: usize = 256 * 1024;
32
33/// Chooses the FIFO message group ID based on network type.
34///
35/// EVM requires per-relayer ordering (nonce management), so uses `relayer_id`.
36/// Non-EVM networks (Stellar, Solana) can safely parallelize per transaction,
37/// so uses `transaction_id` for better throughput.
38/// Falls back to `relayer_id` when network type is unknown (conservative/safe).
39fn transaction_message_group_id(
40    network_type: Option<&NetworkType>,
41    _relayer_id: &str,
42    transaction_id: &str,
43) -> String {
44    match network_type {
45        Some(_) | None => transaction_id.to_string(),
46    }
47}
48
49/// Selects the status-check queue for a given network type.
50///
51/// EVM and Stellar use dedicated queues; all other/unknown network types use
52/// the generic status-check queue.
53fn status_check_queue_type(network_type: Option<&NetworkType>) -> QueueType {
54    match network_type {
55        Some(NetworkType::Evm) => QueueType::StatusCheckEvm,
56        Some(NetworkType::Stellar) => QueueType::StatusCheckStellar,
57        _ => QueueType::StatusCheck,
58    }
59}
60
61/// AWS SQS backend for job queue operations.
62///
63/// Supports both Standard and FIFO queues (auto-detected at startup, or set via `SQS_QUEUE_TYPE`).
64/// FIFO mode provides message ordering and exactly-once delivery;
65/// Standard mode offers higher throughput and native per-message delays.
66#[derive(Clone)]
67pub struct SqsBackend {
68    /// AWS SQS client for all operations (send, delete, poll, change visibility)
69    sqs_client: aws_sdk_sqs::Client,
70    /// Mapping of queue types to SQS queue URLs
71    queue_urls: HashMap<QueueType, String>,
72    /// Cached DLQ URLs resolved once at startup, keyed by the source queue type.
73    /// Avoids repeated `get_queue_url` calls on every health check.
74    dlq_urls: HashMap<QueueType, String>,
75    /// AWS region
76    region: String,
77    /// Shutdown signal sender — sending `true` tells all workers and cron tasks to stop
78    shutdown_tx: Arc<watch::Sender<bool>>,
79}
80
81impl std::fmt::Debug for SqsBackend {
82    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
83        f.debug_struct("SqsBackend")
84            .field("backend_type", &"sqs")
85            .field("region", &self.region)
86            .field("queue_count", &self.queue_urls.len())
87            .finish()
88    }
89}
90
91/// Resolves the queue type from the configured `SQS_QUEUE_TYPE` value and
92/// probe results.
93///
94/// - `"standard"` / `"fifo"` → returns immediately (probes ignored).
95/// - `"auto"` → decides based on which probe succeeded.
96/// - anything else → error.
97///
98/// `probe_results` is `Option<(bool, bool)>`: `Some((standard_ok, fifo_ok))`
99/// when `sqs_queue_type == "auto"`, `None` otherwise.
100fn resolve_queue_type(
101    sqs_queue_type: &str,
102    probe_results: Option<(bool, bool)>,
103    ref_standard_url: &str,
104    ref_fifo_url: &str,
105) -> Result<bool, QueueBackendError> {
106    match sqs_queue_type {
107        "standard" => {
108            info!("Using explicit SQS queue type: standard");
109            Ok(false)
110        }
111        "fifo" => {
112            info!("Using explicit SQS queue type: fifo");
113            Ok(true)
114        }
115        "auto" => {
116            let (standard_exists, fifo_exists) = probe_results.unwrap_or((false, false));
117            match (standard_exists, fifo_exists) {
118                (true, false) => {
119                    info!("Detected SQS queue type: standard");
120                    Ok(false)
121                }
122                (false, true) => {
123                    info!("Detected SQS queue type: fifo");
124                    Ok(true)
125                }
126                (true, true) => Err(QueueBackendError::ConfigError(
127                    "Ambiguous SQS queue type: both standard and FIFO \
128                     'transaction-request' queues exist. Remove one set or set \
129                     SQS_QUEUE_TYPE explicitly."
130                        .to_string(),
131                )),
132                (false, false) => Err(QueueBackendError::ConfigError(format!(
133                    "No SQS queues found. Neither '{ref_standard_url}' nor \
134                     '{ref_fifo_url}' is accessible. Create queues before starting \
135                     the relayer, or set SQS_QUEUE_TYPE explicitly."
136                ))),
137            }
138        }
139        other => Err(QueueBackendError::ConfigError(format!(
140            "Unsupported SQS_QUEUE_TYPE: '{other}'. Must be 'auto', 'standard', or 'fifo'."
141        ))),
142    }
143}
144
145impl SqsBackend {
146    fn is_fifo_queue_url(queue_url: &str) -> bool {
147        queue_url.ends_with(".fifo")
148    }
149
150    /// Creates a new SQS backend.
151    ///
152    /// Loads AWS configuration from environment and builds queue URLs.
153    /// Queue type is determined by `SQS_QUEUE_TYPE`:
154    /// - `auto` (default): probes a reference queue at startup to detect the type
155    /// - `standard` / `fifo`: uses the specified type directly, skipping probing
156    ///
157    /// # Environment Variables
158    /// - `AWS_REGION` - AWS region (required)
159    /// - `SQS_QUEUE_URL_PREFIX` - Optional custom prefix
160    /// - `AWS_ACCOUNT_ID` - Required only when `SQS_QUEUE_URL_PREFIX` is not set
161    /// - `SQS_QUEUE_TYPE` - Queue type: `auto` (default), `standard`, or `fifo`
162    ///
163    /// # Errors
164    /// Returns ConfigError if required environment variables are missing or
165    /// if queue type cannot be determined (no queues found, or both types exist).
166    pub async fn new() -> Result<Self, QueueBackendError> {
167        info!("Initializing SQS queue backend");
168
169        // Load AWS config from environment
170        let config = aws_config::load_from_env().await;
171        let sqs_client = aws_sdk_sqs::Client::new(&config);
172        let region = config
173            .region()
174            .ok_or_else(|| {
175                QueueBackendError::ConfigError(
176                    "AWS_REGION not set. Required for SQS backend.".to_string(),
177                )
178            })?
179            .to_string();
180
181        // Build queue URL prefix.
182        // If an explicit prefix is provided, avoid forcing AWS_ACCOUNT_ID.
183        let prefix = match std::env::var("SQS_QUEUE_URL_PREFIX") {
184            Ok(prefix) => prefix,
185            Err(_) => {
186                let account_id =
187                    ServerConfig::get_aws_account_id().map_err(QueueBackendError::ConfigError)?;
188                format!("https://sqs.{region}.amazonaws.com/{account_id}/relayer-")
189            }
190        };
191        info!(
192            region = %region,
193            queue_url_prefix = %prefix,
194            "Resolved SQS queue URL prefix"
195        );
196
197        // Determine queue type: explicit override or auto-detect by probing.
198        let sqs_queue_type = ServerConfig::get_sqs_queue_type().to_lowercase();
199        let ref_standard_url = format!("{prefix}transaction-request");
200        let ref_fifo_url = format!("{prefix}transaction-request.fifo");
201
202        // Only probe when auto-detecting; explicit values skip the network call.
203        let probe_results = if sqs_queue_type == "auto" {
204            let (standard_probe, fifo_probe) = {
205                let client_s = sqs_client.clone();
206                let client_f = sqs_client.clone();
207                let url_s = ref_standard_url.clone();
208                let url_f = ref_fifo_url.clone();
209                tokio::join!(
210                    async move {
211                        client_s
212                            .get_queue_attributes()
213                            .queue_url(&url_s)
214                            .attribute_names(aws_sdk_sqs::types::QueueAttributeName::QueueArn)
215                            .send()
216                            .await
217                    },
218                    async move {
219                        client_f
220                            .get_queue_attributes()
221                            .queue_url(&url_f)
222                            .attribute_names(aws_sdk_sqs::types::QueueAttributeName::QueueArn)
223                            .send()
224                            .await
225                    }
226                )
227            };
228            Some((standard_probe.is_ok(), fifo_probe.is_ok()))
229        } else {
230            None
231        };
232
233        let is_fifo = resolve_queue_type(
234            &sqs_queue_type,
235            probe_results,
236            &ref_standard_url,
237            &ref_fifo_url,
238        )?;
239        let suffix = if is_fifo { ".fifo" } else { "" };
240
241        // Build queue URL mapping.
242        // Status checks use per-network queues (EVM, Stellar, generic/Solana)
243        // to match the Redis backend's separate worker setup with independent
244        // concurrency pools and network-tuned polling intervals.
245        let queue_urls = HashMap::from([
246            (
247                QueueType::TransactionRequest,
248                format!("{prefix}transaction-request{suffix}"),
249            ),
250            (
251                QueueType::TransactionSubmission,
252                format!("{prefix}transaction-submission{suffix}"),
253            ),
254            (
255                QueueType::StatusCheck,
256                format!("{prefix}status-check{suffix}"),
257            ),
258            (
259                QueueType::StatusCheckEvm,
260                format!("{prefix}status-check-evm{suffix}"),
261            ),
262            (
263                QueueType::StatusCheckStellar,
264                format!("{prefix}status-check-stellar{suffix}"),
265            ),
266            (
267                QueueType::Notification,
268                format!("{prefix}notification{suffix}"),
269            ),
270            (
271                QueueType::TokenSwapRequest,
272                format!("{prefix}token-swap-request{suffix}"),
273            ),
274            (
275                QueueType::RelayerHealthCheck,
276                format!("{prefix}relayer-health-check{suffix}"),
277            ),
278        ]);
279
280        // Fail fast at startup when expected queues are missing/misconfigured.
281        // This avoids silent runtime polling failures and makes infra drift explicit.
282        // Probe all queues concurrently to avoid slow sequential startup under
283        // SQS throttling during scale-out events.
284        let probe_futures: Vec<_> = queue_urls
285            .iter()
286            .map(|(queue_type, queue_url)| {
287                let client = sqs_client.clone();
288                let qt = *queue_type;
289                let url = queue_url.clone();
290                async move {
291                    debug!(
292                        queue_type = %qt,
293                        queue_url = %url,
294                        "Probing SQS queue accessibility at startup"
295                    );
296                    let probe = client
297                        .get_queue_attributes()
298                        .queue_url(&url)
299                        .attribute_names(aws_sdk_sqs::types::QueueAttributeName::QueueArn)
300                        .attribute_names(aws_sdk_sqs::types::QueueAttributeName::RedrivePolicy)
301                        .send()
302                        .await;
303                    (qt, url, probe)
304                }
305            })
306            .collect();
307
308        let probe_results = futures::future::join_all(probe_futures).await;
309
310        let mut missing_queues = Vec::new();
311        let mut dlq_urls: HashMap<QueueType, String> = HashMap::new();
312
313        for (queue_type, queue_url, probe) in probe_results {
314            match probe {
315                Ok(output) => {
316                    debug!(
317                        queue_type = %queue_type,
318                        queue_url = %queue_url,
319                        is_fifo = is_fifo,
320                        "SQS queue probe succeeded"
321                    );
322
323                    // Resolve and cache DLQ URL from the redrive policy while we
324                    // already have the attributes, avoiding per-health-check lookups.
325                    if let Some(dlq_url) =
326                        Self::resolve_dlq_url_from_attrs(&sqs_client, output.attributes()).await
327                    {
328                        dlq_urls.insert(queue_type, dlq_url);
329                    }
330                }
331                Err(err) => {
332                    // Include debug details because Display often collapses to "service error".
333                    error!(
334                        queue_type = %queue_type,
335                        queue_url = %queue_url,
336                        error = ?err,
337                        "SQS queue probe failed"
338                    );
339                    missing_queues.push(format!("{queue_type} ({queue_url}): {err:?}"));
340                }
341            }
342        }
343
344        if !missing_queues.is_empty() {
345            return Err(QueueBackendError::ConfigError(format!(
346                "SQS backend initialization failed. Missing/inaccessible queues: {}",
347                missing_queues.join(", ")
348            )));
349        }
350
351        info!(
352            region = %region,
353            queue_count = queue_urls.len(),
354            "SQS backend initialized"
355        );
356
357        let (shutdown_tx, _) = watch::channel(false);
358
359        Ok(Self {
360            sqs_client,
361            queue_urls,
362            dlq_urls,
363            region,
364            shutdown_tx: Arc::new(shutdown_tx),
365        })
366    }
367
368    /// Sends a message to SQS with FIFO parameters.
369    ///
370    /// # Arguments
371    /// * `queue_url` - SQS queue URL
372    /// * `body` - JSON-serialized job
373    /// * `message_group_id` - FIFO group ID (for ordering)
374    /// * `message_deduplication_id` - Deduplication ID (prevent duplicates)
375    /// * `delay_seconds` - Optional delay (0-900 seconds). Applied only for non-FIFO queues.
376    ///
377    /// # Returns
378    /// SQS message ID on success
379    async fn send_message_to_sqs(
380        &self,
381        queue_url: &str,
382        body: String,
383        message_group_id: String,
384        message_deduplication_id: String,
385        delay_seconds: Option<i32>,
386        target_scheduled_on: Option<i64>,
387    ) -> Result<String, QueueBackendError> {
388        if body.len() > SQS_MAX_MESSAGE_SIZE_BYTES {
389            return Err(QueueBackendError::SqsError(format!(
390                "Message body size ({} bytes) exceeds SQS limit ({} bytes)",
391                body.len(),
392                SQS_MAX_MESSAGE_SIZE_BYTES
393            )));
394        }
395
396        let mut request = self
397            .sqs_client
398            .send_message()
399            .queue_url(queue_url)
400            .message_body(body);
401
402        // FIFO queues require MessageGroupId and MessageDeduplicationId;
403        // standard queues reject these parameters.
404        if Self::is_fifo_queue_url(queue_url) {
405            request = request
406                .message_group_id(message_group_id)
407                .message_deduplication_id(message_deduplication_id);
408        }
409
410        if let Some(timestamp) = target_scheduled_on {
411            request = request.message_attributes(
412                "target_scheduled_on",
413                MessageAttributeValue::builder()
414                    .data_type("Number")
415                    .string_value(timestamp.to_string())
416                    .build()
417                    .map_err(|e| {
418                        QueueBackendError::SqsError(format!(
419                            "Failed to build scheduled-on attribute: {e}"
420                        ))
421                    })?,
422            );
423        }
424
425        // Add delay if specified (max 900 seconds = 15 minutes).
426        // FIFO queues do not support per-message DelaySeconds.
427        if let Some(delay) = delay_seconds {
428            let clamped_delay = delay.clamp(0, 900);
429            if Self::is_fifo_queue_url(queue_url) {
430                debug!(
431                    queue_url = %queue_url,
432                    requested_delay_seconds = delay,
433                    "Skipping per-message DelaySeconds for FIFO queue; worker-side scheduling will enforce target_scheduled_on"
434                );
435            } else {
436                request = request.delay_seconds(clamped_delay);
437                if delay != clamped_delay {
438                    warn!(
439                        requested = delay,
440                        clamped = clamped_delay,
441                        "Delay seconds clamped to SQS limit (0-900)"
442                    );
443                }
444            }
445        }
446
447        let response = request.send().await.map_err(|e| {
448            error!(
449                error.kind = classify_sdk_error(&e),
450                error.detail = %DisplayErrorContext(&e),
451                queue_url = %queue_url,
452                "Failed to send message to SQS"
453            );
454            QueueBackendError::SqsError(format!("SendMessage failed: {}", classify_sdk_error(&e)))
455        })?;
456
457        let message_id = response
458            .message_id()
459            .ok_or_else(|| QueueBackendError::SqsError("No message_id returned".to_string()))?
460            .to_string();
461
462        debug!(
463            message_id = %message_id,
464            queue_url = %queue_url,
465            "Message sent to SQS"
466        );
467
468        Ok(message_id)
469    }
470
471    /// Calculates delay in seconds from Unix timestamp.
472    ///
473    /// Returns None if scheduled_on is in the past or None.
474    fn calculate_delay_seconds(scheduled_on: Option<i64>) -> Option<i32> {
475        scheduled_on.and_then(|timestamp| {
476            let now = SystemTime::now()
477                .duration_since(SystemTime::UNIX_EPOCH)
478                .ok()?
479                .as_secs() as i64;
480
481            let delay = timestamp - now;
482            if delay > 0 {
483                Some(delay.min(900) as i32) // SQS max delay: 900 seconds
484            } else {
485                None // Already past scheduled time
486            }
487        })
488    }
489
490    /// Extracts the DLQ ARN from a redrive policy and resolves its queue URL.
491    ///
492    /// Called once at startup so the URL can be cached in `dlq_urls`.
493    async fn resolve_dlq_url_from_attrs(
494        sqs_client: &aws_sdk_sqs::Client,
495        attrs: Option<&HashMap<aws_sdk_sqs::types::QueueAttributeName, String>>,
496    ) -> Option<String> {
497        let redrive_policy =
498            attrs.and_then(|a| a.get(&aws_sdk_sqs::types::QueueAttributeName::RedrivePolicy))?;
499
500        let dlq_arn = serde_json::from_str::<serde_json::Value>(redrive_policy)
501            .ok()
502            .and_then(|v| v.get("deadLetterTargetArn").cloned())
503            .and_then(|v| v.as_str().map(|s| s.to_string()));
504
505        let dlq_name = dlq_arn.as_deref()?.rsplit(':').next()?;
506
507        match sqs_client.get_queue_url().queue_name(dlq_name).send().await {
508            Ok(output) => output.queue_url().map(str::to_string),
509            Err(err) => {
510                warn!(
511                    error.kind = classify_sdk_error(&err),
512                    error.detail = %DisplayErrorContext(&err),
513                    dlq_name = %dlq_name,
514                    "Failed to resolve DLQ URL at startup"
515                );
516                None
517            }
518        }
519    }
520
521    /// Returns the approximate message count for a cached DLQ URL.
522    ///
523    /// Uses URLs resolved and cached at startup, requiring only a single
524    /// `get_queue_attributes` call per health check (no URL resolution).
525    async fn get_dlq_message_count(&self, queue_type: &QueueType) -> u64 {
526        let Some(dlq_url) = self.dlq_urls.get(queue_type) else {
527            return 0;
528        };
529
530        match self
531            .sqs_client
532            .get_queue_attributes()
533            .queue_url(dlq_url)
534            .attribute_names(aws_sdk_sqs::types::QueueAttributeName::ApproximateNumberOfMessages)
535            .send()
536            .await
537        {
538            Ok(output) => output
539                .attributes()
540                .and_then(|attrs| {
541                    attrs.get(&aws_sdk_sqs::types::QueueAttributeName::ApproximateNumberOfMessages)
542                })
543                .and_then(|value| value.parse::<u64>().ok())
544                .unwrap_or(0),
545            Err(err) => {
546                warn!(
547                    error.kind = classify_sdk_error(&err),
548                    error.detail = %DisplayErrorContext(&err),
549                    dlq_url = %dlq_url,
550                    "Failed to fetch DLQ depth"
551                );
552                0
553            }
554        }
555    }
556}
557
558#[async_trait]
559impl QueueBackend for SqsBackend {
560    async fn produce_transaction_request(
561        &self,
562        job: Job<TransactionRequest>,
563        scheduled_on: Option<i64>,
564    ) -> Result<String, QueueBackendError> {
565        let queue_url = self
566            .queue_urls
567            .get(&QueueType::TransactionRequest)
568            .ok_or_else(|| QueueBackendError::QueueNotFound("TransactionRequest".to_string()))?;
569
570        let body = serde_json::to_string(&job).map_err(|e| {
571            error!(error = %e, "Failed to serialize TransactionRequest job");
572            QueueBackendError::SerializationError(e.to_string())
573        })?;
574
575        let message_group_id = transaction_message_group_id(
576            job.data.network_type.as_ref(),
577            &job.data.relayer_id,
578            &job.data.transaction_id,
579        );
580        let message_deduplication_id = job.message_id.clone();
581        let delay_seconds = Self::calculate_delay_seconds(scheduled_on);
582
583        self.send_message_to_sqs(
584            queue_url,
585            body,
586            message_group_id,
587            message_deduplication_id,
588            delay_seconds,
589            scheduled_on,
590        )
591        .await
592    }
593
594    async fn produce_transaction_submission(
595        &self,
596        job: Job<TransactionSend>,
597        scheduled_on: Option<i64>,
598    ) -> Result<String, QueueBackendError> {
599        let queue_url = self
600            .queue_urls
601            .get(&QueueType::TransactionSubmission)
602            .ok_or_else(|| QueueBackendError::QueueNotFound("TransactionSubmission".to_string()))?;
603
604        let body = serde_json::to_string(&job).map_err(|e| {
605            error!(error = %e, "Failed to serialize TransactionSend job");
606            QueueBackendError::SerializationError(e.to_string())
607        })?;
608
609        let message_group_id = transaction_message_group_id(
610            job.data.network_type.as_ref(),
611            &job.data.relayer_id,
612            &job.data.transaction_id,
613        );
614        let message_deduplication_id = job.message_id.clone();
615        let delay_seconds = Self::calculate_delay_seconds(scheduled_on);
616
617        self.send_message_to_sqs(
618            queue_url,
619            body,
620            message_group_id,
621            message_deduplication_id,
622            delay_seconds,
623            scheduled_on,
624        )
625        .await
626    }
627
628    async fn produce_transaction_status_check(
629        &self,
630        job: Job<TransactionStatusCheck>,
631        scheduled_on: Option<i64>,
632    ) -> Result<String, QueueBackendError> {
633        // Route to network-specific queue based on network type.
634        // EVM and Stellar get dedicated queues with tuned concurrency/polling;
635        // Solana and unknown network types use the generic StatusCheck queue.
636        let queue_type = status_check_queue_type(job.data.network_type.as_ref());
637        let queue_url = self
638            .queue_urls
639            .get(&queue_type)
640            .ok_or_else(|| QueueBackendError::QueueNotFound(format!("{queue_type}")))?;
641
642        let body = serde_json::to_string(&job).map_err(|e| {
643            error!(error = %e, "Failed to serialize TransactionStatusCheck job");
644            QueueBackendError::SerializationError(e.to_string())
645        })?;
646
647        let message_group_id = job.data.transaction_id.clone();
648        let message_deduplication_id = job.message_id.clone();
649        let delay_seconds = Self::calculate_delay_seconds(scheduled_on);
650
651        self.send_message_to_sqs(
652            queue_url,
653            body,
654            message_group_id,
655            message_deduplication_id,
656            delay_seconds,
657            scheduled_on,
658        )
659        .await
660    }
661
662    async fn produce_notification(
663        &self,
664        job: Job<NotificationSend>,
665        scheduled_on: Option<i64>,
666    ) -> Result<String, QueueBackendError> {
667        let queue_url = self
668            .queue_urls
669            .get(&QueueType::Notification)
670            .ok_or_else(|| QueueBackendError::QueueNotFound("Notification".to_string()))?;
671
672        let body = serde_json::to_string(&job).map_err(|e| {
673            error!(error = %e, "Failed to serialize NotificationSend job");
674            QueueBackendError::SerializationError(e.to_string())
675        })?;
676
677        // Notifications use notification_id as the group ID
678        let message_group_id = job.data.notification_id.clone();
679        let message_deduplication_id = job.message_id.clone();
680        let delay_seconds = Self::calculate_delay_seconds(scheduled_on);
681
682        self.send_message_to_sqs(
683            queue_url,
684            body,
685            message_group_id,
686            message_deduplication_id,
687            delay_seconds,
688            scheduled_on,
689        )
690        .await
691    }
692
693    async fn produce_token_swap_request(
694        &self,
695        job: Job<TokenSwapRequest>,
696        scheduled_on: Option<i64>,
697    ) -> Result<String, QueueBackendError> {
698        let queue_url = self
699            .queue_urls
700            .get(&QueueType::TokenSwapRequest)
701            .ok_or_else(|| QueueBackendError::QueueNotFound("TokenSwapRequest".to_string()))?;
702
703        let body = serde_json::to_string(&job).map_err(|e| {
704            error!(error = %e, "Failed to serialize TokenSwapRequest job");
705            QueueBackendError::SerializationError(e.to_string())
706        })?;
707
708        let message_group_id = job.data.relayer_id.clone();
709        let message_deduplication_id = job.message_id.clone();
710        let delay_seconds = Self::calculate_delay_seconds(scheduled_on);
711
712        self.send_message_to_sqs(
713            queue_url,
714            body,
715            message_group_id,
716            message_deduplication_id,
717            delay_seconds,
718            scheduled_on,
719        )
720        .await
721    }
722
723    async fn produce_relayer_health_check(
724        &self,
725        job: Job<RelayerHealthCheck>,
726        scheduled_on: Option<i64>,
727    ) -> Result<String, QueueBackendError> {
728        let queue_url = self
729            .queue_urls
730            .get(&QueueType::RelayerHealthCheck)
731            .ok_or_else(|| QueueBackendError::QueueNotFound("RelayerHealthCheck".to_string()))?;
732
733        let body = serde_json::to_string(&job).map_err(|e| {
734            error!(error = %e, "Failed to serialize RelayerHealthCheck job");
735            QueueBackendError::SerializationError(e.to_string())
736        })?;
737
738        let message_group_id = job.data.relayer_id.clone();
739        let message_deduplication_id = job.message_id.clone();
740        let delay_seconds = Self::calculate_delay_seconds(scheduled_on);
741
742        self.send_message_to_sqs(
743            queue_url,
744            body,
745            message_group_id,
746            message_deduplication_id,
747            delay_seconds,
748            scheduled_on,
749        )
750        .await
751    }
752
753    async fn initialize_workers(
754        &self,
755        app_state: Arc<ThinData<DefaultAppState>>,
756        handle: tokio::runtime::Handle,
757    ) -> Result<Vec<WorkerHandle>, QueueBackendError> {
758        info!(
759            "Initializing SQS workers for {} queues",
760            self.queue_urls.len()
761        );
762
763        let mut handles = Vec::new();
764
765        // Spawn a worker for each queue type onto the pipeline runtime.
766        for (queue_type, queue_url) in &self.queue_urls {
767            let worker_handle = super::sqs_worker::spawn_worker_for_queue(
768                self.sqs_client.clone(),
769                *queue_type,
770                queue_url.clone(),
771                app_state.clone(),
772                self.shutdown_tx.subscribe(),
773                handle.clone(),
774            )
775            .await?;
776
777            handles.push(worker_handle);
778        }
779
780        // Start cron scheduler for periodic tasks (cleanup, token swaps)
781        let cron_scheduler = super::sqs_cron::SqsCronScheduler::new(
782            app_state.clone(),
783            self.shutdown_tx.subscribe(),
784            handle.clone(),
785        );
786        let cron_handles = cron_scheduler.start().await?;
787        handles.extend(cron_handles);
788
789        // Internal shutdown signal handler — listens for SIGINT/SIGTERM and
790        // broadcasts shutdown to all SQS workers and cron tasks.
791        // (Redis/Apalis workers handle signals via their own Monitor.)
792        //
793        // NOT pushed into `handles`: this task only resolves on an OS signal, so on a
794        // programmatic/server-driven shutdown (no signal sent) it would never complete
795        // and `drain_worker_handles` would block for the full drain timeout waiting on
796        // it instead of the real worker/cron tasks.
797        {
798            let shutdown_tx = self.shutdown_tx.clone();
799            handle.spawn(async move {
800                let mut sigint =
801                    tokio::signal::unix::signal(tokio::signal::unix::SignalKind::interrupt())
802                        .expect("Failed to create SIGINT handler");
803                let mut sigterm =
804                    tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
805                        .expect("Failed to create SIGTERM handler");
806
807                tokio::select! {
808                    _ = sigint.recv() => info!("SQS backend: received SIGINT, shutting down workers"),
809                    _ = sigterm.recv() => info!("SQS backend: received SIGTERM, shutting down workers"),
810                }
811
812                let _ = shutdown_tx.send(true);
813            });
814        }
815
816        info!(
817            "Successfully spawned {} SQS workers and cron tasks",
818            handles.len()
819        );
820        Ok(handles)
821    }
822
823    async fn health_check(&self) -> Result<Vec<QueueHealth>, QueueBackendError> {
824        let mut health_statuses = Vec::new();
825
826        for (queue_type, queue_url) in &self.queue_urls {
827            // Get queue attributes to check health
828            let result = self
829                .sqs_client
830                .get_queue_attributes()
831                .queue_url(queue_url)
832                .attribute_names(
833                    aws_sdk_sqs::types::QueueAttributeName::ApproximateNumberOfMessages,
834                )
835                .attribute_names(
836                    aws_sdk_sqs::types::QueueAttributeName::ApproximateNumberOfMessagesNotVisible,
837                )
838                .send()
839                .await;
840
841            let (messages_visible, messages_in_flight, messages_dlq, is_healthy) = match result {
842                Ok(output) => {
843                    let attrs = output.attributes();
844                    let visible = attrs
845                        .and_then(|a| {
846                            a.get(&aws_sdk_sqs::types::QueueAttributeName::ApproximateNumberOfMessages)
847                        })
848                        .and_then(|v| v.parse::<u64>().ok())
849                        .unwrap_or(0);
850                    let in_flight = attrs
851                        .and_then(|a| {
852                            a.get(
853                                &aws_sdk_sqs::types::QueueAttributeName::ApproximateNumberOfMessagesNotVisible,
854                            )
855                        })
856                        .and_then(|v| v.parse::<u64>().ok())
857                        .unwrap_or(0);
858                    let dlq_count = self.get_dlq_message_count(queue_type).await;
859                    (visible, in_flight, dlq_count, true)
860                }
861                Err(e) => {
862                    error!(
863                        error = %e,
864                        queue_type = ?queue_type,
865                        "Failed to get queue attributes"
866                    );
867                    (0, 0, 0, false)
868                }
869            };
870
871            health_statuses.push(QueueHealth {
872                queue_type: *queue_type,
873                // SQS reports an approximate visible count (or 0 on probe error),
874                // both genuine values, so always `Some(..)`.
875                messages_visible: Some(messages_visible),
876                messages_in_flight,
877                messages_dlq,
878                backend: "sqs".to_string(),
879                is_healthy,
880            });
881        }
882
883        Ok(health_statuses)
884    }
885
886    fn backend_type(&self) -> QueueBackendType {
887        QueueBackendType::Sqs
888    }
889
890    fn shutdown(&self) {
891        info!("SQS backend: broadcasting shutdown signal to all workers");
892        let _ = self.shutdown_tx.send(true);
893    }
894}
895
896#[cfg(test)]
897mod tests {
898    use super::*;
899    use crate::jobs::{Job, JobType, TransactionStatusCheck};
900    use crate::models::NetworkType;
901
902    #[test]
903    fn test_calculate_delay_seconds() {
904        // No scheduled time
905        assert_eq!(SqsBackend::calculate_delay_seconds(None), None);
906
907        // Past time
908        let past = SystemTime::now()
909            .duration_since(SystemTime::UNIX_EPOCH)
910            .unwrap()
911            .as_secs() as i64
912            - 10;
913        assert_eq!(SqsBackend::calculate_delay_seconds(Some(past)), None);
914
915        // Future time within SQS limit (< 900s)
916        let future_5s = SystemTime::now()
917            .duration_since(SystemTime::UNIX_EPOCH)
918            .unwrap()
919            .as_secs() as i64
920            + 5;
921        assert_eq!(
922            SqsBackend::calculate_delay_seconds(Some(future_5s)),
923            Some(5)
924        );
925
926        // Future time beyond SQS limit (> 900s) - should clamp to 900
927        let future_1000s = SystemTime::now()
928            .duration_since(SystemTime::UNIX_EPOCH)
929            .unwrap()
930            .as_secs() as i64
931            + 1000;
932        assert_eq!(
933            SqsBackend::calculate_delay_seconds(Some(future_1000s)),
934            Some(900)
935        );
936    }
937
938    #[test]
939    fn test_calculate_delay_seconds_edge_cases() {
940        // Exactly at current time (should return None)
941        let now = SystemTime::now()
942            .duration_since(SystemTime::UNIX_EPOCH)
943            .unwrap()
944            .as_secs() as i64;
945        assert_eq!(SqsBackend::calculate_delay_seconds(Some(now)), None);
946
947        // Exactly at SQS limit (900s)
948        let future_900s = now + 900;
949        assert_eq!(
950            SqsBackend::calculate_delay_seconds(Some(future_900s)),
951            Some(900)
952        );
953
954        // Just over SQS limit (901s) - should clamp to 900
955        let future_901s = now + 901;
956        assert_eq!(
957            SqsBackend::calculate_delay_seconds(Some(future_901s)),
958            Some(900)
959        );
960    }
961
962    #[test]
963    fn test_sqs_backend_type_value() {
964        assert_eq!(QueueBackendType::Sqs.as_str(), "sqs");
965        assert_eq!(QueueBackendType::Sqs.to_string(), "sqs");
966    }
967
968    #[test]
969    fn test_queue_url_construction() {
970        // Test that queue URLs are correctly constructed
971        let mut queue_urls = HashMap::new();
972        let prefix = "https://sqs.us-east-1.amazonaws.com/123456789/relayer-";
973
974        queue_urls.insert(
975            QueueType::TransactionRequest,
976            format!("{prefix}transaction-request.fifo"),
977        );
978        queue_urls.insert(
979            QueueType::TransactionSubmission,
980            format!("{prefix}transaction-submission.fifo"),
981        );
982        queue_urls.insert(QueueType::StatusCheck, format!("{prefix}status-check.fifo"));
983        queue_urls.insert(
984            QueueType::Notification,
985            format!("{prefix}notification.fifo"),
986        );
987        queue_urls.insert(
988            QueueType::TokenSwapRequest,
989            format!("{prefix}token-swap-request.fifo"),
990        );
991        queue_urls.insert(
992            QueueType::RelayerHealthCheck,
993            format!("{prefix}relayer-health-check.fifo"),
994        );
995        queue_urls.insert(
996            QueueType::StatusCheckEvm,
997            format!("{prefix}status-check-evm.fifo"),
998        );
999        queue_urls.insert(
1000            QueueType::StatusCheckStellar,
1001            format!("{prefix}status-check-stellar.fifo"),
1002        );
1003
1004        // Verify all queue types have URLs
1005        assert_eq!(queue_urls.len(), 8);
1006        assert!(queue_urls
1007            .get(&QueueType::TransactionRequest)
1008            .unwrap()
1009            .ends_with(".fifo"));
1010        assert!(queue_urls
1011            .get(&QueueType::TransactionSubmission)
1012            .unwrap()
1013            .contains("transaction-submission"));
1014        assert!(queue_urls
1015            .get(&QueueType::StatusCheckEvm)
1016            .unwrap()
1017            .contains("status-check-evm"));
1018        assert!(queue_urls
1019            .get(&QueueType::StatusCheckStellar)
1020            .unwrap()
1021            .contains("status-check-stellar"));
1022    }
1023
1024    #[test]
1025    fn test_queue_url_construction_standard() {
1026        // Test that standard queue URLs do not have .fifo suffix
1027        let mut queue_urls = HashMap::new();
1028        let prefix = "https://sqs.us-east-1.amazonaws.com/123456789/relayer-";
1029
1030        queue_urls.insert(
1031            QueueType::TransactionRequest,
1032            format!("{prefix}transaction-request"),
1033        );
1034        queue_urls.insert(
1035            QueueType::TransactionSubmission,
1036            format!("{prefix}transaction-submission"),
1037        );
1038        queue_urls.insert(QueueType::StatusCheck, format!("{prefix}status-check"));
1039        queue_urls.insert(QueueType::Notification, format!("{prefix}notification"));
1040        queue_urls.insert(
1041            QueueType::TokenSwapRequest,
1042            format!("{prefix}token-swap-request"),
1043        );
1044        queue_urls.insert(
1045            QueueType::RelayerHealthCheck,
1046            format!("{prefix}relayer-health-check"),
1047        );
1048        queue_urls.insert(
1049            QueueType::StatusCheckEvm,
1050            format!("{prefix}status-check-evm"),
1051        );
1052        queue_urls.insert(
1053            QueueType::StatusCheckStellar,
1054            format!("{prefix}status-check-stellar"),
1055        );
1056
1057        assert_eq!(queue_urls.len(), 8);
1058        // Standard queue URLs should NOT end with .fifo
1059        for (_, url) in &queue_urls {
1060            assert!(
1061                !url.ends_with(".fifo"),
1062                "Standard queue URL should not end with .fifo: {url}"
1063            );
1064        }
1065        assert!(queue_urls
1066            .get(&QueueType::TransactionRequest)
1067            .unwrap()
1068            .contains("transaction-request"));
1069    }
1070
1071    #[test]
1072    fn test_is_fifo_queue_url_standard() {
1073        assert!(!SqsBackend::is_fifo_queue_url(
1074            "https://sqs.us-east-1.amazonaws.com/123/relayer-transaction-request"
1075        ));
1076        assert!(!SqsBackend::is_fifo_queue_url(
1077            "http://localstack:4566/000000000000/relayer-status-check"
1078        ));
1079    }
1080
1081    #[test]
1082    fn test_transaction_message_group_id_evm_uses_transaction() {
1083        let group = transaction_message_group_id(Some(&NetworkType::Evm), "relayer-1", "tx-123");
1084        assert_eq!(group, "tx-123");
1085    }
1086
1087    #[test]
1088    fn test_transaction_message_group_id_stellar_uses_transaction() {
1089        let group =
1090            transaction_message_group_id(Some(&NetworkType::Stellar), "relayer-1", "tx-123");
1091        assert_eq!(group, "tx-123");
1092    }
1093
1094    #[test]
1095    fn test_transaction_message_group_id_solana_uses_transaction() {
1096        let group = transaction_message_group_id(Some(&NetworkType::Solana), "relayer-1", "tx-123");
1097        assert_eq!(group, "tx-123");
1098    }
1099
1100    #[test]
1101    fn test_transaction_message_group_id_none_defaults_to_transaction() {
1102        let group = transaction_message_group_id(None, "relayer-1", "tx-123");
1103        assert_eq!(
1104            group, "tx-123",
1105            "Unknown network should default to transaction id"
1106        );
1107    }
1108
1109    #[test]
1110    fn test_status_check_queue_type_evm() {
1111        assert_eq!(
1112            status_check_queue_type(Some(&NetworkType::Evm)),
1113            QueueType::StatusCheckEvm
1114        );
1115    }
1116
1117    #[test]
1118    fn test_status_check_queue_type_stellar() {
1119        assert_eq!(
1120            status_check_queue_type(Some(&NetworkType::Stellar)),
1121            QueueType::StatusCheckStellar
1122        );
1123    }
1124
1125    #[test]
1126    fn test_status_check_queue_type_solana_defaults_to_generic() {
1127        assert_eq!(
1128            status_check_queue_type(Some(&NetworkType::Solana)),
1129            QueueType::StatusCheck
1130        );
1131    }
1132
1133    #[test]
1134    fn test_status_check_queue_type_none_defaults_to_generic() {
1135        assert_eq!(status_check_queue_type(None), QueueType::StatusCheck);
1136    }
1137
1138    #[test]
1139    fn test_sqs_max_message_size_constant() {
1140        assert_eq!(SQS_MAX_MESSAGE_SIZE_BYTES, 256 * 1024);
1141    }
1142
1143    #[test]
1144    fn test_is_fifo_queue_url() {
1145        assert!(SqsBackend::is_fifo_queue_url(
1146            "https://sqs.us-east-1.amazonaws.com/123/queue.fifo"
1147        ));
1148        assert!(!SqsBackend::is_fifo_queue_url(
1149            "https://sqs.us-east-1.amazonaws.com/123/queue"
1150        ));
1151    }
1152
1153    // --- resolve_queue_type tests ---
1154
1155    const REF_STD: &str = "http://localhost:4566/000000000000/relayer-transaction-request";
1156    const REF_FIFO: &str = "http://localhost:4566/000000000000/relayer-transaction-request.fifo";
1157
1158    #[test]
1159    fn test_resolve_queue_type_explicit_standard() {
1160        let result = resolve_queue_type("standard", None, REF_STD, REF_FIFO);
1161        assert_eq!(result.unwrap(), false);
1162    }
1163
1164    #[test]
1165    fn test_resolve_queue_type_explicit_fifo() {
1166        let result = resolve_queue_type("fifo", None, REF_STD, REF_FIFO);
1167        assert_eq!(result.unwrap(), true);
1168    }
1169
1170    #[test]
1171    fn test_resolve_queue_type_explicit_ignores_probes() {
1172        // Even if probes say FIFO exists, explicit "standard" wins
1173        let result = resolve_queue_type("standard", Some((false, true)), REF_STD, REF_FIFO);
1174        assert_eq!(result.unwrap(), false);
1175    }
1176
1177    #[test]
1178    fn test_resolve_queue_type_auto_standard_only() {
1179        let result = resolve_queue_type("auto", Some((true, false)), REF_STD, REF_FIFO);
1180        assert_eq!(result.unwrap(), false);
1181    }
1182
1183    #[test]
1184    fn test_resolve_queue_type_auto_fifo_only() {
1185        let result = resolve_queue_type("auto", Some((false, true)), REF_STD, REF_FIFO);
1186        assert_eq!(result.unwrap(), true);
1187    }
1188
1189    #[test]
1190    fn test_resolve_queue_type_auto_both_exist_errors() {
1191        let result = resolve_queue_type("auto", Some((true, true)), REF_STD, REF_FIFO);
1192        assert!(result.is_err());
1193        let err = result.unwrap_err().to_string();
1194        assert!(
1195            err.contains("Ambiguous"),
1196            "Expected 'Ambiguous' error, got: {err}"
1197        );
1198    }
1199
1200    #[test]
1201    fn test_resolve_queue_type_auto_neither_exists_errors() {
1202        let result = resolve_queue_type("auto", Some((false, false)), REF_STD, REF_FIFO);
1203        assert!(result.is_err());
1204        let err = result.unwrap_err().to_string();
1205        assert!(
1206            err.contains("No SQS queues found"),
1207            "Expected 'No SQS queues found' error, got: {err}"
1208        );
1209    }
1210
1211    #[test]
1212    fn test_resolve_queue_type_auto_no_probes_defaults_to_neither() {
1213        // None probe results (shouldn't happen in practice) treated as (false, false)
1214        let result = resolve_queue_type("auto", None, REF_STD, REF_FIFO);
1215        assert!(result.is_err());
1216    }
1217
1218    #[test]
1219    fn test_resolve_queue_type_unknown_value_errors() {
1220        let result = resolve_queue_type("invalid", None, REF_STD, REF_FIFO);
1221        assert!(result.is_err());
1222        let err = result.unwrap_err().to_string();
1223        assert!(
1224            err.contains("Unsupported SQS_QUEUE_TYPE"),
1225            "Expected unsupported error, got: {err}"
1226        );
1227    }
1228
1229    // ── resolve_queue_type: error variant and message checks ──────────
1230
1231    #[test]
1232    fn test_resolve_queue_type_auto_neither_error_includes_urls() {
1233        let std_url = "https://sqs.us-east-1.amazonaws.com/123/relayer-transaction-request";
1234        let fifo_url = "https://sqs.us-east-1.amazonaws.com/123/relayer-transaction-request.fifo";
1235        let result = resolve_queue_type("auto", Some((false, false)), std_url, fifo_url);
1236        let err = result.unwrap_err();
1237        let msg = err.to_string();
1238        assert!(
1239            msg.contains(std_url),
1240            "Error should include standard URL: {msg}"
1241        );
1242        assert!(
1243            msg.contains(fifo_url),
1244            "Error should include FIFO URL: {msg}"
1245        );
1246    }
1247
1248    #[test]
1249    fn test_resolve_queue_type_returns_config_error_variant() {
1250        let result = resolve_queue_type("invalid", None, REF_STD, REF_FIFO);
1251        assert!(
1252            matches!(result, Err(QueueBackendError::ConfigError(_))),
1253            "Expected ConfigError variant"
1254        );
1255
1256        let result = resolve_queue_type("auto", Some((true, true)), REF_STD, REF_FIFO);
1257        assert!(
1258            matches!(result, Err(QueueBackendError::ConfigError(_))),
1259            "Ambiguous case should be ConfigError"
1260        );
1261
1262        let result = resolve_queue_type("auto", Some((false, false)), REF_STD, REF_FIFO);
1263        assert!(
1264            matches!(result, Err(QueueBackendError::ConfigError(_))),
1265            "No queues case should be ConfigError"
1266        );
1267    }
1268
1269    #[test]
1270    fn test_resolve_queue_type_unknown_includes_value_in_error() {
1271        let result = resolve_queue_type("redis", None, REF_STD, REF_FIFO);
1272        let msg = result.unwrap_err().to_string();
1273        assert!(
1274            msg.contains("redis"),
1275            "Error should echo the invalid value: {msg}"
1276        );
1277
1278        let result = resolve_queue_type("", None, REF_STD, REF_FIFO);
1279        assert!(result.is_err(), "Empty string should be rejected");
1280    }
1281
1282    #[test]
1283    fn test_resolve_queue_type_case_sensitive() {
1284        // The function matches exact lowercase strings; mixed case is unsupported
1285        assert!(resolve_queue_type("Standard", None, REF_STD, REF_FIFO).is_err());
1286        assert!(resolve_queue_type("FIFO", None, REF_STD, REF_FIFO).is_err());
1287        assert!(resolve_queue_type("Auto", None, REF_STD, REF_FIFO).is_err());
1288    }
1289
1290    // ── calculate_delay_seconds: additional edge cases ────────────────
1291
1292    #[test]
1293    fn test_calculate_delay_seconds_one_second_future() {
1294        let future_1s = SystemTime::now()
1295            .duration_since(SystemTime::UNIX_EPOCH)
1296            .unwrap()
1297            .as_secs() as i64
1298            + 2; // +2 to avoid race with timing
1299        let result = SqsBackend::calculate_delay_seconds(Some(future_1s));
1300        assert!(result.is_some(), "1-2s in future should yield Some");
1301        assert!(result.unwrap() > 0, "Delay should be positive");
1302        assert!(result.unwrap() <= 2, "Delay should be at most 2s");
1303    }
1304
1305    #[test]
1306    fn test_calculate_delay_seconds_far_past() {
1307        // Unix epoch itself
1308        assert_eq!(SqsBackend::calculate_delay_seconds(Some(0)), None);
1309        // Negative timestamp (before epoch)
1310        assert_eq!(SqsBackend::calculate_delay_seconds(Some(-1000)), None);
1311    }
1312
1313    #[test]
1314    fn test_calculate_delay_seconds_very_far_future() {
1315        // Year ~2100 — should still clamp to 900
1316        let far_future = 4_102_444_800_i64; // 2100-01-01T00:00:00Z
1317        assert_eq!(
1318            SqsBackend::calculate_delay_seconds(Some(far_future)),
1319            Some(900)
1320        );
1321    }
1322
1323    #[test]
1324    fn test_calculate_delay_seconds_exactly_900_boundary() {
1325        let now = SystemTime::now()
1326            .duration_since(SystemTime::UNIX_EPOCH)
1327            .unwrap()
1328            .as_secs() as i64;
1329
1330        // 899s future → should return 899 (under the cap)
1331        let result = SqsBackend::calculate_delay_seconds(Some(now + 899));
1332        assert!(result.is_some());
1333        // Allow ±1 for timing
1334        let val = result.unwrap();
1335        assert!((898..=899).contains(&val), "Expected ~899, got {val}");
1336    }
1337
1338    // ── is_fifo_queue_url: edge cases ─────────────────────────────────
1339
1340    #[test]
1341    fn test_is_fifo_queue_url_empty() {
1342        assert!(!SqsBackend::is_fifo_queue_url(""));
1343    }
1344
1345    #[test]
1346    fn test_is_fifo_queue_url_case_sensitive() {
1347        assert!(!SqsBackend::is_fifo_queue_url(
1348            "https://sqs.us-east-1.amazonaws.com/123/queue.FIFO"
1349        ));
1350        assert!(!SqsBackend::is_fifo_queue_url(
1351            "https://sqs.us-east-1.amazonaws.com/123/queue.Fifo"
1352        ));
1353    }
1354
1355    #[test]
1356    fn test_is_fifo_queue_url_fifo_in_middle() {
1357        assert!(!SqsBackend::is_fifo_queue_url(
1358            "https://sqs.us-east-1.amazonaws.com/123/.fifo/queue"
1359        ));
1360    }
1361
1362    #[test]
1363    fn test_is_fifo_queue_url_just_suffix() {
1364        assert!(SqsBackend::is_fifo_queue_url(".fifo"));
1365    }
1366
1367    #[test]
1368    fn test_is_fifo_queue_url_localstack() {
1369        assert!(SqsBackend::is_fifo_queue_url(
1370            "http://localhost:4566/000000000000/relayer-tx.fifo"
1371        ));
1372    }
1373
1374    // ── transaction_message_group_id: consistency ─────────────────────
1375
1376    #[test]
1377    fn test_transaction_message_group_id_always_returns_transaction_id() {
1378        // Current implementation returns transaction_id for all network types.
1379        // This test documents that behavior and catches accidental changes.
1380        let networks: &[Option<NetworkType>] = &[
1381            Some(NetworkType::Evm),
1382            Some(NetworkType::Stellar),
1383            Some(NetworkType::Solana),
1384            None,
1385        ];
1386
1387        for network in networks {
1388            let group = transaction_message_group_id(network.as_ref(), "relayer-99", "tx-abc");
1389            assert_eq!(
1390                group, "tx-abc",
1391                "Expected transaction_id for network {network:?}"
1392            );
1393        }
1394    }
1395
1396    // ── status_check_queue_type: returned queue names ──────────────────
1397
1398    #[test]
1399    fn test_status_check_queue_type_returns_distinct_queue_names() {
1400        let evm = status_check_queue_type(Some(&NetworkType::Evm));
1401        let stellar = status_check_queue_type(Some(&NetworkType::Stellar));
1402        let generic = status_check_queue_type(None);
1403
1404        assert_ne!(evm.queue_name(), stellar.queue_name());
1405        assert_ne!(evm.queue_name(), generic.queue_name());
1406        assert_ne!(stellar.queue_name(), generic.queue_name());
1407    }
1408
1409    #[test]
1410    fn test_status_check_queue_type_all_are_status_checks() {
1411        let networks: &[Option<&NetworkType>] = &[
1412            Some(&NetworkType::Evm),
1413            Some(&NetworkType::Stellar),
1414            Some(&NetworkType::Solana),
1415            None,
1416        ];
1417
1418        for network in networks {
1419            let qt = status_check_queue_type(*network);
1420            assert!(
1421                qt.is_status_check(),
1422                "{qt:?} should be a status check variant"
1423            );
1424        }
1425    }
1426
1427    // ── Queue URL construction algorithm ──────────────────────────────
1428
1429    #[test]
1430    fn test_queue_url_construction_algorithm_fifo() {
1431        // Replicate the algorithm from SqsBackend::new()
1432        let prefix = "https://sqs.us-east-1.amazonaws.com/123456789/relayer-";
1433        let suffix = ".fifo";
1434
1435        let expected_urls = [
1436            (
1437                QueueType::TransactionRequest,
1438                format!("{prefix}transaction-request{suffix}"),
1439            ),
1440            (
1441                QueueType::TransactionSubmission,
1442                format!("{prefix}transaction-submission{suffix}"),
1443            ),
1444            (
1445                QueueType::StatusCheck,
1446                format!("{prefix}status-check{suffix}"),
1447            ),
1448            (
1449                QueueType::StatusCheckEvm,
1450                format!("{prefix}status-check-evm{suffix}"),
1451            ),
1452            (
1453                QueueType::StatusCheckStellar,
1454                format!("{prefix}status-check-stellar{suffix}"),
1455            ),
1456            (
1457                QueueType::Notification,
1458                format!("{prefix}notification{suffix}"),
1459            ),
1460            (
1461                QueueType::TokenSwapRequest,
1462                format!("{prefix}token-swap-request{suffix}"),
1463            ),
1464            (
1465                QueueType::RelayerHealthCheck,
1466                format!("{prefix}relayer-health-check{suffix}"),
1467            ),
1468        ];
1469
1470        for (qt, expected) in &expected_urls {
1471            assert!(
1472                SqsBackend::is_fifo_queue_url(expected),
1473                "{qt:?}: URL should be FIFO: {expected}"
1474            );
1475        }
1476    }
1477
1478    #[test]
1479    fn test_queue_url_construction_algorithm_standard() {
1480        let prefix = "https://sqs.us-east-1.amazonaws.com/123456789/relayer-";
1481        let suffix = "";
1482
1483        let urls = [
1484            format!("{prefix}transaction-request{suffix}"),
1485            format!("{prefix}transaction-submission{suffix}"),
1486            format!("{prefix}status-check{suffix}"),
1487            format!("{prefix}notification{suffix}"),
1488        ];
1489
1490        for url in &urls {
1491            assert!(
1492                !SqsBackend::is_fifo_queue_url(url),
1493                "Standard URL should not be FIFO: {url}"
1494            );
1495        }
1496    }
1497
1498    // ── SQS_MAX_MESSAGE_SIZE_BYTES ────────────────────────────────────
1499
1500    #[test]
1501    fn test_sqs_max_message_size_matches_aws_limit() {
1502        // AWS SQS maximum message body size is exactly 256 KiB
1503        assert_eq!(SQS_MAX_MESSAGE_SIZE_BYTES, 262_144);
1504    }
1505
1506    // ── resolve_queue_type: all branches produce expected is_fifo ─────
1507
1508    #[test]
1509    fn test_resolve_queue_type_suffix_logic() {
1510        // Verify the suffix derived from resolve_queue_type produces correct URLs
1511        let test_cases = [
1512            ("standard", None, false),
1513            ("fifo", None, true),
1514            ("auto", Some((true, false)), false),
1515            ("auto", Some((false, true)), true),
1516        ];
1517
1518        for (sqs_type, probes, expected_fifo) in test_cases {
1519            let is_fifo = resolve_queue_type(sqs_type, probes, REF_STD, REF_FIFO).unwrap();
1520            assert_eq!(is_fifo, expected_fifo, "sqs_type={sqs_type}");
1521
1522            let suffix = if is_fifo { ".fifo" } else { "" };
1523            let url = format!("https://sqs.us-east-1.amazonaws.com/123/relayer-tx{suffix}");
1524            assert_eq!(
1525                SqsBackend::is_fifo_queue_url(&url),
1526                expected_fifo,
1527                "URL FIFO detection mismatch for sqs_type={sqs_type}"
1528            );
1529        }
1530    }
1531
1532    #[tokio::test]
1533    #[ignore]
1534    async fn smoke_push_status_check_to_sqs() {
1535        // Requires real AWS credentials and queue env config.
1536        // Expected env:
1537        // - AWS_REGION
1538        // - SQS_QUEUE_URL_PREFIX
1539        // - optional AWS_ACCESS_KEY_ID/AWS_SECRET_ACCESS_KEY/AWS_SESSION_TOKEN
1540        let backend = SqsBackend::new()
1541            .await
1542            .expect("SQS backend initialization failed");
1543        let job = Job::new(
1544            JobType::TransactionStatusCheck,
1545            TransactionStatusCheck::new("smoke-tx-id", "smoke-relayer", NetworkType::Stellar),
1546        );
1547        let now = SystemTime::now()
1548            .duration_since(SystemTime::UNIX_EPOCH)
1549            .expect("system time before unix epoch")
1550            .as_secs() as i64;
1551        let scheduled_on = Some(now + 2);
1552        let result = backend
1553            .produce_transaction_status_check(job, scheduled_on)
1554            .await;
1555        assert!(
1556            result.is_ok(),
1557            "Expected SendMessage via SQS backend to succeed, got: {result:?}"
1558        );
1559    }
1560}