1use crate::{
10 jobs::{
11 Job, NotificationSend, RelayerHealthCheck, TransactionRequest, TransactionSend,
12 TransactionStatusCheck,
13 },
14 models::RelayerError,
15 observability::request_id::get_request_id,
16 queues::{QueueBackend, QueueBackendStorage, QueueBackendType},
17};
18use async_trait::async_trait;
19use serde::Serialize;
20use std::sync::Arc;
21use thiserror::Error;
22use tracing::{debug, error};
23
24use super::{JobType, TokenSwapRequest};
25
26#[cfg(test)]
27use mockall::automock;
28
29#[derive(Debug, Error, Serialize, Clone)]
30pub enum JobProducerError {
31 #[error("Queue error: {0}")]
32 QueueError(String),
33}
34
35impl From<JobProducerError> for RelayerError {
36 fn from(err: JobProducerError) -> Self {
37 RelayerError::QueueError(err.to_string())
38 }
39}
40
41#[derive(Debug, Clone)]
43pub struct JobProducer {
44 queue_backend: Arc<QueueBackendStorage>,
45}
46
47#[async_trait]
48#[cfg_attr(test, automock)]
49pub trait JobProducerTrait: Send + Sync {
50 async fn produce_transaction_request_job(
51 &self,
52 transaction_process_job: TransactionRequest,
53 scheduled_on: Option<i64>,
54 ) -> Result<(), JobProducerError>;
55
56 async fn produce_submit_transaction_job(
57 &self,
58 transaction_submit_job: TransactionSend,
59 scheduled_on: Option<i64>,
60 ) -> Result<(), JobProducerError>;
61
62 async fn produce_check_transaction_status_job(
63 &self,
64 transaction_status_check_job: TransactionStatusCheck,
65 scheduled_on: Option<i64>,
66 ) -> Result<(), JobProducerError>;
67
68 async fn produce_send_notification_job(
69 &self,
70 notification_send_job: NotificationSend,
71 scheduled_on: Option<i64>,
72 ) -> Result<(), JobProducerError>;
73
74 async fn produce_token_swap_request_job(
75 &self,
76 swap_request_job: TokenSwapRequest,
77 scheduled_on: Option<i64>,
78 ) -> Result<(), JobProducerError>;
79
80 async fn produce_relayer_health_check_job(
81 &self,
82 relayer_health_check_job: RelayerHealthCheck,
83 scheduled_on: Option<i64>,
84 ) -> Result<(), JobProducerError>;
85
86 fn get_queue_backend(&self) -> Option<Arc<QueueBackendStorage>> {
88 None
89 }
90
91 fn backend_type(&self) -> QueueBackendType {
93 QueueBackendType::Redis
94 }
95}
96
97impl JobProducer {
98 pub fn new(queue_backend: Arc<QueueBackendStorage>) -> Self {
99 Self { queue_backend }
100 }
101
102 pub fn queue_backend(&self) -> Arc<QueueBackendStorage> {
103 self.queue_backend.clone()
104 }
105}
106
107#[async_trait]
108impl JobProducerTrait for JobProducer {
109 fn get_queue_backend(&self) -> Option<Arc<QueueBackendStorage>> {
110 Some(self.queue_backend())
111 }
112
113 fn backend_type(&self) -> QueueBackendType {
114 self.queue_backend.backend_type()
115 }
116
117 async fn produce_transaction_request_job(
118 &self,
119 transaction_process_job: TransactionRequest,
120 scheduled_on: Option<i64>,
121 ) -> Result<(), JobProducerError> {
122 debug!(
123 "Producing transaction request job: {:?}",
124 transaction_process_job
125 );
126 let job = Job::new(JobType::TransactionRequest, transaction_process_job)
127 .with_request_id(get_request_id())
128 .with_scheduled_on(scheduled_on);
129 let request_id = job.request_id.clone();
130 let tx_id = job.data.transaction_id.clone();
131 let relayer_id = job.data.relayer_id.clone();
132
133 let backend = self.queue_backend();
134 let job_id = backend
135 .produce_transaction_request(job, scheduled_on)
136 .await
137 .map_err(|e| JobProducerError::QueueError(e.to_string()))?;
138
139 debug!(
140 job_type = %JobType::TransactionRequest,
141 backend = %backend.backend_type(),
142 job_id = %job_id,
143 request_id = ?request_id,
144 tx_id = %tx_id,
145 relayer_id = %relayer_id,
146 scheduled_on = ?scheduled_on,
147 "transaction request job produced"
148 );
149
150 Ok(())
151 }
152
153 async fn produce_submit_transaction_job(
154 &self,
155 transaction_submit_job: TransactionSend,
156 scheduled_on: Option<i64>,
157 ) -> Result<(), JobProducerError> {
158 let job = Job::new(JobType::TransactionSend, transaction_submit_job)
159 .with_request_id(get_request_id())
160 .with_scheduled_on(scheduled_on);
161 let request_id = job.request_id.clone();
162 let tx_id = job.data.transaction_id.clone();
163 let relayer_id = job.data.relayer_id.clone();
164 let command = job.data.command.clone();
165
166 let backend = self.queue_backend();
167 let job_id = backend
168 .produce_transaction_submission(job, scheduled_on)
169 .await
170 .map_err(|e| JobProducerError::QueueError(e.to_string()))?;
171
172 debug!(
173 job_type = %JobType::TransactionSend,
174 backend = %backend.backend_type(),
175 job_id = %job_id,
176 request_id = ?request_id,
177 tx_id = %tx_id,
178 relayer_id = %relayer_id,
179 command = ?command,
180 scheduled_on = ?scheduled_on,
181 "transaction submission job produced"
182 );
183
184 Ok(())
185 }
186
187 async fn produce_check_transaction_status_job(
188 &self,
189 transaction_status_check_job: TransactionStatusCheck,
190 scheduled_on: Option<i64>,
191 ) -> Result<(), JobProducerError> {
192 let job = Job::new(
193 JobType::TransactionStatusCheck,
194 transaction_status_check_job.clone(),
195 )
196 .with_request_id(get_request_id())
197 .with_scheduled_on(scheduled_on);
198 let request_id = job.request_id.clone();
199 let tx_id = job.data.transaction_id.clone();
200 let relayer_id = job.data.relayer_id.clone();
201
202 let backend = self.queue_backend();
203 let job_id = backend
204 .produce_transaction_status_check(job, scheduled_on)
205 .await
206 .map_err(|e| JobProducerError::QueueError(e.to_string()))?;
207
208 debug!(
209 job_type = %JobType::TransactionStatusCheck,
210 backend = %backend.backend_type(),
211 job_id = %job_id,
212 request_id = ?request_id,
213 tx_id = %tx_id,
214 relayer_id = %relayer_id,
215 network_type = ?transaction_status_check_job.network_type,
216 scheduled_on = ?scheduled_on,
217 "Transaction Status Check job produced successfully"
218 );
219 Ok(())
220 }
221
222 async fn produce_send_notification_job(
223 &self,
224 notification_send_job: NotificationSend,
225 scheduled_on: Option<i64>,
226 ) -> Result<(), JobProducerError> {
227 let job = Job::new(JobType::NotificationSend, notification_send_job)
228 .with_request_id(get_request_id())
229 .with_scheduled_on(scheduled_on);
230 let request_id = job.request_id.clone();
231 let notification_id = job.data.notification_id.clone();
232
233 let backend = self.queue_backend();
234 let job_id = backend
235 .produce_notification(job, scheduled_on)
236 .await
237 .map_err(|e| JobProducerError::QueueError(e.to_string()))?;
238
239 debug!(
240 job_type = %JobType::NotificationSend,
241 backend = %backend.backend_type(),
242 job_id = %job_id,
243 request_id = ?request_id,
244 notification_id = %notification_id,
245 scheduled_on = ?scheduled_on,
246 "notification send job produced"
247 );
248 Ok(())
249 }
250
251 async fn produce_token_swap_request_job(
252 &self,
253 swap_request_job: TokenSwapRequest,
254 scheduled_on: Option<i64>,
255 ) -> Result<(), JobProducerError> {
256 let job = Job::new(JobType::TokenSwapRequest, swap_request_job)
257 .with_request_id(get_request_id())
258 .with_scheduled_on(scheduled_on);
259 let request_id = job.request_id.clone();
260 let relayer_id = job.data.relayer_id.clone();
261 let backend = self.queue_backend();
262 let job_id = backend
263 .produce_token_swap_request(job, scheduled_on)
264 .await
265 .map_err(|e| JobProducerError::QueueError(e.to_string()))?;
266
267 debug!(
268 job_type = %JobType::TokenSwapRequest,
269 backend = %backend.backend_type(),
270 job_id = %job_id,
271 request_id = ?request_id,
272 relayer_id = %relayer_id,
273 scheduled_on = ?scheduled_on,
274 "token swap job produced"
275 );
276 Ok(())
277 }
278
279 async fn produce_relayer_health_check_job(
280 &self,
281 relayer_health_check_job: RelayerHealthCheck,
282 scheduled_on: Option<i64>,
283 ) -> Result<(), JobProducerError> {
284 let job = Job::new(
285 JobType::RelayerHealthCheck,
286 relayer_health_check_job.clone(),
287 )
288 .with_request_id(get_request_id())
289 .with_scheduled_on(scheduled_on);
290 let request_id = job.request_id.clone();
291 let relayer_id = job.data.relayer_id.clone();
292 let backend = self.queue_backend();
293 let job_id = backend
294 .produce_relayer_health_check(job, scheduled_on)
295 .await
296 .map_err(|e| JobProducerError::QueueError(e.to_string()))?;
297
298 debug!(
299 job_type = %JobType::RelayerHealthCheck,
300 backend = %backend.backend_type(),
301 job_id = %job_id,
302 request_id = ?request_id,
303 relayer_id = %relayer_id,
304 scheduled_on = ?scheduled_on,
305 "relayer health check job produced"
306 );
307 Ok(())
308 }
309}
310
311#[cfg(test)]
312mod tests {
313 use super::*;
314 use crate::models::{
315 EvmTransactionResponse, TransactionResponse, TransactionStatus, WebhookNotification,
316 WebhookPayload, U256,
317 };
318 use crate::utils::calculate_scheduled_timestamp;
319 use tokio::sync::Mutex;
320
321 #[derive(Clone, Debug)]
322 struct TestRedisStorage<T> {
324 pub push_called: bool,
325 pub schedule_called: bool,
326 pub last_job: Option<T>,
327 pub last_scheduled_timestamp: Option<i64>,
328 _phantom: std::marker::PhantomData<T>,
329 }
330
331 impl<T> TestRedisStorage<T> {
332 fn new() -> Self {
333 Self {
334 push_called: false,
335 schedule_called: false,
336 last_job: None,
337 last_scheduled_timestamp: None,
338 _phantom: std::marker::PhantomData,
339 }
340 }
341 }
342
343 impl<T: Clone> TestRedisStorage<T> {
344 async fn push(&mut self, job: T) -> Result<(), JobProducerError> {
345 self.push_called = true;
346 self.last_job = Some(job);
347 Ok(())
348 }
349
350 async fn schedule(&mut self, job: T, timestamp: i64) -> Result<(), JobProducerError> {
351 self.schedule_called = true;
352 self.last_job = Some(job);
353 self.last_scheduled_timestamp = Some(timestamp);
354 Ok(())
355 }
356 }
357
358 #[derive(Clone, Debug)]
360 struct TestQueue {
361 pub transaction_request_queue: TestRedisStorage<Job<TransactionRequest>>,
362 pub transaction_submission_queue: TestRedisStorage<Job<TransactionSend>>,
363 pub transaction_status_queue: TestRedisStorage<Job<TransactionStatusCheck>>,
364 pub transaction_status_queue_evm: TestRedisStorage<Job<TransactionStatusCheck>>,
365 pub transaction_status_queue_stellar: TestRedisStorage<Job<TransactionStatusCheck>>,
366 pub notification_queue: TestRedisStorage<Job<NotificationSend>>,
367 pub token_swap_request_queue: TestRedisStorage<Job<TokenSwapRequest>>,
368 pub relayer_health_check_queue: TestRedisStorage<Job<RelayerHealthCheck>>,
369 }
370
371 impl TestQueue {
372 fn new() -> Self {
373 Self {
374 transaction_request_queue: TestRedisStorage::new(),
375 transaction_submission_queue: TestRedisStorage::new(),
376 transaction_status_queue: TestRedisStorage::new(),
377 transaction_status_queue_evm: TestRedisStorage::new(),
378 transaction_status_queue_stellar: TestRedisStorage::new(),
379 notification_queue: TestRedisStorage::new(),
380 token_swap_request_queue: TestRedisStorage::new(),
381 relayer_health_check_queue: TestRedisStorage::new(),
382 }
383 }
384 }
385
386 struct TestJobProducer {
388 queue: Mutex<TestQueue>,
389 }
390
391 impl Clone for TestJobProducer {
392 fn clone(&self) -> Self {
393 let queue = self
394 .queue
395 .try_lock()
396 .expect("Failed to lock queue for cloning")
397 .clone();
398 Self {
399 queue: Mutex::new(queue),
400 }
401 }
402 }
403
404 impl TestJobProducer {
405 fn new() -> Self {
406 Self {
407 queue: Mutex::new(TestQueue::new()),
408 }
409 }
410
411 async fn get_queue(&self) -> TestQueue {
412 self.queue.lock().await.clone()
413 }
414 }
415
416 #[async_trait]
417 impl JobProducerTrait for TestJobProducer {
418 fn get_queue_backend(&self) -> Option<Arc<QueueBackendStorage>> {
419 None
420 }
421
422 async fn produce_transaction_request_job(
423 &self,
424 transaction_process_job: TransactionRequest,
425 scheduled_on: Option<i64>,
426 ) -> Result<(), JobProducerError> {
427 let mut queue = self.queue.lock().await;
428 let job = Job::new(JobType::TransactionRequest, transaction_process_job)
429 .with_scheduled_on(scheduled_on);
430
431 match scheduled_on {
432 Some(scheduled_on) => {
433 queue
434 .transaction_request_queue
435 .schedule(job, scheduled_on)
436 .await?;
437 }
438 None => {
439 queue.transaction_request_queue.push(job).await?;
440 }
441 }
442
443 Ok(())
444 }
445
446 async fn produce_submit_transaction_job(
447 &self,
448 transaction_submit_job: TransactionSend,
449 scheduled_on: Option<i64>,
450 ) -> Result<(), JobProducerError> {
451 let mut queue = self.queue.lock().await;
452 let job = Job::new(JobType::TransactionSend, transaction_submit_job)
453 .with_scheduled_on(scheduled_on);
454
455 match scheduled_on {
456 Some(on) => {
457 queue.transaction_submission_queue.schedule(job, on).await?;
458 }
459 None => {
460 queue.transaction_submission_queue.push(job).await?;
461 }
462 }
463
464 Ok(())
465 }
466
467 async fn produce_check_transaction_status_job(
468 &self,
469 transaction_status_check_job: TransactionStatusCheck,
470 scheduled_on: Option<i64>,
471 ) -> Result<(), JobProducerError> {
472 let mut queue = self.queue.lock().await;
473 let job = Job::new(
474 JobType::TransactionStatusCheck,
475 transaction_status_check_job.clone(),
476 )
477 .with_scheduled_on(scheduled_on);
478
479 use crate::models::NetworkType;
481 let status_queue = match transaction_status_check_job.network_type {
482 Some(NetworkType::Evm) => &mut queue.transaction_status_queue_evm,
483 Some(NetworkType::Stellar) => &mut queue.transaction_status_queue_stellar,
484 Some(NetworkType::Solana) => &mut queue.transaction_status_queue, None => &mut queue.transaction_status_queue, };
487
488 match scheduled_on {
489 Some(on) => {
490 status_queue.schedule(job, on).await?;
491 }
492 None => {
493 status_queue.push(job).await?;
494 }
495 }
496
497 Ok(())
498 }
499
500 async fn produce_send_notification_job(
501 &self,
502 notification_send_job: NotificationSend,
503 scheduled_on: Option<i64>,
504 ) -> Result<(), JobProducerError> {
505 let mut queue = self.queue.lock().await;
506 let job = Job::new(JobType::NotificationSend, notification_send_job)
507 .with_scheduled_on(scheduled_on);
508
509 match scheduled_on {
510 Some(on) => {
511 queue.notification_queue.schedule(job, on).await?;
512 }
513 None => {
514 queue.notification_queue.push(job).await?;
515 }
516 }
517
518 Ok(())
519 }
520
521 async fn produce_token_swap_request_job(
522 &self,
523 swap_request_job: TokenSwapRequest,
524 scheduled_on: Option<i64>,
525 ) -> Result<(), JobProducerError> {
526 let mut queue = self.queue.lock().await;
527 let job = Job::new(JobType::TokenSwapRequest, swap_request_job)
528 .with_scheduled_on(scheduled_on);
529
530 match scheduled_on {
531 Some(on) => {
532 queue.token_swap_request_queue.schedule(job, on).await?;
533 }
534 None => {
535 queue.token_swap_request_queue.push(job).await?;
536 }
537 }
538
539 Ok(())
540 }
541
542 async fn produce_relayer_health_check_job(
543 &self,
544 relayer_health_check_job: RelayerHealthCheck,
545 scheduled_on: Option<i64>,
546 ) -> Result<(), JobProducerError> {
547 let mut queue = self.queue.lock().await;
548 let job = Job::new(JobType::RelayerHealthCheck, relayer_health_check_job)
549 .with_scheduled_on(scheduled_on);
550
551 match scheduled_on {
552 Some(scheduled_on) => {
553 queue
554 .relayer_health_check_queue
555 .schedule(job, scheduled_on)
556 .await?;
557 }
558 None => {
559 queue.relayer_health_check_queue.push(job).await?;
560 }
561 }
562
563 Ok(())
564 }
565 }
566
567 #[tokio::test]
568 async fn test_job_producer_operations() {
569 let producer = TestJobProducer::new();
570
571 let request = TransactionRequest::new("tx123", "relayer-1");
573 let result = producer
574 .produce_transaction_request_job(request, None)
575 .await;
576 assert!(result.is_ok());
577
578 let queue = producer.get_queue().await;
579 assert!(queue.transaction_request_queue.push_called);
580
581 let producer = TestJobProducer::new();
583 let request = TransactionRequest::new("tx123", "relayer-1");
584 let scheduled_timestamp = calculate_scheduled_timestamp(10); let result = producer
586 .produce_transaction_request_job(request, Some(scheduled_timestamp))
587 .await;
588 assert!(result.is_ok());
589
590 let queue = producer.get_queue().await;
591 assert!(queue.transaction_request_queue.schedule_called);
592 }
593
594 #[tokio::test]
595 async fn test_submit_transaction_job() {
596 let producer = TestJobProducer::new();
597
598 let submit_job = TransactionSend::submit("tx123", "relayer-1");
600 let result = producer
601 .produce_submit_transaction_job(submit_job, None)
602 .await;
603 assert!(result.is_ok());
604
605 let queue = producer.get_queue().await;
606 assert!(queue.transaction_submission_queue.push_called);
607 }
608
609 #[tokio::test]
610 async fn test_check_status_job() {
611 use crate::models::NetworkType;
612 let producer = TestJobProducer::new();
613
614 let status_job = TransactionStatusCheck::new("tx123", "relayer-1", NetworkType::Evm);
616 let result = producer
617 .produce_check_transaction_status_job(status_job, None)
618 .await;
619 assert!(result.is_ok());
620
621 let queue = producer.get_queue().await;
622 assert!(queue.transaction_status_queue_evm.push_called);
623 }
624
625 #[tokio::test]
626 async fn test_notification_job() {
627 let producer = TestJobProducer::new();
628
629 let notification = WebhookNotification::new(
631 "test_event".to_string(),
632 WebhookPayload::Transaction(TransactionResponse::Evm(Box::new(
633 EvmTransactionResponse {
634 id: "tx123".to_string(),
635 hash: Some("0x123".to_string()),
636 status: TransactionStatus::Confirmed,
637 status_reason: None,
638 created_at: "2025-01-27T15:31:10.777083+00:00".to_string(),
639 sent_at: Some("2025-01-27T15:31:10.777083+00:00".to_string()),
640 confirmed_at: Some("2025-01-27T15:31:10.777083+00:00".to_string()),
641 gas_price: Some(1000000000),
642 gas_limit: Some(21000),
643 nonce: Some(1),
644 value: U256::from(1000000000000000000_u64),
645 from: "0xabc".to_string(),
646 to: Some("0xdef".to_string()),
647 relayer_id: "relayer-1".to_string(),
648 data: None,
649 max_fee_per_gas: None,
650 max_priority_fee_per_gas: None,
651 signature: None,
652 speed: None,
653 is_canceled: None,
654 },
655 ))),
656 );
657 let job = NotificationSend::new("notification-1".to_string(), notification);
658
659 let result = producer.produce_send_notification_job(job, None).await;
660 assert!(result.is_ok());
661
662 let queue = producer.get_queue().await;
663 assert!(queue.notification_queue.push_called);
664 }
665
666 #[tokio::test]
667 async fn test_relayer_health_check_job() {
668 let producer = TestJobProducer::new();
669
670 let health_check = RelayerHealthCheck::new("relayer-1".to_string());
672 let result = producer
673 .produce_relayer_health_check_job(health_check, None)
674 .await;
675 assert!(result.is_ok());
676
677 let queue = producer.get_queue().await;
678 assert!(queue.relayer_health_check_queue.push_called);
679
680 let producer = TestJobProducer::new();
682 let health_check = RelayerHealthCheck::new("relayer-1".to_string());
683 let scheduled_timestamp = calculate_scheduled_timestamp(60);
684 let result = producer
685 .produce_relayer_health_check_job(health_check, Some(scheduled_timestamp))
686 .await;
687 assert!(result.is_ok());
688
689 let queue = producer.get_queue().await;
690 assert!(queue.relayer_health_check_queue.schedule_called);
691 }
692
693 #[test]
694 fn test_job_producer_error_conversion() {
695 let job_error = JobProducerError::QueueError("Test error".to_string());
697 let relayer_error: RelayerError = job_error.into();
698
699 match relayer_error {
700 RelayerError::QueueError(msg) => {
701 assert_eq!(msg, "Queue error: Test error");
702 }
703 _ => panic!("Unexpected error type"),
704 }
705 }
706
707 #[tokio::test]
708 async fn test_get_queue() {
709 let producer = TestJobProducer::new();
710
711 let queue = producer.get_queue().await;
713
714 assert!(!queue.transaction_request_queue.push_called);
716 assert!(!queue.transaction_request_queue.schedule_called);
717 assert!(!queue.transaction_submission_queue.push_called);
718 assert!(!queue.notification_queue.push_called);
719 assert!(!queue.token_swap_request_queue.push_called);
720 assert!(!queue.relayer_health_check_queue.push_called);
721 }
722
723 #[tokio::test]
724 async fn test_produce_relayer_health_check_job_immediate() {
725 let producer = TestJobProducer::new();
726
727 let health_check = RelayerHealthCheck::new("relayer-1".to_string());
729 let result = producer
730 .produce_relayer_health_check_job(health_check, None)
731 .await;
732
733 assert!(result.is_ok());
735
736 let queue = producer.get_queue().await;
738 assert!(queue.relayer_health_check_queue.push_called);
739 assert!(!queue.relayer_health_check_queue.schedule_called);
740
741 assert!(!queue.transaction_request_queue.push_called);
743 assert!(!queue.transaction_submission_queue.push_called);
744 assert!(!queue.transaction_status_queue.push_called);
745 assert!(!queue.notification_queue.push_called);
746 assert!(!queue.token_swap_request_queue.push_called);
747 }
748
749 #[tokio::test]
750 async fn test_produce_relayer_health_check_job_scheduled() {
751 let producer = TestJobProducer::new();
752
753 let health_check = RelayerHealthCheck::new("relayer-2".to_string());
755 let scheduled_timestamp = calculate_scheduled_timestamp(300); let result = producer
757 .produce_relayer_health_check_job(health_check, Some(scheduled_timestamp))
758 .await;
759
760 assert!(result.is_ok());
762
763 let queue = producer.get_queue().await;
765 assert!(queue.relayer_health_check_queue.schedule_called);
766 assert!(!queue.relayer_health_check_queue.push_called);
767
768 assert!(!queue.transaction_request_queue.push_called);
770 assert!(!queue.transaction_submission_queue.push_called);
771 assert!(!queue.transaction_status_queue.push_called);
772 assert!(!queue.notification_queue.push_called);
773 assert!(!queue.token_swap_request_queue.push_called);
774 }
775
776 #[tokio::test]
777 async fn test_produce_relayer_health_check_job_multiple_relayers() {
778 let producer = TestJobProducer::new();
779
780 let relayer_ids = vec!["relayer-1", "relayer-2", "relayer-3"];
782
783 for relayer_id in &relayer_ids {
784 let health_check = RelayerHealthCheck::new(relayer_id.to_string());
785 let result = producer
786 .produce_relayer_health_check_job(health_check, None)
787 .await;
788 assert!(result.is_ok());
789 }
790
791 let queue = producer.get_queue().await;
793 assert!(queue.relayer_health_check_queue.push_called);
794 }
795
796 #[tokio::test]
797 async fn test_status_check_routes_to_evm_queue() {
798 use crate::models::NetworkType;
799 let producer = TestJobProducer::new();
800
801 let status_job = TransactionStatusCheck::new("tx-evm", "relayer-1", NetworkType::Evm);
802 let result = producer
803 .produce_check_transaction_status_job(status_job, None)
804 .await;
805
806 assert!(result.is_ok());
807 let queue = producer.get_queue().await;
808 assert!(queue.transaction_status_queue_evm.push_called);
809 assert!(!queue.transaction_status_queue_stellar.push_called);
810 assert!(!queue.transaction_status_queue.push_called);
811 }
812
813 #[tokio::test]
814 async fn test_status_check_routes_to_stellar_queue() {
815 use crate::models::NetworkType;
816 let producer = TestJobProducer::new();
817
818 let status_job =
819 TransactionStatusCheck::new("tx-stellar", "relayer-2", NetworkType::Stellar);
820 let result = producer
821 .produce_check_transaction_status_job(status_job, None)
822 .await;
823
824 assert!(result.is_ok());
825 let queue = producer.get_queue().await;
826 assert!(queue.transaction_status_queue_stellar.push_called);
827 assert!(!queue.transaction_status_queue_evm.push_called);
828 assert!(!queue.transaction_status_queue.push_called);
829 }
830
831 #[tokio::test]
832 async fn test_status_check_routes_to_default_queue_for_solana() {
833 use crate::models::NetworkType;
834 let producer = TestJobProducer::new();
835
836 let status_job = TransactionStatusCheck::new("tx-solana", "relayer-3", NetworkType::Solana);
837 let result = producer
838 .produce_check_transaction_status_job(status_job, None)
839 .await;
840
841 assert!(result.is_ok());
842 let queue = producer.get_queue().await;
843 assert!(queue.transaction_status_queue.push_called);
844 assert!(!queue.transaction_status_queue_evm.push_called);
845 assert!(!queue.transaction_status_queue_stellar.push_called);
846 }
847
848 #[tokio::test]
849 async fn test_status_check_scheduled_evm() {
850 use crate::models::NetworkType;
851 let producer = TestJobProducer::new();
852
853 let status_job =
854 TransactionStatusCheck::new("tx-evm-scheduled", "relayer-1", NetworkType::Evm);
855 let scheduled_timestamp = calculate_scheduled_timestamp(30);
856 let result = producer
857 .produce_check_transaction_status_job(status_job, Some(scheduled_timestamp))
858 .await;
859
860 assert!(result.is_ok());
861 let queue = producer.get_queue().await;
862 assert!(queue.transaction_status_queue_evm.schedule_called);
863 assert!(!queue.transaction_status_queue_evm.push_called);
864 }
865
866 #[tokio::test]
867 async fn test_scheduled_status_check_persists_available_at() {
868 use crate::models::NetworkType;
869 let producer = TestJobProducer::new();
870
871 let status_job =
872 TransactionStatusCheck::new("tx-evm-scheduled", "relayer-1", NetworkType::Evm);
873 let scheduled_timestamp = calculate_scheduled_timestamp(30);
874 producer
875 .produce_check_transaction_status_job(status_job, Some(scheduled_timestamp))
876 .await
877 .unwrap();
878
879 let queue = producer.get_queue().await;
880 let stored_job = queue
881 .transaction_status_queue_evm
882 .last_job
883 .expect("scheduled status check should be stored");
884 let expected_available_at = scheduled_timestamp.to_string();
885
886 assert_eq!(
887 stored_job.available_at.as_deref(),
888 Some(expected_available_at.as_str())
889 );
890 assert_eq!(
891 queue.transaction_status_queue_evm.last_scheduled_timestamp,
892 Some(scheduled_timestamp)
893 );
894 }
895
896 #[tokio::test]
897 async fn test_submit_transaction_scheduled() {
898 let producer = TestJobProducer::new();
899
900 let submit_job = TransactionSend::submit("tx-scheduled", "relayer-1");
901 let scheduled_timestamp = calculate_scheduled_timestamp(15);
902 let result = producer
903 .produce_submit_transaction_job(submit_job, Some(scheduled_timestamp))
904 .await;
905
906 assert!(result.is_ok());
907 let queue = producer.get_queue().await;
908 assert!(queue.transaction_submission_queue.schedule_called);
909 assert!(!queue.transaction_submission_queue.push_called);
910 }
911
912 #[tokio::test]
913 async fn test_scheduled_submit_job_persists_available_at() {
914 let producer = TestJobProducer::new();
915
916 let submit_job = TransactionSend::submit("tx-scheduled", "relayer-1");
917 let scheduled_timestamp = calculate_scheduled_timestamp(15);
918 producer
919 .produce_submit_transaction_job(submit_job, Some(scheduled_timestamp))
920 .await
921 .unwrap();
922
923 let queue = producer.get_queue().await;
924 let stored_job = queue
925 .transaction_submission_queue
926 .last_job
927 .expect("scheduled submission should be stored");
928 let expected_available_at = scheduled_timestamp.to_string();
929
930 assert_eq!(
931 stored_job.available_at.as_deref(),
932 Some(expected_available_at.as_str())
933 );
934 }
935
936 #[tokio::test]
937 async fn test_notification_job_scheduled() {
938 let producer = TestJobProducer::new();
939
940 let notification = WebhookNotification::new(
941 "test_scheduled_event".to_string(),
942 WebhookPayload::Transaction(TransactionResponse::Evm(Box::new(
943 EvmTransactionResponse {
944 id: "tx-notify-scheduled".to_string(),
945 hash: Some("0xabc123".to_string()),
946 status: TransactionStatus::Confirmed,
947 status_reason: None,
948 created_at: "2025-01-27T15:31:10.777083+00:00".to_string(),
949 sent_at: Some("2025-01-27T15:31:10.777083+00:00".to_string()),
950 confirmed_at: Some("2025-01-27T15:31:10.777083+00:00".to_string()),
951 gas_price: Some(1000000000),
952 gas_limit: Some(21000),
953 nonce: Some(1),
954 value: U256::from(1000000000000000000_u64),
955 from: "0xabc".to_string(),
956 to: Some("0xdef".to_string()),
957 relayer_id: "relayer-1".to_string(),
958 data: None,
959 max_fee_per_gas: None,
960 max_priority_fee_per_gas: None,
961 signature: None,
962 speed: None,
963 is_canceled: None,
964 },
965 ))),
966 );
967 let job = NotificationSend::new("notification-scheduled".to_string(), notification);
968
969 let scheduled_timestamp = calculate_scheduled_timestamp(5);
970 let result = producer
971 .produce_send_notification_job(job, Some(scheduled_timestamp))
972 .await;
973
974 assert!(result.is_ok());
975 let queue = producer.get_queue().await;
976 assert!(queue.notification_queue.schedule_called);
977 assert!(!queue.notification_queue.push_called);
978 }
979
980 #[tokio::test]
981 async fn test_solana_swap_job_immediate() {
982 let producer = TestJobProducer::new();
983
984 let swap_job = TokenSwapRequest::new("relayer-solana".to_string());
985 let result = producer
986 .produce_token_swap_request_job(swap_job, None)
987 .await;
988
989 assert!(result.is_ok());
990 let queue = producer.get_queue().await;
991 assert!(queue.token_swap_request_queue.push_called);
992 assert!(!queue.token_swap_request_queue.schedule_called);
993 }
994
995 #[tokio::test]
996 async fn test_token_swap_job_scheduled() {
997 let producer = TestJobProducer::new();
998
999 let swap_job = TokenSwapRequest::new("relayer-solana".to_string());
1000 let scheduled_timestamp = calculate_scheduled_timestamp(20);
1001 let result = producer
1002 .produce_token_swap_request_job(swap_job, Some(scheduled_timestamp))
1003 .await;
1004
1005 assert!(result.is_ok());
1006 let queue = producer.get_queue().await;
1007 assert!(queue.token_swap_request_queue.schedule_called);
1008 assert!(!queue.token_swap_request_queue.push_called);
1009 }
1010
1011 #[tokio::test]
1012 async fn test_transaction_send_cancel_job() {
1013 let producer = TestJobProducer::new();
1014
1015 let cancel_job = TransactionSend::cancel("tx-cancel", "relayer-1", "user requested");
1016 let result = producer
1017 .produce_submit_transaction_job(cancel_job, None)
1018 .await;
1019
1020 assert!(result.is_ok());
1021 let queue = producer.get_queue().await;
1022 assert!(queue.transaction_submission_queue.push_called);
1023 }
1024
1025 #[tokio::test]
1026 async fn test_transaction_send_resubmit_job() {
1027 let producer = TestJobProducer::new();
1028
1029 let resubmit_job = TransactionSend::resubmit("tx-resubmit", "relayer-1");
1030 let result = producer
1031 .produce_submit_transaction_job(resubmit_job, None)
1032 .await;
1033
1034 assert!(result.is_ok());
1035 let queue = producer.get_queue().await;
1036 assert!(queue.transaction_submission_queue.push_called);
1037 }
1038
1039 #[tokio::test]
1040 async fn test_transaction_send_resend_job() {
1041 let producer = TestJobProducer::new();
1042
1043 let resend_job = TransactionSend::resend("tx-resend", "relayer-1");
1044 let result = producer
1045 .produce_submit_transaction_job(resend_job, None)
1046 .await;
1047
1048 assert!(result.is_ok());
1049 let queue = producer.get_queue().await;
1050 assert!(queue.transaction_submission_queue.push_called);
1051 }
1052
1053 #[tokio::test]
1054 async fn test_multiple_jobs_different_queues() {
1055 let producer = TestJobProducer::new();
1056
1057 let request = TransactionRequest::new("tx1", "relayer-1");
1059 producer
1060 .produce_transaction_request_job(request, None)
1061 .await
1062 .unwrap();
1063
1064 let submit = TransactionSend::submit("tx2", "relayer-1");
1065 producer
1066 .produce_submit_transaction_job(submit, None)
1067 .await
1068 .unwrap();
1069
1070 use crate::models::NetworkType;
1071 let status = TransactionStatusCheck::new("tx3", "relayer-1", NetworkType::Evm);
1072 producer
1073 .produce_check_transaction_status_job(status, None)
1074 .await
1075 .unwrap();
1076
1077 let queue = producer.get_queue().await;
1079 assert!(queue.transaction_request_queue.push_called);
1080 assert!(queue.transaction_submission_queue.push_called);
1081 assert!(queue.transaction_status_queue_evm.push_called);
1082 }
1083
1084 #[test]
1085 fn test_job_producer_clone() {
1086 let producer = TestJobProducer::new();
1087 let cloned_producer = producer.clone();
1088
1089 assert!(std::ptr::addr_of!(producer) != std::ptr::addr_of!(cloned_producer));
1092 }
1093
1094 #[tokio::test]
1095 async fn test_transaction_request_with_metadata() {
1096 let producer = TestJobProducer::new();
1097
1098 let mut metadata = std::collections::HashMap::new();
1099 metadata.insert("retry_count".to_string(), "3".to_string());
1100
1101 let request = TransactionRequest::new("tx-meta", "relayer-1").with_metadata(metadata);
1102
1103 let result = producer
1104 .produce_transaction_request_job(request, None)
1105 .await;
1106
1107 assert!(result.is_ok());
1108 let queue = producer.get_queue().await;
1109 assert!(queue.transaction_request_queue.push_called);
1110 }
1111
1112 #[tokio::test]
1113 async fn test_status_check_with_metadata() {
1114 use crate::models::NetworkType;
1115 let producer = TestJobProducer::new();
1116
1117 let mut metadata = std::collections::HashMap::new();
1118 metadata.insert("attempt".to_string(), "2".to_string());
1119
1120 let status =
1121 TransactionStatusCheck::new("tx-status-meta", "relayer-1", NetworkType::Stellar)
1122 .with_metadata(metadata);
1123
1124 let result = producer
1125 .produce_check_transaction_status_job(status, None)
1126 .await;
1127
1128 assert!(result.is_ok());
1129 let queue = producer.get_queue().await;
1130 assert!(queue.transaction_status_queue_stellar.push_called);
1131 }
1132
1133 #[tokio::test]
1134 async fn test_scheduled_jobs_with_different_delays() {
1135 let producer = TestJobProducer::new();
1136
1137 let delays = [1, 10, 60, 300, 3600]; for (idx, delay) in delays.iter().enumerate() {
1141 let request = TransactionRequest::new(format!("tx-delay-{idx}"), "relayer-1");
1142 let timestamp = calculate_scheduled_timestamp(*delay);
1143
1144 let result = producer
1145 .produce_transaction_request_job(request, Some(timestamp))
1146 .await;
1147
1148 assert!(result.is_ok(), "Failed to schedule job with delay {delay}");
1149 }
1150 }
1151
1152 #[test]
1153 fn test_job_producer_error_display() {
1154 let error = JobProducerError::QueueError("Test queue error".to_string());
1155 let error_string = error.to_string();
1156
1157 assert!(error_string.contains("Queue error"));
1158 assert!(error_string.contains("Test queue error"));
1159 }
1160
1161 #[test]
1162 fn test_job_producer_error_to_relayer_error() {
1163 let job_error = JobProducerError::QueueError("Connection failed".to_string());
1165 let relayer_error: RelayerError = job_error.into();
1166
1167 match relayer_error {
1168 RelayerError::QueueError(msg) => {
1169 assert_eq!(msg, "Queue error: Connection failed");
1170 }
1171 _ => panic!("Expected QueueError variant"),
1172 }
1173 }
1174}