1use crate::{
7 constants::DEFAULT_STELLAR_CONCURRENT_TRANSACTIONS,
8 domain::transaction::{stellar::fetch_next_sequence_from_chain, Transaction},
9 jobs::{JobProducer, JobProducerTrait, StatusCheckContext, TransactionRequest},
10 models::{
11 produce_transaction_update_notification_payload, NetworkTransactionRequest,
12 PaginationQuery, RelayerNetworkPolicy, RelayerRepoModel, TransactionError,
13 TransactionRepoModel, TransactionStatus, TransactionUpdateRequest,
14 },
15 repositories::{
16 RelayerRepositoryStorage, Repository, TransactionCounterRepositoryStorage,
17 TransactionCounterTrait, TransactionRepository, TransactionRepositoryStorage,
18 },
19 services::{
20 provider::{StellarProvider, StellarProviderTrait},
21 signer::{Signer, StellarSignTrait, StellarSigner},
22 stellar_dex::{StellarDexService, StellarDexServiceTrait},
23 },
24 utils::calculate_scheduled_timestamp,
25};
26use async_trait::async_trait;
27use std::sync::Arc;
28use tracing::{error, info, warn};
29
30use super::lane_gate;
31
32#[allow(dead_code)]
33pub struct StellarRelayerTransaction<R, T, J, S, P, C, D>
34where
35 R: Repository<RelayerRepoModel, String>,
36 T: TransactionRepository,
37 J: JobProducerTrait,
38 S: Signer + StellarSignTrait,
39 P: StellarProviderTrait,
40 C: TransactionCounterTrait,
41 D: StellarDexServiceTrait + Send + Sync + 'static,
42{
43 relayer: RelayerRepoModel,
44 relayer_repository: Arc<R>,
45 transaction_repository: Arc<T>,
46 job_producer: Arc<J>,
47 signer: Arc<S>,
48 provider: P,
49 transaction_counter_service: Arc<C>,
50 dex_service: Arc<D>,
51}
52
53#[allow(dead_code)]
54impl<R, T, J, S, P, C, D> StellarRelayerTransaction<R, T, J, S, P, C, D>
55where
56 R: Repository<RelayerRepoModel, String>,
57 T: TransactionRepository,
58 J: JobProducerTrait,
59 S: Signer + StellarSignTrait,
60 P: StellarProviderTrait,
61 C: TransactionCounterTrait,
62 D: StellarDexServiceTrait + Send + Sync + 'static,
63{
64 #[allow(clippy::too_many_arguments)]
81 pub fn new(
82 relayer: RelayerRepoModel,
83 relayer_repository: Arc<R>,
84 transaction_repository: Arc<T>,
85 job_producer: Arc<J>,
86 signer: Arc<S>,
87 provider: P,
88 transaction_counter_service: Arc<C>,
89 dex_service: Arc<D>,
90 ) -> Result<Self, TransactionError> {
91 Ok(Self {
92 relayer,
93 relayer_repository,
94 transaction_repository,
95 job_producer,
96 signer,
97 provider,
98 transaction_counter_service,
99 dex_service,
100 })
101 }
102
103 pub fn provider(&self) -> &P {
104 &self.provider
105 }
106
107 pub fn relayer(&self) -> &RelayerRepoModel {
108 &self.relayer
109 }
110
111 pub fn job_producer(&self) -> &J {
112 &self.job_producer
113 }
114
115 pub fn transaction_repository(&self) -> &T {
116 &self.transaction_repository
117 }
118
119 pub fn signer(&self) -> &S {
120 &self.signer
121 }
122
123 pub fn transaction_counter_service(&self) -> &C {
124 &self.transaction_counter_service
125 }
126
127 pub fn dex_service(&self) -> &D {
128 &self.dex_service
129 }
130
131 pub fn concurrent_transactions_enabled(&self) -> bool {
132 if let RelayerNetworkPolicy::Stellar(policy) = &self.relayer().policies {
133 policy
134 .concurrent_transactions
135 .unwrap_or(DEFAULT_STELLAR_CONCURRENT_TRANSACTIONS)
136 } else {
137 DEFAULT_STELLAR_CONCURRENT_TRANSACTIONS
138 }
139 }
140
141 pub async fn send_transaction_request_job(
143 &self,
144 tx: &TransactionRepoModel,
145 delay_seconds: Option<i64>,
146 ) -> Result<(), TransactionError> {
147 let job = TransactionRequest::new(tx.id.clone(), tx.relayer_id.clone());
148 let scheduled_on = delay_seconds.map(calculate_scheduled_timestamp);
149 self.job_producer()
150 .produce_transaction_request_job(job, scheduled_on)
151 .await?;
152 Ok(())
153 }
154
155 pub(super) async fn send_transaction_update_notification(&self, tx: &TransactionRepoModel) {
160 if let Some(notification_id) = &self.relayer().notification_id {
161 if let Err(e) = self
162 .job_producer()
163 .produce_send_notification_job(
164 produce_transaction_update_notification_payload(notification_id, tx),
165 None,
166 )
167 .await
168 {
169 error!(error = %e, "failed to produce notification job");
170 }
171 }
172 }
173
174 pub async fn finalize_transaction_state(
176 &self,
177 tx_id: String,
178 update_req: TransactionUpdateRequest,
179 ) -> Result<TransactionRepoModel, TransactionError> {
180 let updated_tx = self
181 .transaction_repository()
182 .partial_update(tx_id, update_req)
183 .await?;
184
185 self.send_transaction_update_notification(&updated_tx).await;
186 Ok(updated_tx)
187 }
188
189 pub async fn enqueue_next_pending_transaction(
190 &self,
191 finished_tx_id: &str,
192 ) -> Result<(), TransactionError> {
193 if !self.concurrent_transactions_enabled() {
194 if let Some(next) = self
195 .find_oldest_pending_for_relayer(&self.relayer().id)
196 .await?
197 {
198 info!(to_tx_id = %next.id, finished_tx_id = %finished_tx_id, "handing over lane");
200 lane_gate::pass_to(&self.relayer().id, finished_tx_id, &next.id);
201 self.send_transaction_request_job(&next, None).await?;
202 } else {
203 info!(finished_tx_id = %finished_tx_id, "releasing relayer lane");
204 lane_gate::free(&self.relayer().id, finished_tx_id);
205 }
206 }
207 Ok(())
208 }
209
210 async fn find_oldest_pending_for_relayer(
215 &self,
216 relayer_id: &str,
217 ) -> Result<Option<TransactionRepoModel>, TransactionError> {
218 let result = self
219 .transaction_repository()
220 .find_by_status_paginated(
221 relayer_id,
222 &[TransactionStatus::Pending],
223 PaginationQuery {
224 page: 1,
225 per_page: 1,
226 },
227 true, )
229 .await
230 .map_err(TransactionError::from)?;
231
232 Ok(result.items.into_iter().next())
234 }
235
236 pub async fn sync_sequence_from_chain(
239 &self,
240 relayer_address: &str,
241 ) -> Result<(), TransactionError> {
242 info!(address = %relayer_address, "syncing sequence number from chain");
243
244 let next_usable_seq = fetch_next_sequence_from_chain(self.provider(), relayer_address)
246 .await
247 .map_err(|e| {
248 warn!(
249 address = %relayer_address,
250 error = %e,
251 "failed to fetch sequence from chain in sync_sequence_from_chain"
252 );
253 TransactionError::UnexpectedError(format!(
254 "Failed to sync sequence from chain: {e}"
255 ))
256 })?;
257
258 let effective_seq = self
264 .transaction_counter_service()
265 .sync_floor(&self.relayer().id, relayer_address, next_usable_seq)
266 .await
267 .map_err(|e| {
268 TransactionError::UnexpectedError(format!("Failed to update sequence counter: {e}"))
269 })?;
270
271 info!(
272 chain_sequence = %next_usable_seq,
273 effective_sequence = %effective_seq,
274 "synced local sequence counter to chain floor"
275 );
276 Ok(())
277 }
278
279 pub async fn reset_transaction_for_retry(
282 &self,
283 tx: TransactionRepoModel,
284 ) -> Result<TransactionRepoModel, TransactionError> {
285 info!("resetting transaction for retry through pipeline");
286
287 let update_req = tx.create_reset_update_request()?;
289
290 let reset_tx = self
292 .transaction_repository()
293 .partial_update(tx.id.clone(), update_req)
294 .await?;
295
296 info!("transaction reset successfully to pre-prepare state");
297 Ok(reset_tx)
298 }
299}
300
301#[async_trait]
302impl<R, T, J, S, P, C, D> Transaction for StellarRelayerTransaction<R, T, J, S, P, C, D>
303where
304 R: Repository<RelayerRepoModel, String> + Send + Sync,
305 T: TransactionRepository + Send + Sync,
306 J: JobProducerTrait + Send + Sync,
307 S: Signer + StellarSignTrait + Send + Sync,
308 P: StellarProviderTrait + Send + Sync,
309 C: TransactionCounterTrait + Send + Sync,
310 D: StellarDexServiceTrait + Send + Sync + 'static,
311{
312 async fn prepare_transaction(
313 &self,
314 tx: TransactionRepoModel,
315 ) -> Result<TransactionRepoModel, TransactionError> {
316 self.prepare_transaction_impl(tx).await
317 }
318
319 async fn submit_transaction(
320 &self,
321 tx: TransactionRepoModel,
322 ) -> Result<TransactionRepoModel, TransactionError> {
323 self.submit_transaction_impl(tx).await
324 }
325
326 async fn resubmit_transaction(
327 &self,
328 tx: TransactionRepoModel,
329 ) -> Result<TransactionRepoModel, TransactionError> {
330 Ok(tx)
331 }
332
333 async fn handle_transaction_status(
334 &self,
335 tx: TransactionRepoModel,
336 context: Option<StatusCheckContext>,
337 ) -> Result<TransactionRepoModel, TransactionError> {
338 self.handle_transaction_status_impl(tx, context).await
339 }
340
341 async fn cancel_transaction(
342 &self,
343 tx: TransactionRepoModel,
344 ) -> Result<TransactionRepoModel, TransactionError> {
345 Ok(tx)
346 }
347
348 async fn replace_transaction(
349 &self,
350 _old_tx: TransactionRepoModel,
351 _new_tx_request: NetworkTransactionRequest,
352 ) -> Result<TransactionRepoModel, TransactionError> {
353 Ok(_old_tx)
354 }
355
356 async fn sign_transaction(
357 &self,
358 tx: TransactionRepoModel,
359 ) -> Result<TransactionRepoModel, TransactionError> {
360 Ok(tx)
361 }
362
363 async fn validate_transaction(
364 &self,
365 _tx: TransactionRepoModel,
366 ) -> Result<bool, TransactionError> {
367 Ok(true)
368 }
369}
370
371pub type DefaultStellarTransaction = StellarRelayerTransaction<
372 RelayerRepositoryStorage,
373 TransactionRepositoryStorage,
374 JobProducer,
375 StellarSigner,
376 StellarProvider,
377 TransactionCounterRepositoryStorage,
378 StellarDexService<StellarProvider, StellarSigner>,
379>;
380
381#[cfg(test)]
382mod tests {
383 use super::*;
384 use crate::repositories::transaction::RedisTransactionRepository;
385 use crate::utils::RedisConnections;
386 use crate::{
387 models::{NetworkTransactionData, RepositoryError},
388 services::provider::ProviderError,
389 };
390 use deadpool_redis::{Config, Runtime};
391 use std::sync::Arc;
392 use uuid::Uuid;
393
394 use crate::domain::transaction::stellar::test_helpers::*;
395
396 #[test]
397 fn new_returns_ok() {
398 let relayer = create_test_relayer();
399 let mocks = default_test_mocks();
400 let result = StellarRelayerTransaction::new(
401 relayer,
402 Arc::new(mocks.relayer_repo),
403 Arc::new(mocks.tx_repo),
404 Arc::new(mocks.job_producer),
405 Arc::new(mocks.signer),
406 mocks.provider,
407 Arc::new(mocks.counter),
408 Arc::new(mocks.dex_service),
409 );
410 assert!(result.is_ok());
411 }
412
413 #[test]
414 fn accessor_methods_return_correct_references() {
415 let relayer = create_test_relayer();
416 let mocks = default_test_mocks();
417 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
418
419 assert_eq!(handler.relayer().id, "relayer-1");
421 assert_eq!(handler.relayer().address, TEST_PK);
422
423 let _ = handler.provider();
425 let _ = handler.job_producer();
426 let _ = handler.transaction_repository();
427 let _ = handler.signer();
428 let _ = handler.transaction_counter_service();
429 }
430
431 #[tokio::test]
432 async fn send_transaction_request_job_success() {
433 let relayer = create_test_relayer();
434 let mut mocks = default_test_mocks();
435
436 mocks
437 .job_producer
438 .expect_produce_transaction_request_job()
439 .withf(|job, delay| {
440 job.transaction_id == "tx-1" && job.relayer_id == "relayer-1" && delay.is_none()
441 })
442 .times(1)
443 .returning(|_, _| Box::pin(async { Ok(()) }));
444
445 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
446 let tx = create_test_transaction(&relayer.id);
447
448 let result = handler.send_transaction_request_job(&tx, None).await;
449 assert!(result.is_ok());
450 }
451
452 #[tokio::test]
453 async fn send_transaction_request_job_with_delay() {
454 let relayer = create_test_relayer();
455 let mut mocks = default_test_mocks();
456
457 mocks
458 .job_producer
459 .expect_produce_transaction_request_job()
460 .withf(|job, delay| {
461 job.transaction_id == "tx-1"
462 && job.relayer_id == "relayer-1"
463 && delay.is_some()
464 && delay.unwrap() > chrono::Utc::now().timestamp()
465 })
466 .times(1)
467 .returning(|_, _| Box::pin(async { Ok(()) }));
468
469 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
470 let tx = create_test_transaction(&relayer.id);
471
472 let result = handler.send_transaction_request_job(&tx, Some(60)).await;
473 assert!(result.is_ok());
474 }
475
476 #[tokio::test]
477 async fn finalize_transaction_state_success() {
478 let relayer = create_test_relayer();
479 let mut mocks = default_test_mocks();
480
481 mocks
483 .tx_repo
484 .expect_partial_update()
485 .withf(|tx_id, update| {
486 tx_id == "tx-1"
487 && update.status == Some(TransactionStatus::Confirmed)
488 && update.status_reason == Some("Transaction confirmed".to_string())
489 })
490 .times(1)
491 .returning(|tx_id, update| {
492 let mut tx = create_test_transaction("relayer-1");
493 tx.id = tx_id;
494 tx.status = update.status.unwrap();
495 tx.status_reason = update.status_reason;
496 tx.confirmed_at = update.confirmed_at;
497 Ok::<_, RepositoryError>(tx)
498 });
499
500 mocks
502 .job_producer
503 .expect_produce_send_notification_job()
504 .times(1)
505 .returning(|_, _| Box::pin(async { Ok(()) }));
506
507 let handler = make_stellar_tx_handler(relayer, mocks);
508
509 let update_request = TransactionUpdateRequest {
510 status: Some(TransactionStatus::Confirmed),
511 status_reason: Some("Transaction confirmed".to_string()),
512 confirmed_at: Some("2023-01-01T00:00:00Z".to_string()),
513 ..Default::default()
514 };
515
516 let result = handler
517 .finalize_transaction_state("tx-1".to_string(), update_request)
518 .await;
519
520 assert!(result.is_ok());
521 let updated_tx = result.unwrap();
522 assert_eq!(updated_tx.status, TransactionStatus::Confirmed);
523 assert_eq!(
524 updated_tx.status_reason,
525 Some("Transaction confirmed".to_string())
526 );
527 }
528
529 #[tokio::test]
530 async fn enqueue_next_pending_transaction_with_pending_tx() {
531 let relayer = create_test_relayer();
532 let mut mocks = default_test_mocks();
533
534 let mut pending_tx = create_test_transaction(&relayer.id);
536 pending_tx.id = "pending-tx-1".to_string();
537
538 mocks
539 .tx_repo
540 .expect_find_by_status_paginated()
541 .withf(|_relayer_id, statuses, query, oldest_first| {
542 statuses == [TransactionStatus::Pending]
543 && query.page == 1
544 && query.per_page == 1
545 && *oldest_first
546 })
547 .times(1)
548 .returning(move |_, _, _, _| {
549 let mut tx = create_test_transaction("relayer-1");
550 tx.id = "pending-tx-1".to_string();
551 Ok(crate::repositories::PaginatedResult {
552 items: vec![tx],
553 total: 1,
554 page: 1,
555 per_page: 1,
556 })
557 });
558
559 mocks
561 .job_producer
562 .expect_produce_transaction_request_job()
563 .withf(|job, delay| job.transaction_id == "pending-tx-1" && delay.is_none())
564 .times(1)
565 .returning(|_, _| Box::pin(async { Ok(()) }));
566
567 let handler = make_stellar_tx_handler(relayer, mocks);
568
569 let result = handler
570 .enqueue_next_pending_transaction("finished-tx")
571 .await;
572 assert!(result.is_ok());
573 }
574
575 #[tokio::test]
576 async fn enqueue_next_pending_transaction_no_pending_tx() {
577 let relayer = create_test_relayer();
578 let mut mocks = default_test_mocks();
579
580 mocks
582 .tx_repo
583 .expect_find_by_status_paginated()
584 .withf(|_relayer_id, statuses, query, oldest_first| {
585 statuses == [TransactionStatus::Pending]
586 && query.page == 1
587 && query.per_page == 1
588 && *oldest_first
589 })
590 .times(1)
591 .returning(|_, _, _, _| {
592 Ok(crate::repositories::PaginatedResult {
593 items: vec![],
594 total: 0,
595 page: 1,
596 per_page: 1,
597 })
598 });
599
600 let handler = make_stellar_tx_handler(relayer, mocks);
601
602 let result = handler
603 .enqueue_next_pending_transaction("finished-tx")
604 .await;
605 assert!(result.is_ok());
606 }
607
608 #[tokio::test]
609 async fn test_sync_sequence_from_chain() {
610 let relayer = create_test_relayer();
611 let mut mocks = default_test_mocks();
612
613 mocks
615 .provider
616 .expect_get_account()
617 .withf(|addr| addr == TEST_PK)
618 .times(1)
619 .returning(|_| {
620 Box::pin(async {
621 use soroban_rs::xdr::{
622 AccountEntry, AccountEntryExt, AccountId, PublicKey, SequenceNumber,
623 String32, Thresholds, Uint256,
624 };
625 use stellar_strkey::ed25519;
626
627 let pk = ed25519::PublicKey::from_string(TEST_PK).unwrap();
629 let account_id = AccountId(PublicKey::PublicKeyTypeEd25519(Uint256(pk.0)));
630
631 Ok(AccountEntry {
632 account_id,
633 balance: 1000000,
634 seq_num: SequenceNumber(100),
635 num_sub_entries: 0,
636 inflation_dest: None,
637 flags: 0,
638 home_domain: String32::default(),
639 thresholds: Thresholds([1, 1, 1, 1]),
640 signers: Default::default(),
641 ext: AccountEntryExt::V0,
642 })
643 })
644 });
645
646 mocks
648 .counter
649 .expect_sync_floor()
650 .withf(|relayer_id, addr, floor| {
651 relayer_id == "relayer-1" && addr == TEST_PK && *floor == 101
652 })
653 .times(1)
654 .returning(|_, _, floor| Box::pin(async move { Ok(floor) }));
655
656 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
657
658 let result = handler.sync_sequence_from_chain(&relayer.address).await;
659 assert!(result.is_ok());
660 }
661
662 #[tokio::test]
663 async fn test_sync_sequence_from_chain_provider_error() {
664 let relayer = create_test_relayer();
665 let mut mocks = default_test_mocks();
666
667 mocks.provider.expect_get_account().times(1).returning(|_| {
669 Box::pin(async { Err(ProviderError::Other("Account not found".to_string())) })
670 });
671
672 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
673
674 let result = handler.sync_sequence_from_chain(&relayer.address).await;
675 assert!(result.is_err());
676 match result.unwrap_err() {
677 TransactionError::UnexpectedError(msg) => {
678 assert!(msg.contains("Failed to fetch account from chain"));
679 }
680 _ => panic!("Expected UnexpectedError"),
681 }
682 }
683
684 #[tokio::test]
685 async fn test_sync_sequence_from_chain_counter_error() {
686 let relayer = create_test_relayer();
687 let mut mocks = default_test_mocks();
688
689 mocks.provider.expect_get_account().times(1).returning(|_| {
691 Box::pin(async {
692 use soroban_rs::xdr::{
693 AccountEntry, AccountEntryExt, AccountId, PublicKey, SequenceNumber, String32,
694 Thresholds, Uint256,
695 };
696 use stellar_strkey::ed25519;
697
698 let pk = ed25519::PublicKey::from_string(TEST_PK).unwrap();
700 let account_id = AccountId(PublicKey::PublicKeyTypeEd25519(Uint256(pk.0)));
701
702 Ok(AccountEntry {
703 account_id,
704 balance: 1000000,
705 seq_num: SequenceNumber(100),
706 num_sub_entries: 0,
707 inflation_dest: None,
708 flags: 0,
709 home_domain: String32::default(),
710 thresholds: Thresholds([1, 1, 1, 1]),
711 signers: Default::default(),
712 ext: AccountEntryExt::V0,
713 })
714 })
715 });
716
717 mocks
719 .counter
720 .expect_sync_floor()
721 .times(1)
722 .returning(|_, _, _| {
723 Box::pin(async {
724 Err(RepositoryError::Unknown(
725 "Counter update failed".to_string(),
726 ))
727 })
728 });
729
730 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
731
732 let result = handler.sync_sequence_from_chain(&relayer.address).await;
733 assert!(result.is_err());
734 match result.unwrap_err() {
735 TransactionError::UnexpectedError(msg) => {
736 assert!(msg.contains("Failed to update sequence counter"));
737 }
738 _ => panic!("Expected UnexpectedError"),
739 }
740 }
741
742 #[test]
743 fn test_concurrent_transactions_enabled() {
744 let mut relayer = create_test_relayer();
746 if let RelayerNetworkPolicy::Stellar(ref mut policy) = relayer.policies {
747 policy.concurrent_transactions = Some(true);
748 }
749 let mocks = default_test_mocks();
750 let handler = make_stellar_tx_handler(relayer, mocks);
751 assert!(handler.concurrent_transactions_enabled());
752
753 let mut relayer = create_test_relayer();
755 if let RelayerNetworkPolicy::Stellar(ref mut policy) = relayer.policies {
756 policy.concurrent_transactions = Some(false);
757 }
758 let mocks = default_test_mocks();
759 let handler = make_stellar_tx_handler(relayer, mocks);
760 assert!(!handler.concurrent_transactions_enabled());
761
762 let relayer = create_test_relayer();
764 let mocks = default_test_mocks();
765 let handler = make_stellar_tx_handler(relayer, mocks);
766 assert_eq!(
767 handler.concurrent_transactions_enabled(),
768 DEFAULT_STELLAR_CONCURRENT_TRANSACTIONS
769 );
770 }
771
772 #[tokio::test]
773 async fn test_enqueue_next_pending_transaction_with_concurrency_enabled() {
774 let mut relayer = create_test_relayer();
776 if let RelayerNetworkPolicy::Stellar(ref mut policy) = relayer.policies {
777 policy.concurrent_transactions = Some(true);
778 }
779 let mut mocks = default_test_mocks();
780
781 mocks.tx_repo.expect_find_by_status_paginated().times(0); mocks
786 .job_producer
787 .expect_produce_transaction_request_job()
788 .times(0); let handler = make_stellar_tx_handler(relayer, mocks);
791
792 let result = handler
793 .enqueue_next_pending_transaction("finished-tx")
794 .await;
795 assert!(result.is_ok());
796 }
797
798 #[tokio::test]
799 async fn test_reset_transaction_for_retry() {
800 let relayer = create_test_relayer();
801 let mut mocks = default_test_mocks();
802
803 let mut tx = create_test_transaction(&relayer.id);
805 if let NetworkTransactionData::Stellar(ref mut data) = tx.network_data {
806 data.sequence_number = Some(42);
807 data.signatures.push(dummy_signature());
808 data.hash = Some("test-hash".to_string());
809 data.signed_envelope_xdr = Some("test-xdr".to_string());
810 }
811
812 mocks
814 .tx_repo
815 .expect_partial_update()
816 .withf(|tx_id, upd| {
817 tx_id == "tx-1"
818 && upd.status == Some(TransactionStatus::Pending)
819 && upd.sent_at.is_none()
820 && upd.confirmed_at.is_none()
821 })
822 .times(1)
823 .returning(|id, upd| {
824 let mut tx = create_test_transaction("relayer-1");
825 tx.id = id;
826 tx.status = upd.status.unwrap();
827 if let Some(network_data) = upd.network_data {
828 tx.network_data = network_data;
829 }
830 Ok::<_, RepositoryError>(tx)
831 });
832
833 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
834
835 let result = handler.reset_transaction_for_retry(tx).await;
836 assert!(result.is_ok());
837
838 let reset_tx = result.unwrap();
839 assert_eq!(reset_tx.status, TransactionStatus::Pending);
840
841 if let NetworkTransactionData::Stellar(data) = &reset_tx.network_data {
843 assert!(data.sequence_number.is_none());
844 assert!(data.signatures.is_empty());
845 assert!(data.hash.is_none());
846 assert!(data.signed_envelope_xdr.is_none());
847 } else {
848 panic!("Expected Stellar transaction data");
849 }
850 }
851
852 #[tokio::test]
853 #[ignore = "Requires active Redis instance"]
854 async fn test_find_oldest_pending_for_relayer_with_redis() {
855 let redis_url = std::env::var("REDIS_TEST_URL")
857 .unwrap_or_else(|_| "redis://127.0.0.1:6379".to_string());
858 let pool = Arc::new(
859 Config::from_url(&redis_url)
860 .builder()
861 .expect("Failed to create Redis pool builder")
862 .max_size(16)
863 .runtime(Runtime::Tokio1)
864 .build()
865 .expect("Failed to build Redis pool"),
866 );
867 let connections = Arc::new(RedisConnections::new_single_pool(pool));
868
869 let random_id = Uuid::new_v4().to_string();
870 let key_prefix = format!("test_stellar:{random_id}");
871 let tx_repo = Arc::new(
872 RedisTransactionRepository::new(connections, key_prefix)
873 .expect("Failed to create RedisTransactionRepository"),
874 );
875
876 let relayer_id = format!("relayer-{}", Uuid::new_v4());
877
878 let mut tx1 = create_test_transaction(&relayer_id);
881 tx1.id = format!("tx-1-{}", Uuid::new_v4());
882 tx1.status = TransactionStatus::Pending;
883 tx1.created_at = "2025-01-27T10:00:00.000000+00:00".to_string();
884
885 let mut tx2 = create_test_transaction(&relayer_id);
887 tx2.id = format!("tx-2-{}", Uuid::new_v4());
888 tx2.status = TransactionStatus::Pending;
889 tx2.created_at = "2025-01-27T11:00:00.000000+00:00".to_string();
890
891 let mut tx3 = create_test_transaction(&relayer_id);
893 tx3.id = format!("tx-3-{}", Uuid::new_v4());
894 tx3.status = TransactionStatus::Pending;
895 tx3.created_at = "2025-01-27T12:00:00.000000+00:00".to_string();
896
897 tx_repo.create(tx1.clone()).await.unwrap();
899 tx_repo.create(tx2.clone()).await.unwrap();
900 tx_repo.create(tx3.clone()).await.unwrap();
901
902 let relayer = create_test_relayer();
905 let mut relayer_model = relayer.clone();
906 relayer_model.id = relayer_id.clone();
907
908 let mocks = default_test_mocks();
909 let handler = StellarRelayerTransaction::new(
910 relayer_model,
911 Arc::new(mocks.relayer_repo),
912 tx_repo.clone(),
913 Arc::new(mocks.job_producer),
914 Arc::new(mocks.signer),
915 mocks.provider,
916 Arc::new(mocks.counter),
917 Arc::new(mocks.dex_service),
918 )
919 .unwrap();
920
921 let result = handler
923 .find_oldest_pending_for_relayer(&relayer_id)
924 .await
925 .unwrap();
926
927 assert!(result.is_some(), "Should find a pending transaction");
929 let found_tx = result.unwrap();
930
931 assert_eq!(
932 found_tx.id, tx1.id,
933 "Should get oldest transaction (tx1) - oldest_first=true so .next() yields oldest"
934 );
935 assert_eq!(
936 found_tx.created_at, tx1.created_at,
937 "Should match oldest transaction's created_at"
938 );
939
940 let _ = tx_repo.delete_by_id(tx1.id.clone()).await;
942 let _ = tx_repo.delete_by_id(tx2.id.clone()).await;
943 let _ = tx_repo.delete_by_id(tx3.id.clone()).await;
944 }
945
946 #[tokio::test]
947 async fn test_find_oldest_pending_for_relayer_with_in_memory() {
948 use crate::repositories::transaction::InMemoryTransactionRepository;
949 use uuid::Uuid;
950
951 let tx_repo = Arc::new(InMemoryTransactionRepository::new());
953
954 let relayer_id = format!("relayer-{}", Uuid::new_v4());
955
956 let mut tx1 = create_test_transaction(&relayer_id);
959 tx1.id = format!("tx-1-{}", Uuid::new_v4());
960 tx1.status = TransactionStatus::Pending;
961 tx1.created_at = "2025-01-27T10:00:00.000000+00:00".to_string();
962
963 let mut tx2 = create_test_transaction(&relayer_id);
965 tx2.id = format!("tx-2-{}", Uuid::new_v4());
966 tx2.status = TransactionStatus::Pending;
967 tx2.created_at = "2025-01-27T11:00:00.000000+00:00".to_string();
968
969 let mut tx3 = create_test_transaction(&relayer_id);
971 tx3.id = format!("tx-3-{}", Uuid::new_v4());
972 tx3.status = TransactionStatus::Pending;
973 tx3.created_at = "2025-01-27T12:00:00.000000+00:00".to_string();
974
975 tx_repo.create(tx1.clone()).await.unwrap();
977 tx_repo.create(tx2.clone()).await.unwrap();
978 tx_repo.create(tx3.clone()).await.unwrap();
979
980 let relayer = create_test_relayer();
983 let mut relayer_model = relayer.clone();
984 relayer_model.id = relayer_id.clone();
985
986 let mocks = default_test_mocks();
987 let handler = StellarRelayerTransaction::new(
988 relayer_model,
989 Arc::new(mocks.relayer_repo),
990 tx_repo.clone(),
991 Arc::new(mocks.job_producer),
992 Arc::new(mocks.signer),
993 mocks.provider,
994 Arc::new(mocks.counter),
995 Arc::new(mocks.dex_service),
996 )
997 .unwrap();
998
999 let result = handler
1001 .find_oldest_pending_for_relayer(&relayer_id)
1002 .await
1003 .unwrap();
1004
1005 assert!(result.is_some(), "Should find a pending transaction");
1007 let found_tx = result.unwrap();
1008
1009 assert_eq!(
1011 found_tx.id, tx1.id,
1012 "Should get oldest transaction (tx1) - oldest_first=true so .next() yields oldest"
1013 );
1014 assert_eq!(
1015 found_tx.created_at, tx1.created_at,
1016 "Should match oldest transaction's created_at"
1017 );
1018 }
1019}