1use 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
30const SQS_MAX_MESSAGE_SIZE_BYTES: usize = 256 * 1024;
32
33fn 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
49fn 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#[derive(Clone)]
67pub struct SqsBackend {
68 sqs_client: aws_sdk_sqs::Client,
70 queue_urls: HashMap<QueueType, String>,
72 dlq_urls: HashMap<QueueType, String>,
75 region: String,
77 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
91fn 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 pub async fn new() -> Result<Self, QueueBackendError> {
167 info!("Initializing SQS queue backend");
168
169 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 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 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 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 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 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 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 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 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 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 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 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) } else {
485 None }
487 })
488 }
489
490 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 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 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 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 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 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 {
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 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 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 assert_eq!(SqsBackend::calculate_delay_seconds(None), None);
906
907 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 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 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 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 let future_900s = now + 900;
949 assert_eq!(
950 SqsBackend::calculate_delay_seconds(Some(future_900s)),
951 Some(900)
952 );
953
954 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 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 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 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 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 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 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 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 #[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 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 #[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; 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 assert_eq!(SqsBackend::calculate_delay_seconds(Some(0)), None);
1309 assert_eq!(SqsBackend::calculate_delay_seconds(Some(-1000)), None);
1311 }
1312
1313 #[test]
1314 fn test_calculate_delay_seconds_very_far_future() {
1315 let far_future = 4_102_444_800_i64; 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 let result = SqsBackend::calculate_delay_seconds(Some(now + 899));
1332 assert!(result.is_some());
1333 let val = result.unwrap();
1335 assert!((898..=899).contains(&val), "Expected ~899, got {val}");
1336 }
1337
1338 #[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 #[test]
1377 fn test_transaction_message_group_id_always_returns_transaction_id() {
1378 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 #[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 #[test]
1430 fn test_queue_url_construction_algorithm_fifo() {
1431 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 #[test]
1501 fn test_sqs_max_message_size_matches_aws_limit() {
1502 assert_eq!(SQS_MAX_MESSAGE_SIZE_BYTES, 262_144);
1504 }
1505
1506 #[test]
1509 fn test_resolve_queue_type_suffix_logic() {
1510 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 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}