1use chrono::Utc;
6use tracing::{debug, info, warn};
7
8use super::{
9 is_final_state,
10 utils::{decode_transaction_result_code, is_bad_sequence_error, is_insufficient_fee_error},
11 StellarRelayerTransaction,
12};
13use crate::{
14 constants::{
15 STELLAR_FAST_RESUBMIT_BASE_DELAY_SECONDS, STELLAR_INSUFFICIENT_FEE_MAX_RETRIES,
16 STELLAR_TRY_AGAIN_LATER_FAST_RETRIES,
17 },
18 domain::transaction::stellar::prepare::common::send_submit_transaction_job,
19 jobs::JobProducerTrait,
20 metrics::{STELLAR_SUBMISSION_FAILURES, TRANSACTIONS_INSUFFICIENT_FEE},
21 models::{
22 NetworkTransactionData, RelayerRepoModel, TransactionError, TransactionRepoModel,
23 TransactionStatus, TransactionUpdateRequest,
24 },
25 repositories::{Repository, TransactionCounterTrait, TransactionRepository},
26 services::{
27 provider::StellarProviderTrait,
28 signer::{Signer, StellarSignTrait},
29 },
30};
31
32impl<R, T, J, S, P, C, D> StellarRelayerTransaction<R, T, J, S, P, C, D>
33where
34 R: Repository<RelayerRepoModel, String> + Send + Sync,
35 T: TransactionRepository + Send + Sync,
36 J: JobProducerTrait + Send + Sync,
37 S: Signer + StellarSignTrait + Send + Sync,
38 P: StellarProviderTrait + Send + Sync,
39 C: TransactionCounterTrait + Send + Sync,
40 D: crate::services::stellar_dex::StellarDexServiceTrait + Send + Sync + 'static,
41{
42 pub async fn submit_transaction_impl(
45 &self,
46 tx: TransactionRepoModel,
47 ) -> Result<TransactionRepoModel, TransactionError> {
48 info!(
49 tx_id = %tx.id,
50 relayer_id = %tx.relayer_id,
51 status = ?tx.status,
52 "submitting stellar transaction"
53 );
54
55 if is_final_state(&tx.status) {
57 warn!(
58 tx_id = %tx.id,
59 relayer_id = %tx.relayer_id,
60 status = ?tx.status,
61 "transaction already in final state, skipping submission"
62 );
63 return Ok(tx);
64 }
65
66 if self.is_transaction_expired(&tx)? {
68 info!(
69 tx_id = %tx.id,
70 relayer_id = %tx.relayer_id,
71 valid_until = ?tx.valid_until,
72 "transaction has expired, marking as Expired"
73 );
74 return self
75 .mark_as_expired(tx, "Transaction time_bounds expired".to_string())
76 .await;
77 }
78
79 match self.submit_core(tx.clone()).await {
81 Ok(submitted_tx) => Ok(submitted_tx),
82 Err(error) => {
83 self.handle_submit_failure(tx, error).await
85 }
86 }
87 }
88
89 async fn submit_core(
104 &self,
105 tx: TransactionRepoModel,
106 ) -> Result<TransactionRepoModel, TransactionError> {
107 let stellar_data = tx.network_data.get_stellar_transaction_data()?;
108 let tx_envelope = stellar_data
109 .get_envelope_for_submission()
110 .map_err(TransactionError::from)?;
111
112 let response = self
114 .provider()
115 .send_transaction_with_status(&tx_envelope)
116 .await
117 .map_err(|e| {
118 STELLAR_SUBMISSION_FAILURES
119 .with_label_values(&["provider_error", "n/a"])
120 .inc();
121 TransactionError::from(e)
122 })?;
123
124 match response.status.as_str() {
126 "PENDING" | "DUPLICATE" => {
127 if response.status == "DUPLICATE" {
129 info!(
130 tx_id = %tx.id,
131 relayer_id = %tx.relayer_id,
132 hash = %response.hash,
133 "transaction already submitted (DUPLICATE status)"
134 );
135 }
136 let tx_hash_hex = response.hash.clone();
137 let updated_stellar_data = stellar_data.with_hash(tx_hash_hex.clone());
138
139 let mut hashes = tx.hashes.clone();
140 if !hashes.contains(&tx_hash_hex) {
141 hashes.push(tx_hash_hex);
142 }
143
144 let update_req = TransactionUpdateRequest {
145 status: Some(TransactionStatus::Submitted),
146 sent_at: Some(Utc::now().to_rfc3339()),
147 network_data: Some(NetworkTransactionData::Stellar(updated_stellar_data)),
148 hashes: Some(hashes),
149 ..Default::default()
150 };
151
152 let updated_tx = self
153 .transaction_repository()
154 .partial_update(tx.id.clone(), update_req)
155 .await?;
156
157 if response.status == "PENDING" {
159 info!(
160 tx_id = %tx.id,
161 relayer_id = %tx.relayer_id,
162 "sending transaction update notification for pending transaction"
163 );
164 self.send_transaction_update_notification(&updated_tx).await;
165 }
166
167 Ok(updated_tx)
168 }
169 "TRY_AGAIN_LATER" => {
170 let updated_tx = self
177 .transaction_repository()
178 .record_stellar_try_again_later_retry(tx.id.clone(), Utc::now().to_rfc3339())
179 .await?;
180
181 if is_final_state(&updated_tx.status) {
185 debug!(
186 tx_id = %updated_tx.id,
187 relayer_id = %updated_tx.relayer_id,
188 status = ?updated_tx.status,
189 "transaction reached final state during retry recording; skipping fast resubmit"
190 );
191 return Ok(updated_tx);
192 }
193
194 let retries = updated_tx
195 .metadata
196 .as_ref()
197 .map_or(0, |m| m.try_again_later_retries);
198
199 if retries == 1 {
201 crate::metrics::STELLAR_TRY_AGAIN_LATER
202 .with_label_values(&[&tx.relayer_id, &tx.status.to_string()])
203 .inc();
204 }
205
206 if retries <= STELLAR_TRY_AGAIN_LATER_FAST_RETRIES {
207 let delay_seconds =
208 STELLAR_FAST_RESUBMIT_BASE_DELAY_SECONDS * i64::from(retries);
209
210 info!(
211 tx_id = %updated_tx.id,
212 relayer_id = %updated_tx.relayer_id,
213 status = ?updated_tx.status,
214 try_again_later_retries = retries,
215 delay_seconds,
216 "enqueueing fast resubmit after TRY_AGAIN_LATER"
217 );
218
219 if let Err(error) = send_submit_transaction_job(
220 self.job_producer(),
221 &updated_tx,
222 Some(delay_seconds),
223 )
224 .await
225 {
226 warn!(
227 tx_id = %updated_tx.id,
228 relayer_id = %updated_tx.relayer_id,
229 error = %error,
230 try_again_later_retries = retries,
231 delay_seconds,
232 "failed to enqueue fast resubmit after TRY_AGAIN_LATER"
233 );
234 }
235 } else {
236 debug!(
237 tx_id = %tx.id,
238 relayer_id = %tx.relayer_id,
239 status = ?tx.status,
240 try_again_later_retries = retries,
241 "TRY_AGAIN_LATER — status checker will retry"
242 );
243 }
244 Ok(updated_tx)
245 }
246 "ERROR" => {
247 let error_detail = response
249 .error_result_xdr
250 .unwrap_or_else(|| "No error details provided".to_string());
251 let decoded_result_code = decode_transaction_result_code(&error_detail);
252
253 if decoded_result_code
257 .as_deref()
258 .is_some_and(is_insufficient_fee_error)
259 {
260 let mut meta = tx.metadata.clone().unwrap_or_default();
261 meta.insufficient_fee_retries = meta.insufficient_fee_retries.saturating_add(1);
262
263 if meta.insufficient_fee_retries == 1 {
265 TRANSACTIONS_INSUFFICIENT_FEE
266 .with_label_values(&[tx.relayer_id.as_str(), "stellar"])
267 .inc();
268 }
269
270 if meta.insufficient_fee_retries > STELLAR_INSUFFICIENT_FEE_MAX_RETRIES {
271 STELLAR_SUBMISSION_FAILURES
272 .with_label_values(&["error", "tx_insufficient_fee"])
273 .inc();
274 return Err(TransactionError::UnexpectedError(format!(
275 "Transaction submission error: insufficient fee retry limit exceeded ({STELLAR_INSUFFICIENT_FEE_MAX_RETRIES})"
276 )));
277 }
278
279 let updated_tx = self
281 .transaction_repository()
282 .record_stellar_insufficient_fee_retry(
283 tx.id.clone(),
284 Utc::now().to_rfc3339(),
285 )
286 .await?;
287
288 if is_final_state(&updated_tx.status) {
292 debug!(
293 tx_id = %updated_tx.id,
294 relayer_id = %updated_tx.relayer_id,
295 status = ?updated_tx.status,
296 "transaction reached final state during retry recording; skipping fast resubmit"
297 );
298 return Ok(updated_tx);
299 }
300
301 let retries = updated_tx
302 .metadata
303 .as_ref()
304 .map_or(0, |metadata| metadata.insufficient_fee_retries);
305
306 if retries > STELLAR_INSUFFICIENT_FEE_MAX_RETRIES {
310 warn!(
311 tx_id = %updated_tx.id,
312 relayer_id = %updated_tx.relayer_id,
313 insufficient_fee_retries = retries,
314 "post-record retry count exceeds cap; skipping fast resubmit - concurrent submission detected"
315 );
316 return Ok(updated_tx);
317 }
318
319 let delay_seconds =
320 STELLAR_FAST_RESUBMIT_BASE_DELAY_SECONDS * i64::from(retries);
321
322 info!(
323 tx_id = %updated_tx.id,
324 relayer_id = %updated_tx.relayer_id,
325 status = ?updated_tx.status,
326 insufficient_fee_retries = retries,
327 delay_seconds,
328 result_code = decoded_result_code.as_deref().unwrap_or("Unknown"),
329 "enqueueing fast resubmit after insufficient fee"
330 );
331
332 if let Err(error) = send_submit_transaction_job(
333 self.job_producer(),
334 &updated_tx,
335 Some(delay_seconds),
336 )
337 .await
338 {
339 warn!(
340 tx_id = %updated_tx.id,
341 relayer_id = %updated_tx.relayer_id,
342 error = %error,
343 insufficient_fee_retries = retries,
344 delay_seconds,
345 "failed to enqueue fast resubmit after insufficient fee"
346 );
347 }
348 return Ok(updated_tx);
349 }
350 STELLAR_SUBMISSION_FAILURES
351 .with_label_values(&[
352 "error",
353 decoded_result_code.as_deref().unwrap_or("unknown"),
354 ])
355 .inc();
356 Err(TransactionError::UnexpectedError(format!(
357 "Transaction submission error: {}",
358 decoded_result_code.unwrap_or(error_detail)
359 )))
360 }
361 unknown => {
362 STELLAR_SUBMISSION_FAILURES
364 .with_label_values(&["unknown_status", "n/a"])
365 .inc();
366 warn!(
367 tx_id = %tx.id,
368 relayer_id = %tx.relayer_id,
369 status = %unknown,
370 "received unknown transaction status from RPC"
371 );
372 Err(TransactionError::UnexpectedError(format!(
373 "Unknown transaction status: {unknown}"
374 )))
375 }
376 }
377 }
378
379 async fn handle_submit_failure(
382 &self,
383 tx: TransactionRepoModel,
384 error: TransactionError,
385 ) -> Result<TransactionRepoModel, TransactionError> {
386 let error_reason = format!("Submission failed: {error}");
387 let tx_id = tx.id.clone();
388 let relayer_id = tx.relayer_id.clone();
389 warn!(
390 tx_id = %tx_id,
391 relayer_id = %relayer_id,
392 reason = %error_reason,
393 "transaction submission failed"
394 );
395
396 if error.is_concurrent_update_conflict() {
401 info!(
402 tx_id = %tx_id,
403 relayer_id = %relayer_id,
404 "concurrent transaction update detected during submission, reloading latest state"
405 );
406 return self
407 .transaction_repository()
408 .get_by_id(tx_id)
409 .await
410 .map_err(TransactionError::from);
411 }
412
413 if is_bad_sequence_error(&error_reason) {
414 if let Ok(stellar_data) = tx.network_data.get_stellar_transaction_data() {
416 info!(
417 tx_id = %tx_id,
418 relayer_id = %relayer_id,
419 "syncing sequence from chain after bad sequence error"
420 );
421 match self
422 .sync_sequence_from_chain(&stellar_data.source_account)
423 .await
424 {
425 Ok(()) => {
426 info!(
427 tx_id = %tx_id,
428 relayer_id = %relayer_id,
429 "successfully synced sequence from chain"
430 );
431 }
432 Err(sync_error) => {
433 warn!(
434 tx_id = %tx_id,
435 relayer_id = %relayer_id,
436 error = %sync_error,
437 "failed to sync sequence from chain"
438 );
439 }
440 }
441 }
442
443 info!(
446 tx_id = %tx_id,
447 relayer_id = %relayer_id,
448 "bad sequence error detected, resetting transaction to pending state"
449 );
450 match self.reset_transaction_for_retry(tx.clone()).await {
451 Ok(reset_tx) => {
452 info!(
453 tx_id = %tx_id,
454 relayer_id = %relayer_id,
455 "transaction reset to pending, status check will handle resubmission"
456 );
457 return Ok(reset_tx);
461 }
462 Err(reset_error) => {
463 warn!(
464 tx_id = %tx_id,
465 relayer_id = %relayer_id,
466 error = %reset_error,
467 "failed to reset transaction for retry"
468 );
469 }
471 }
472 }
473
474 let update_request = TransactionUpdateRequest {
477 status: Some(TransactionStatus::Failed),
478 status_reason: Some(error_reason.clone()),
479 ..Default::default()
480 };
481 let failed_tx = match self
482 .finalize_transaction_state(tx_id.clone(), update_request)
483 .await
484 {
485 Ok(updated_tx) => updated_tx,
486 Err(finalize_error) => {
487 warn!(
488 tx_id = %tx_id,
489 relayer_id = %relayer_id,
490 error = %finalize_error,
491 "failed to mark transaction as failed, continuing with lane cleanup"
492 );
493 return Err(error);
496 }
497 };
498
499 if let Err(enqueue_error) = self.enqueue_next_pending_transaction(&tx_id).await {
501 warn!(
502 tx_id = %tx_id,
503 relayer_id = %relayer_id,
504 error = %enqueue_error,
505 "failed to enqueue next pending transaction after submission failure"
506 );
507 }
508
509 info!(
510 tx_id = %tx_id,
511 relayer_id = %relayer_id,
512 error = %error_reason,
513 "transaction submission failure handled, marked as failed"
514 );
515
516 Ok(failed_tx)
520 }
521
522 pub async fn resubmit_transaction_impl(
524 &self,
525 tx: TransactionRepoModel,
526 ) -> Result<TransactionRepoModel, TransactionError> {
527 self.submit_transaction_impl(tx).await
528 }
529}
530
531#[cfg(test)]
532mod tests {
533 use super::*;
534 use soroban_rs::stellar_rpc_client::SendTransactionResponse;
535 use soroban_rs::xdr::WriteXdr;
536
537 use crate::domain::transaction::stellar::test_helpers::*;
538 use crate::models::TransactionMetadata;
539
540 fn create_send_tx_response(status: &str, hash: &str) -> SendTransactionResponse {
542 SendTransactionResponse {
543 status: status.to_string(),
544 hash: hash.to_string(),
545 error_result_xdr: None,
546 latest_ledger: 100,
547 latest_ledger_close_time: 1700000000,
548 }
549 }
550
551 mod submit_transaction_tests {
552 use crate::{
553 models::RepositoryError, repositories::PaginatedResult,
554 services::provider::ProviderError,
555 };
556
557 use super::*;
558
559 #[tokio::test]
560 async fn submit_transaction_happy_path() {
561 let relayer = create_test_relayer();
562 let mut mocks = default_test_mocks();
563
564 let response = create_send_tx_response(
566 "PENDING",
567 "0101010101010101010101010101010101010101010101010101010101010101",
568 );
569 mocks
570 .provider
571 .expect_send_transaction_with_status()
572 .returning(move |_| {
573 let r = response.clone();
574 Box::pin(async move { Ok(r) })
575 });
576
577 mocks
579 .tx_repo
580 .expect_partial_update()
581 .withf(|_, upd| upd.status == Some(TransactionStatus::Submitted))
582 .returning(|id, upd| {
583 let mut tx = create_test_transaction("relayer-1");
584 tx.id = id;
585 tx.status = upd.status.unwrap();
586 Ok::<_, RepositoryError>(tx)
587 });
588
589 mocks
591 .job_producer
592 .expect_produce_send_notification_job()
593 .times(1)
594 .returning(|_, _| Box::pin(async { Ok(()) }));
595
596 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
597
598 let mut tx = create_test_transaction(&relayer.id);
599 tx.status = TransactionStatus::Sent; if let NetworkTransactionData::Stellar(ref mut d) = tx.network_data {
601 d.signatures.push(dummy_signature());
602 d.signed_envelope_xdr = Some(create_signed_xdr(TEST_PK, TEST_PK_2));
603 }
605
606 let res = handler.submit_transaction_impl(tx).await.unwrap();
607 assert_eq!(res.status, TransactionStatus::Submitted);
608 }
609
610 #[tokio::test]
611 async fn submit_transaction_provider_error_marks_failed() {
612 let relayer = create_test_relayer();
613 let mut mocks = default_test_mocks();
614
615 mocks
617 .provider
618 .expect_send_transaction_with_status()
619 .returning(|_| {
620 Box::pin(async { Err(ProviderError::Other("Network error".to_string())) })
621 });
622
623 mocks
625 .tx_repo
626 .expect_partial_update()
627 .withf(|_, upd| upd.status == Some(TransactionStatus::Failed))
628 .returning(|id, upd| {
629 let mut tx = create_test_transaction("relayer-1");
630 tx.id = id;
631 tx.status = upd.status.unwrap();
632 Ok::<_, RepositoryError>(tx)
633 });
634
635 mocks
637 .job_producer
638 .expect_produce_send_notification_job()
639 .times(1)
640 .returning(|_, _| Box::pin(async { Ok(()) }));
641
642 mocks
644 .tx_repo
645 .expect_find_by_status_paginated()
646 .returning(move |_, _, _, _| {
647 Ok(PaginatedResult {
648 items: vec![],
649 total: 0,
650 page: 1,
651 per_page: 1,
652 })
653 }); let handler = make_stellar_tx_handler(relayer.clone(), mocks);
656 let mut tx = create_test_transaction(&relayer.id);
657 tx.status = TransactionStatus::Sent; if let NetworkTransactionData::Stellar(ref mut data) = tx.network_data {
659 data.signatures.push(dummy_signature());
660 data.sequence_number = Some(42); data.signed_envelope_xdr = Some("test-xdr".to_string()); }
663
664 let res = handler.submit_transaction_impl(tx).await;
665
666 let failed_tx = res.unwrap();
668 assert_eq!(failed_tx.status, TransactionStatus::Failed);
669 }
670
671 #[tokio::test]
672 async fn submit_transaction_repository_error_marks_failed() {
673 let relayer = create_test_relayer();
674 let mut mocks = default_test_mocks();
675
676 let response = create_send_tx_response(
678 "PENDING",
679 "0101010101010101010101010101010101010101010101010101010101010101",
680 );
681 mocks
682 .provider
683 .expect_send_transaction_with_status()
684 .returning(move |_| {
685 let r = response.clone();
686 Box::pin(async move { Ok(r) })
687 });
688
689 mocks
691 .tx_repo
692 .expect_partial_update()
693 .withf(|_, upd| upd.status == Some(TransactionStatus::Submitted))
694 .returning(|_, _| Err(RepositoryError::Unknown("Database error".to_string())));
695
696 mocks
698 .tx_repo
699 .expect_partial_update()
700 .withf(|_, upd| upd.status == Some(TransactionStatus::Failed))
701 .returning(|id, upd| {
702 let mut tx = create_test_transaction("relayer-1");
703 tx.id = id;
704 tx.status = upd.status.unwrap();
705 Ok::<_, RepositoryError>(tx)
706 });
707
708 mocks
710 .job_producer
711 .expect_produce_send_notification_job()
712 .times(1)
713 .returning(|_, _| Box::pin(async { Ok(()) }));
714
715 mocks
717 .tx_repo
718 .expect_find_by_status_paginated()
719 .returning(move |_, _, _, _| {
720 Ok(PaginatedResult {
721 items: vec![],
722 total: 0,
723 page: 1,
724 per_page: 1,
725 })
726 }); let handler = make_stellar_tx_handler(relayer.clone(), mocks);
729 let mut tx = create_test_transaction(&relayer.id);
730 tx.status = TransactionStatus::Sent; if let NetworkTransactionData::Stellar(ref mut data) = tx.network_data {
732 data.signatures.push(dummy_signature());
733 data.sequence_number = Some(42); data.signed_envelope_xdr = Some("test-xdr".to_string()); }
736
737 let res = handler.submit_transaction_impl(tx).await;
738
739 let failed_tx = res.unwrap();
742 assert_eq!(failed_tx.status, TransactionStatus::Failed);
743 }
744
745 #[tokio::test]
746 async fn submit_transaction_uses_signed_envelope_xdr() {
747 let relayer = create_test_relayer();
748 let mut mocks = default_test_mocks();
749
750 let mut tx = create_test_transaction(&relayer.id);
752 tx.status = TransactionStatus::Sent; if let NetworkTransactionData::Stellar(ref mut data) = tx.network_data {
754 data.signatures.push(dummy_signature());
755 let envelope = data.get_envelope_for_submission().unwrap();
757 let xdr = envelope
758 .to_xdr_base64(soroban_rs::xdr::Limits::none())
759 .unwrap();
760 data.signed_envelope_xdr = Some(xdr);
761 }
762
763 let response = create_send_tx_response(
765 "PENDING",
766 "0202020202020202020202020202020202020202020202020202020202020202",
767 );
768 mocks
769 .provider
770 .expect_send_transaction_with_status()
771 .returning(move |_| {
772 let r = response.clone();
773 Box::pin(async move { Ok(r) })
774 });
775
776 mocks
778 .tx_repo
779 .expect_partial_update()
780 .withf(|_, upd| upd.status == Some(TransactionStatus::Submitted))
781 .returning(|id, upd| {
782 let mut tx = create_test_transaction("relayer-1");
783 tx.id = id;
784 tx.status = upd.status.unwrap();
785 Ok::<_, RepositoryError>(tx)
786 });
787
788 mocks
790 .job_producer
791 .expect_produce_send_notification_job()
792 .times(1)
793 .returning(|_, _| Box::pin(async { Ok(()) }));
794
795 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
796 let res = handler.submit_transaction_impl(tx).await.unwrap();
797
798 assert_eq!(res.status, TransactionStatus::Submitted);
799 }
800
801 #[tokio::test]
802 async fn resubmit_transaction_delegates_to_submit() {
803 let relayer = create_test_relayer();
804 let mut mocks = default_test_mocks();
805
806 let response = create_send_tx_response(
808 "PENDING",
809 "0101010101010101010101010101010101010101010101010101010101010101",
810 );
811 mocks
812 .provider
813 .expect_send_transaction_with_status()
814 .returning(move |_| {
815 let r = response.clone();
816 Box::pin(async move { Ok(r) })
817 });
818
819 mocks
821 .tx_repo
822 .expect_partial_update()
823 .withf(|_, upd| upd.status == Some(TransactionStatus::Submitted))
824 .returning(|id, upd| {
825 let mut tx = create_test_transaction("relayer-1");
826 tx.id = id;
827 tx.status = upd.status.unwrap();
828 Ok::<_, RepositoryError>(tx)
829 });
830
831 mocks
833 .job_producer
834 .expect_produce_send_notification_job()
835 .times(1)
836 .returning(|_, _| Box::pin(async { Ok(()) }));
837
838 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
839
840 let mut tx = create_test_transaction(&relayer.id);
841 tx.status = TransactionStatus::Sent; if let NetworkTransactionData::Stellar(ref mut d) = tx.network_data {
843 d.signatures.push(dummy_signature());
844 d.signed_envelope_xdr = Some(create_signed_xdr(TEST_PK, TEST_PK_2));
845 }
847
848 let res = handler.resubmit_transaction_impl(tx).await.unwrap();
849 assert_eq!(res.status, TransactionStatus::Submitted);
850 }
851
852 #[tokio::test]
853 async fn submit_transaction_failure_enqueues_next_transaction() {
854 let relayer = create_test_relayer();
855 let mut mocks = default_test_mocks();
856
857 mocks
859 .provider
860 .expect_send_transaction_with_status()
861 .returning(|_| {
862 Box::pin(async { Err(ProviderError::Other("Network error".to_string())) })
863 });
864
865 mocks
869 .tx_repo
870 .expect_partial_update()
871 .withf(|_, upd| upd.status == Some(TransactionStatus::Failed))
872 .returning(|id, upd| {
873 let mut tx = create_test_transaction("relayer-1");
874 tx.id = id;
875 tx.status = upd.status.unwrap();
876 Ok::<_, RepositoryError>(tx)
877 });
878
879 mocks
881 .job_producer
882 .expect_produce_send_notification_job()
883 .times(1)
884 .returning(|_, _| Box::pin(async { Ok(()) }));
885
886 let mut pending_tx = create_test_transaction(&relayer.id);
888 pending_tx.id = "next-pending-tx".to_string();
889 pending_tx.status = TransactionStatus::Pending;
890 let captured_pending_tx = pending_tx.clone();
891 let relayer_id_clone = relayer.id.clone();
892 mocks
893 .tx_repo
894 .expect_find_by_status_paginated()
895 .withf(move |relayer_id, statuses, query, oldest_first| {
896 *relayer_id == relayer_id_clone
897 && statuses == [TransactionStatus::Pending]
898 && query.page == 1
899 && query.per_page == 1
900 && *oldest_first
901 })
902 .times(1)
903 .returning(move |_, _, _, _| {
904 Ok(PaginatedResult {
905 items: vec![captured_pending_tx.clone()],
906 total: 1,
907 page: 1,
908 per_page: 1,
909 })
910 });
911
912 mocks
914 .job_producer
915 .expect_produce_transaction_request_job()
916 .withf(move |job, _delay| job.transaction_id == "next-pending-tx")
917 .times(1)
918 .returning(|_, _| Box::pin(async { Ok(()) }));
919
920 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
921 let mut tx = create_test_transaction(&relayer.id);
922 tx.status = TransactionStatus::Sent; if let NetworkTransactionData::Stellar(ref mut data) = tx.network_data {
924 data.signatures.push(dummy_signature());
925 data.sequence_number = Some(42); data.signed_envelope_xdr = Some("test-xdr".to_string()); }
928
929 let res = handler.submit_transaction_impl(tx).await;
930
931 let failed_tx = res.unwrap();
933 assert_eq!(failed_tx.status, TransactionStatus::Failed);
934 }
935
936 #[tokio::test]
937 async fn test_submit_bad_sequence_resets_and_retries() {
938 let relayer = create_test_relayer();
939 let mut mocks = default_test_mocks();
940
941 mocks
943 .provider
944 .expect_send_transaction_with_status()
945 .returning(|_| {
946 Box::pin(async {
947 Err(ProviderError::Other(
948 "transaction submission failed: TxBadSeq".to_string(),
949 ))
950 })
951 });
952
953 mocks.provider.expect_get_account().times(1).returning(|_| {
955 Box::pin(async {
956 use soroban_rs::xdr::{
957 AccountEntry, AccountEntryExt, AccountId, PublicKey, SequenceNumber,
958 String32, Thresholds, Uint256,
959 };
960 use stellar_strkey::ed25519;
961
962 let pk = ed25519::PublicKey::from_string(TEST_PK).unwrap();
963 let account_id = AccountId(PublicKey::PublicKeyTypeEd25519(Uint256(pk.0)));
964
965 Ok(AccountEntry {
966 account_id,
967 balance: 1000000,
968 seq_num: SequenceNumber(100),
969 num_sub_entries: 0,
970 inflation_dest: None,
971 flags: 0,
972 home_domain: String32::default(),
973 thresholds: Thresholds([1, 1, 1, 1]),
974 signers: Default::default(),
975 ext: AccountEntryExt::V0,
976 })
977 })
978 });
979
980 mocks
982 .counter
983 .expect_sync_floor()
984 .times(1)
985 .returning(|_, _, floor| Box::pin(async move { Ok(floor) }));
986
987 mocks
989 .tx_repo
990 .expect_partial_update()
991 .withf(|_, upd| upd.status == Some(TransactionStatus::Pending))
992 .times(1)
993 .returning(|id, upd| {
994 let mut tx = create_test_transaction("relayer-1");
995 tx.id = id;
996 tx.status = upd.status.unwrap();
997 if let Some(network_data) = upd.network_data {
998 tx.network_data = network_data;
999 }
1000 Ok::<_, RepositoryError>(tx)
1001 });
1002
1003 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
1007 let mut tx = create_test_transaction(&relayer.id);
1008 tx.status = TransactionStatus::Sent; if let NetworkTransactionData::Stellar(ref mut data) = tx.network_data {
1010 data.signatures.push(dummy_signature());
1011 data.sequence_number = Some(42);
1012 data.signed_envelope_xdr = Some(create_signed_xdr(TEST_PK, TEST_PK_2));
1013 }
1015
1016 let result = handler.submit_transaction_impl(tx).await;
1017
1018 assert!(result.is_ok());
1020 let reset_tx = result.unwrap();
1021 assert_eq!(reset_tx.status, TransactionStatus::Pending);
1022
1023 if let NetworkTransactionData::Stellar(data) = &reset_tx.network_data {
1025 assert!(data.sequence_number.is_none());
1026 assert!(data.signatures.is_empty());
1027 assert!(data.hash.is_none());
1028 assert!(data.signed_envelope_xdr.is_none());
1029 } else {
1030 panic!("Expected Stellar transaction data");
1031 }
1032 }
1033
1034 #[tokio::test]
1035 async fn submit_transaction_duplicate_status_succeeds() {
1036 let relayer = create_test_relayer();
1037 let mut mocks = default_test_mocks();
1038
1039 let response = create_send_tx_response(
1041 "DUPLICATE",
1042 "0101010101010101010101010101010101010101010101010101010101010101",
1043 );
1044 mocks
1045 .provider
1046 .expect_send_transaction_with_status()
1047 .returning(move |_| {
1048 let r = response.clone();
1049 Box::pin(async move { Ok(r) })
1050 });
1051
1052 mocks
1054 .tx_repo
1055 .expect_partial_update()
1056 .withf(|_, upd| upd.status == Some(TransactionStatus::Submitted))
1057 .returning(|id, upd| {
1058 let mut tx = create_test_transaction("relayer-1");
1059 tx.id = id;
1060 tx.status = upd.status.unwrap();
1061 Ok::<_, RepositoryError>(tx)
1062 });
1063
1064 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
1065
1066 let mut tx = create_test_transaction(&relayer.id);
1067 tx.status = TransactionStatus::Sent;
1068 if let NetworkTransactionData::Stellar(ref mut d) = tx.network_data {
1069 d.signatures.push(dummy_signature());
1070 d.signed_envelope_xdr = Some(create_signed_xdr(TEST_PK, TEST_PK_2));
1071 }
1072
1073 let res = handler.submit_transaction_impl(tx).await.unwrap();
1074 assert_eq!(res.status, TransactionStatus::Submitted);
1075 }
1076
1077 #[tokio::test]
1078 async fn submit_transaction_try_again_later_keeps_tx_alive() {
1079 let relayer = create_test_relayer();
1080 let mut mocks = default_test_mocks();
1081
1082 let response = create_send_tx_response(
1084 "TRY_AGAIN_LATER",
1085 "0101010101010101010101010101010101010101010101010101010101010101",
1086 );
1087 mocks
1088 .provider
1089 .expect_send_transaction_with_status()
1090 .returning(move |_| {
1091 let r = response.clone();
1092 Box::pin(async move { Ok(r) })
1093 });
1094
1095 mocks
1096 .tx_repo
1097 .expect_record_stellar_try_again_later_retry()
1098 .withf(|id, sent_at| id == "tx-1" && !sent_at.is_empty())
1099 .returning(|id, _| {
1100 let mut tx = create_test_transaction("relayer-1");
1101 tx.id = id;
1102 tx.status = TransactionStatus::Sent;
1103 tx.metadata = Some(TransactionMetadata {
1104 consecutive_failures: 0,
1105 total_failures: 0,
1106 insufficient_fee_retries: 0,
1107 try_again_later_retries: 1,
1108 nonce_too_high_retries: 0,
1109 });
1110 Ok::<_, RepositoryError>(tx)
1111 });
1112
1113 mocks
1114 .job_producer
1115 .expect_produce_submit_transaction_job()
1116 .withf(|_, scheduled_on| scheduled_on.is_some())
1117 .times(1)
1118 .returning(|_, _| Box::pin(async { Ok(()) }));
1119
1120 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
1121 let mut tx = create_test_transaction(&relayer.id);
1122 tx.status = TransactionStatus::Sent;
1123 if let NetworkTransactionData::Stellar(ref mut data) = tx.network_data {
1124 data.signatures.push(dummy_signature());
1125 data.signed_envelope_xdr = Some(create_signed_xdr(TEST_PK, TEST_PK_2));
1126 }
1127
1128 let res = handler.submit_transaction_impl(tx).await;
1129
1130 let returned_tx = res.unwrap();
1132 assert_eq!(returned_tx.status, TransactionStatus::Sent);
1133 }
1134
1135 #[tokio::test]
1136 async fn submit_try_again_later_then_status_checker_reenqueues_submit() {
1137 let relayer = create_test_relayer();
1138
1139 let mut submit_mocks = default_test_mocks();
1141 let response = create_send_tx_response(
1142 "TRY_AGAIN_LATER",
1143 "0101010101010101010101010101010101010101010101010101010101010101",
1144 );
1145 submit_mocks
1146 .provider
1147 .expect_send_transaction_with_status()
1148 .times(1)
1149 .returning(move |_| {
1150 let r = response.clone();
1151 Box::pin(async move { Ok(r) })
1152 });
1153 submit_mocks
1154 .tx_repo
1155 .expect_record_stellar_try_again_later_retry()
1156 .withf(|id, sent_at| id == "tx-1" && !sent_at.is_empty())
1157 .times(1)
1158 .returning(|id, sent_at| {
1159 let mut tx = create_test_transaction("relayer-1");
1160 tx.id = id;
1161 tx.status = TransactionStatus::Sent;
1162 tx.sent_at = Some(sent_at);
1163 tx.metadata = Some(TransactionMetadata {
1164 consecutive_failures: 0,
1165 total_failures: 0,
1166 insufficient_fee_retries: 0,
1167 try_again_later_retries: 1,
1168 nonce_too_high_retries: 0,
1169 });
1170 Ok::<_, RepositoryError>(tx)
1171 });
1172
1173 submit_mocks
1174 .job_producer
1175 .expect_produce_submit_transaction_job()
1176 .withf(|_, scheduled_on| scheduled_on.is_some())
1177 .times(1)
1178 .returning(|_, _| Box::pin(async { Ok(()) }));
1179
1180 let submit_handler = make_stellar_tx_handler(relayer.clone(), submit_mocks);
1181 let mut sent_tx = create_test_transaction(&relayer.id);
1182 sent_tx.status = TransactionStatus::Sent;
1183 if let NetworkTransactionData::Stellar(ref mut data) = sent_tx.network_data {
1184 data.signatures.push(dummy_signature());
1185 data.signed_envelope_xdr = Some(create_signed_xdr(TEST_PK, TEST_PK_2));
1186 }
1187
1188 let mut returned_tx = submit_handler
1189 .submit_transaction_impl(sent_tx)
1190 .await
1191 .unwrap();
1192 assert_eq!(returned_tx.status, TransactionStatus::Sent);
1193 assert!(returned_tx.sent_at.is_some());
1194
1195 use crate::constants::STELLAR_RESUBMIT_BASE_INTERVAL_SECONDS;
1200 let buffer = 2;
1201 let created_at = (Utc::now()
1202 - chrono::Duration::seconds(STELLAR_RESUBMIT_BASE_INTERVAL_SECONDS + buffer))
1203 .to_rfc3339();
1204 let sent_at = (Utc::now()
1205 - chrono::Duration::seconds(STELLAR_RESUBMIT_BASE_INTERVAL_SECONDS + 1))
1206 .to_rfc3339();
1207 returned_tx.created_at = created_at;
1208 returned_tx.sent_at = Some(sent_at);
1209
1210 let mut status_mocks = default_test_mocks();
1211 status_mocks
1212 .job_producer
1213 .expect_produce_submit_transaction_job()
1214 .times(1)
1215 .returning(|_, _| Box::pin(async { Ok(()) }));
1216
1217 let status_handler = make_stellar_tx_handler(relayer.clone(), status_mocks);
1218 let status_result = status_handler
1219 .handle_transaction_status_impl(returned_tx, None)
1220 .await
1221 .unwrap();
1222 assert_eq!(status_result.status, TransactionStatus::Sent);
1223 }
1224
1225 #[tokio::test]
1226 async fn resubmit_try_again_later_returns_ok_for_submitted_tx() {
1227 let relayer = create_test_relayer();
1228 let mut mocks = default_test_mocks();
1229
1230 let response = create_send_tx_response(
1232 "TRY_AGAIN_LATER",
1233 "0101010101010101010101010101010101010101010101010101010101010101",
1234 );
1235 mocks
1236 .provider
1237 .expect_send_transaction_with_status()
1238 .returning(move |_| {
1239 let r = response.clone();
1240 Box::pin(async move { Ok(r) })
1241 });
1242
1243 mocks
1244 .tx_repo
1245 .expect_record_stellar_try_again_later_retry()
1246 .withf(|id, sent_at| id == "tx-1" && !sent_at.is_empty())
1247 .returning(|id, _| {
1248 let mut tx = create_test_transaction("relayer-1");
1249 tx.id = id;
1250 tx.status = TransactionStatus::Submitted;
1251 tx.metadata = Some(TransactionMetadata {
1252 consecutive_failures: 0,
1253 total_failures: 0,
1254 insufficient_fee_retries: 0,
1255 try_again_later_retries: 1,
1256 nonce_too_high_retries: 0,
1257 });
1258 Ok::<_, RepositoryError>(tx)
1259 });
1260
1261 mocks
1262 .job_producer
1263 .expect_produce_submit_transaction_job()
1264 .withf(|_, scheduled_on| scheduled_on.is_some())
1265 .times(1)
1266 .returning(|_, _| Box::pin(async { Ok(()) }));
1267
1268 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
1269 let mut tx = create_test_transaction(&relayer.id);
1270 tx.status = TransactionStatus::Submitted; if let NetworkTransactionData::Stellar(ref mut data) = tx.network_data {
1272 data.signatures.push(dummy_signature());
1273 data.signed_envelope_xdr = Some(create_signed_xdr(TEST_PK, TEST_PK_2));
1274 }
1275
1276 let res = handler.submit_transaction_impl(tx).await;
1277
1278 let returned_tx = res.unwrap();
1280 assert_eq!(returned_tx.status, TransactionStatus::Submitted);
1281 }
1282
1283 #[tokio::test]
1284 async fn submit_transaction_try_again_later_fast_window_closes() {
1285 let relayer = create_test_relayer();
1286 let mut mocks = default_test_mocks();
1287
1288 let response = create_send_tx_response(
1289 "TRY_AGAIN_LATER",
1290 "0101010101010101010101010101010101010101010101010101010101010101",
1291 );
1292 mocks
1293 .provider
1294 .expect_send_transaction_with_status()
1295 .returning(move |_| {
1296 let r = response.clone();
1297 Box::pin(async move { Ok(r) })
1298 });
1299
1300 mocks
1301 .tx_repo
1302 .expect_record_stellar_try_again_later_retry()
1303 .withf(|id, sent_at| id == "tx-1" && !sent_at.is_empty())
1304 .returning(|id, _| {
1305 let mut tx = create_test_transaction("relayer-1");
1306 tx.id = id;
1307 tx.status = TransactionStatus::Sent;
1308 tx.metadata = Some(TransactionMetadata {
1309 consecutive_failures: 0,
1310 total_failures: 0,
1311 insufficient_fee_retries: 0,
1312 try_again_later_retries: 4,
1313 nonce_too_high_retries: 0,
1314 });
1315 Ok::<_, RepositoryError>(tx)
1316 });
1317
1318 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
1319 let mut tx = create_test_transaction(&relayer.id);
1320 tx.status = TransactionStatus::Sent;
1321 if let NetworkTransactionData::Stellar(ref mut data) = tx.network_data {
1322 data.signatures.push(dummy_signature());
1323 data.signed_envelope_xdr = Some(create_signed_xdr(TEST_PK, TEST_PK_2));
1324 }
1325
1326 let res = handler.submit_transaction_impl(tx).await;
1327
1328 let returned_tx = res.unwrap();
1329 assert_eq!(returned_tx.status, TransactionStatus::Sent);
1330 assert_eq!(
1331 returned_tx
1332 .metadata
1333 .as_ref()
1334 .map(|metadata| metadata.try_again_later_retries),
1335 Some(4)
1336 );
1337 }
1338
1339 #[tokio::test]
1340 async fn submit_transaction_try_again_later_final_state_skips_enqueue() {
1341 let relayer = create_test_relayer();
1342 let mut mocks = default_test_mocks();
1343
1344 let response = create_send_tx_response(
1345 "TRY_AGAIN_LATER",
1346 "0101010101010101010101010101010101010101010101010101010101010101",
1347 );
1348 mocks
1349 .provider
1350 .expect_send_transaction_with_status()
1351 .returning(move |_| {
1352 let r = response.clone();
1353 Box::pin(async move { Ok(r) })
1354 });
1355
1356 mocks
1357 .tx_repo
1358 .expect_record_stellar_try_again_later_retry()
1359 .withf(|id, sent_at| id == "tx-1" && !sent_at.is_empty())
1360 .returning(|id, _| {
1361 let mut tx = create_test_transaction("relayer-1");
1362 tx.id = id;
1363 tx.status = TransactionStatus::Confirmed;
1364 tx.metadata = Some(TransactionMetadata {
1365 consecutive_failures: 0,
1366 total_failures: 0,
1367 insufficient_fee_retries: 0,
1368 try_again_later_retries: 0,
1369 nonce_too_high_retries: 0,
1370 });
1371 Ok::<_, RepositoryError>(tx)
1372 })
1373 .times(1);
1374
1375 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
1376 let mut tx = create_test_transaction(&relayer.id);
1377 tx.status = TransactionStatus::Sent;
1378 tx.metadata = Some(TransactionMetadata {
1379 consecutive_failures: 0,
1380 total_failures: 0,
1381 insufficient_fee_retries: 0,
1382 try_again_later_retries: 0,
1383 nonce_too_high_retries: 0,
1384 });
1385 if let NetworkTransactionData::Stellar(ref mut data) = tx.network_data {
1386 data.signatures.push(dummy_signature());
1387 data.signed_envelope_xdr = Some(create_signed_xdr(TEST_PK, TEST_PK_2));
1388 }
1389
1390 let res = handler.submit_core(tx).await;
1391
1392 assert!(res.is_ok());
1393 let returned_tx = res.unwrap();
1394 assert_eq!(returned_tx.status, TransactionStatus::Confirmed);
1395 }
1396
1397 #[tokio::test]
1398 async fn submit_transaction_error_status_fails() {
1399 let relayer = create_test_relayer();
1400 let mut mocks = default_test_mocks();
1401
1402 let mut response = create_send_tx_response(
1404 "ERROR",
1405 "0101010101010101010101010101010101010101010101010101010101010101",
1406 );
1407 response.error_result_xdr = Some("not-base64".to_string());
1408 mocks
1409 .provider
1410 .expect_send_transaction_with_status()
1411 .returning(move |_| {
1412 let r = response.clone();
1413 Box::pin(async move { Ok(r) })
1414 });
1415
1416 mocks
1418 .tx_repo
1419 .expect_partial_update()
1420 .withf(|_, upd| upd.status == Some(TransactionStatus::Failed))
1421 .returning(|id, upd| {
1422 let mut tx = create_test_transaction("relayer-1");
1423 tx.id = id;
1424 tx.status = upd.status.unwrap();
1425 Ok::<_, RepositoryError>(tx)
1426 });
1427
1428 mocks
1430 .job_producer
1431 .expect_produce_send_notification_job()
1432 .times(1)
1433 .returning(|_, _| Box::pin(async { Ok(()) }));
1434
1435 mocks
1437 .tx_repo
1438 .expect_find_by_status_paginated()
1439 .returning(move |_, _, _, _| {
1440 Ok(PaginatedResult {
1441 items: vec![],
1442 total: 0,
1443 page: 1,
1444 per_page: 1,
1445 })
1446 });
1447
1448 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
1449 let mut tx = create_test_transaction(&relayer.id);
1450 tx.status = TransactionStatus::Sent;
1451 if let NetworkTransactionData::Stellar(ref mut data) = tx.network_data {
1452 data.signatures.push(dummy_signature());
1453 data.signed_envelope_xdr = Some(create_signed_xdr(TEST_PK, TEST_PK_2));
1454 }
1455
1456 let res = handler.submit_transaction_impl(tx).await;
1457
1458 let failed_tx = res.unwrap();
1460 assert_eq!(failed_tx.status, TransactionStatus::Failed);
1461 }
1462
1463 #[tokio::test]
1464 async fn submit_transaction_insufficient_fee_keeps_tx_alive() {
1465 let relayer = create_test_relayer();
1466 let mut mocks = default_test_mocks();
1467
1468 let mut response = create_send_tx_response(
1470 "ERROR",
1471 "0101010101010101010101010101010101010101010101010101010101010101",
1472 );
1473 response.error_result_xdr = Some("AAAAAAAAY/n////3AAAAAA==".to_string());
1474 mocks
1475 .provider
1476 .expect_send_transaction_with_status()
1477 .returning(move |_| {
1478 let r = response.clone();
1479 Box::pin(async move { Ok(r) })
1480 });
1481
1482 mocks
1484 .tx_repo
1485 .expect_record_stellar_insufficient_fee_retry()
1486 .withf(|id, sent_at| id == "tx-1" && !sent_at.is_empty())
1487 .returning(|id, _| {
1488 let mut tx = create_test_transaction("relayer-1");
1489 tx.id = id;
1490 tx.status = TransactionStatus::Sent;
1491 tx.metadata = Some(TransactionMetadata {
1492 consecutive_failures: 0,
1493 total_failures: 0,
1494 insufficient_fee_retries: 1,
1495 try_again_later_retries: 0,
1496 nonce_too_high_retries: 0,
1497 });
1498 Ok::<_, RepositoryError>(tx)
1499 })
1500 .times(1);
1501
1502 mocks
1503 .job_producer
1504 .expect_produce_submit_transaction_job()
1505 .withf(|_, scheduled_on| scheduled_on.is_some())
1506 .times(1)
1507 .returning(|_, _| Box::pin(async { Ok(()) }));
1508
1509 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
1510 let mut tx = create_test_transaction(&relayer.id);
1511 tx.status = TransactionStatus::Sent;
1512 if let NetworkTransactionData::Stellar(ref mut data) = tx.network_data {
1513 data.signatures.push(dummy_signature());
1514 data.signed_envelope_xdr = Some(create_signed_xdr(TEST_PK, TEST_PK_2));
1515 }
1516
1517 let res = handler.submit_transaction_impl(tx).await;
1518
1519 let returned_tx = res.unwrap();
1521 assert_eq!(returned_tx.status, TransactionStatus::Sent);
1522 assert_eq!(
1523 returned_tx
1524 .metadata
1525 .as_ref()
1526 .map(|metadata| metadata.insufficient_fee_retries),
1527 Some(1)
1528 );
1529 }
1530
1531 #[tokio::test]
1532 async fn submit_transaction_insufficient_fee_escalates_delay() {
1533 let relayer = create_test_relayer();
1534 let mut mocks = default_test_mocks();
1535 let before_submit = std::sync::Arc::new(std::sync::atomic::AtomicI64::new(0));
1536 let before_submit_for_mock = before_submit.clone();
1537 let expected_delay = STELLAR_FAST_RESUBMIT_BASE_DELAY_SECONDS * 3;
1538
1539 let mut response = create_send_tx_response(
1540 "ERROR",
1541 "0101010101010101010101010101010101010101010101010101010101010101",
1542 );
1543 response.error_result_xdr = Some("AAAAAAAAY/n////3AAAAAA==".to_string());
1544 mocks
1545 .provider
1546 .expect_send_transaction_with_status()
1547 .returning(move |_| {
1548 let r = response.clone();
1549 Box::pin(async move { Ok(r) })
1550 });
1551
1552 mocks
1553 .tx_repo
1554 .expect_record_stellar_insufficient_fee_retry()
1555 .withf(|id, sent_at| id == "tx-1" && !sent_at.is_empty())
1556 .returning(|id, _| {
1557 let mut tx = create_test_transaction("relayer-1");
1558 tx.id = id;
1559 tx.status = TransactionStatus::Sent;
1560 tx.metadata = Some(TransactionMetadata {
1561 consecutive_failures: 0,
1562 total_failures: 0,
1563 insufficient_fee_retries: 3,
1564 try_again_later_retries: 0,
1565 nonce_too_high_retries: 0,
1566 });
1567 Ok::<_, RepositoryError>(tx)
1568 })
1569 .times(1);
1570
1571 mocks
1572 .job_producer
1573 .expect_produce_submit_transaction_job()
1574 .withf(move |_, scheduled_on| {
1575 let before_submit =
1576 before_submit_for_mock.load(std::sync::atomic::Ordering::SeqCst);
1577 let now = Utc::now().timestamp();
1578 scheduled_on.is_some_and(|scheduled_on| {
1579 scheduled_on >= before_submit + expected_delay - 1
1580 && scheduled_on <= now + expected_delay + 1
1581 })
1582 })
1583 .times(1)
1584 .returning(|_, _| Box::pin(async { Ok(()) }));
1585
1586 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
1587 let mut tx = create_test_transaction(&relayer.id);
1588 tx.status = TransactionStatus::Sent;
1589 if let NetworkTransactionData::Stellar(ref mut data) = tx.network_data {
1590 data.signatures.push(dummy_signature());
1591 data.signed_envelope_xdr = Some(create_signed_xdr(TEST_PK, TEST_PK_2));
1592 }
1593
1594 before_submit.store(Utc::now().timestamp(), std::sync::atomic::Ordering::SeqCst);
1595 let res = handler.submit_transaction_impl(tx).await;
1596
1597 let returned_tx = res.unwrap();
1598 assert_eq!(returned_tx.status, TransactionStatus::Sent);
1599 assert_eq!(
1600 returned_tx
1601 .metadata
1602 .as_ref()
1603 .map(|metadata| metadata.insufficient_fee_retries),
1604 Some(3)
1605 );
1606 }
1607
1608 #[tokio::test]
1609 async fn submit_transaction_insufficient_fee_post_record_over_cap_skips_enqueue() {
1610 let relayer = create_test_relayer();
1611 let mut mocks = default_test_mocks();
1612
1613 let mut response = create_send_tx_response(
1614 "ERROR",
1615 "0101010101010101010101010101010101010101010101010101010101010101",
1616 );
1617 response.error_result_xdr = Some("AAAAAAAAY/n////3AAAAAA==".to_string());
1618 mocks
1619 .provider
1620 .expect_send_transaction_with_status()
1621 .returning(move |_| {
1622 let r = response.clone();
1623 Box::pin(async move { Ok(r) })
1624 });
1625
1626 mocks
1627 .tx_repo
1628 .expect_record_stellar_insufficient_fee_retry()
1629 .withf(|id, sent_at| id == "tx-1" && !sent_at.is_empty())
1630 .returning(|id, _| {
1631 let mut tx = create_test_transaction("relayer-1");
1632 tx.id = id;
1633 tx.status = TransactionStatus::Sent;
1634 tx.metadata = Some(TransactionMetadata {
1635 consecutive_failures: 0,
1636 total_failures: 0,
1637 insufficient_fee_retries: STELLAR_INSUFFICIENT_FEE_MAX_RETRIES + 1,
1638 try_again_later_retries: 0,
1639 nonce_too_high_retries: 0,
1640 });
1641 Ok::<_, RepositoryError>(tx)
1642 })
1643 .times(1);
1644
1645 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
1646 let mut tx = create_test_transaction(&relayer.id);
1647 tx.status = TransactionStatus::Sent;
1648 if let NetworkTransactionData::Stellar(ref mut data) = tx.network_data {
1649 data.signatures.push(dummy_signature());
1650 data.signed_envelope_xdr = Some(create_signed_xdr(TEST_PK, TEST_PK_2));
1651 }
1652
1653 let res = handler.submit_core(tx).await;
1654
1655 let returned_tx = res.unwrap();
1656 assert_eq!(returned_tx.status, TransactionStatus::Sent);
1657 assert_eq!(
1658 returned_tx
1659 .metadata
1660 .as_ref()
1661 .map(|metadata| metadata.insufficient_fee_retries),
1662 Some(STELLAR_INSUFFICIENT_FEE_MAX_RETRIES + 1)
1663 );
1664 }
1665
1666 #[tokio::test]
1667 async fn submit_transaction_insufficient_fee_final_state_skips_enqueue() {
1668 let relayer = create_test_relayer();
1669 let mut mocks = default_test_mocks();
1670
1671 let mut response = create_send_tx_response(
1672 "ERROR",
1673 "0101010101010101010101010101010101010101010101010101010101010101",
1674 );
1675 response.error_result_xdr = Some("AAAAAAAAY/n////3AAAAAA==".to_string());
1676 mocks
1677 .provider
1678 .expect_send_transaction_with_status()
1679 .returning(move |_| {
1680 let r = response.clone();
1681 Box::pin(async move { Ok(r) })
1682 });
1683
1684 mocks
1685 .tx_repo
1686 .expect_record_stellar_insufficient_fee_retry()
1687 .withf(|id, sent_at| id == "tx-1" && !sent_at.is_empty())
1688 .returning(|id, _| {
1689 let mut tx = create_test_transaction("relayer-1");
1690 tx.id = id;
1691 tx.status = TransactionStatus::Confirmed;
1692 tx.metadata = Some(TransactionMetadata {
1693 consecutive_failures: 0,
1694 total_failures: 0,
1695 insufficient_fee_retries: 0,
1696 try_again_later_retries: 0,
1697 nonce_too_high_retries: 0,
1698 });
1699 Ok::<_, RepositoryError>(tx)
1700 })
1701 .times(1);
1702
1703 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
1704 let mut tx = create_test_transaction(&relayer.id);
1705 tx.status = TransactionStatus::Sent;
1706 tx.metadata = Some(TransactionMetadata {
1707 consecutive_failures: 0,
1708 total_failures: 0,
1709 insufficient_fee_retries: 0,
1710 try_again_later_retries: 0,
1711 nonce_too_high_retries: 0,
1712 });
1713 if let NetworkTransactionData::Stellar(ref mut data) = tx.network_data {
1714 data.signatures.push(dummy_signature());
1715 data.signed_envelope_xdr = Some(create_signed_xdr(TEST_PK, TEST_PK_2));
1716 }
1717
1718 let res = handler.submit_core(tx).await;
1719
1720 assert!(res.is_ok());
1721 let returned_tx = res.unwrap();
1722 assert_eq!(returned_tx.status, TransactionStatus::Confirmed);
1723 }
1724
1725 #[tokio::test]
1726 async fn submit_transaction_fast_resubmit_enqueue_failure_is_swallowed() {
1727 let relayer = create_test_relayer();
1728 let mut mocks = default_test_mocks();
1729
1730 let mut response = create_send_tx_response(
1731 "ERROR",
1732 "0101010101010101010101010101010101010101010101010101010101010101",
1733 );
1734 response.error_result_xdr = Some("AAAAAAAAY/n////3AAAAAA==".to_string());
1735 mocks
1736 .provider
1737 .expect_send_transaction_with_status()
1738 .returning(move |_| {
1739 let r = response.clone();
1740 Box::pin(async move { Ok(r) })
1741 });
1742
1743 mocks
1744 .tx_repo
1745 .expect_record_stellar_insufficient_fee_retry()
1746 .withf(|id, sent_at| id == "tx-1" && !sent_at.is_empty())
1747 .returning(|id, _| {
1748 let mut tx = create_test_transaction("relayer-1");
1749 tx.id = id;
1750 tx.status = TransactionStatus::Sent;
1751 tx.metadata = Some(TransactionMetadata {
1752 consecutive_failures: 0,
1753 total_failures: 0,
1754 insufficient_fee_retries: 1,
1755 try_again_later_retries: 0,
1756 nonce_too_high_retries: 0,
1757 });
1758 Ok::<_, RepositoryError>(tx)
1759 })
1760 .times(1);
1761
1762 mocks
1763 .job_producer
1764 .expect_produce_submit_transaction_job()
1765 .withf(|_, scheduled_on| scheduled_on.is_some())
1766 .times(1)
1767 .returning(|_, _| {
1768 Box::pin(async {
1769 Err(crate::jobs::JobProducerError::QueueError(
1770 "queue unavailable".to_string(),
1771 ))
1772 })
1773 });
1774
1775 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
1776 let mut tx = create_test_transaction(&relayer.id);
1777 tx.status = TransactionStatus::Sent;
1778 if let NetworkTransactionData::Stellar(ref mut data) = tx.network_data {
1779 data.signatures.push(dummy_signature());
1780 data.signed_envelope_xdr = Some(create_signed_xdr(TEST_PK, TEST_PK_2));
1781 }
1782
1783 let res = handler.submit_transaction_impl(tx).await;
1784
1785 let returned_tx = res.unwrap();
1786 assert_eq!(returned_tx.status, TransactionStatus::Sent);
1787 assert_eq!(
1788 returned_tx
1789 .metadata
1790 .as_ref()
1791 .map(|metadata| metadata.insufficient_fee_retries),
1792 Some(1)
1793 );
1794 }
1795
1796 #[tokio::test]
1797 async fn submit_transaction_try_again_later_enqueue_failure_is_swallowed() {
1798 let relayer = create_test_relayer();
1799 let mut mocks = default_test_mocks();
1800
1801 let response = create_send_tx_response(
1802 "TRY_AGAIN_LATER",
1803 "0101010101010101010101010101010101010101010101010101010101010101",
1804 );
1805 mocks
1806 .provider
1807 .expect_send_transaction_with_status()
1808 .returning(move |_| {
1809 let r = response.clone();
1810 Box::pin(async move { Ok(r) })
1811 });
1812
1813 mocks
1814 .tx_repo
1815 .expect_record_stellar_try_again_later_retry()
1816 .withf(|id, sent_at| id == "tx-1" && !sent_at.is_empty())
1817 .returning(|id, _| {
1818 let mut tx = create_test_transaction("relayer-1");
1819 tx.id = id;
1820 tx.status = TransactionStatus::Sent;
1821 tx.metadata = Some(TransactionMetadata {
1822 consecutive_failures: 0,
1823 total_failures: 0,
1824 insufficient_fee_retries: 0,
1825 try_again_later_retries: 1,
1826 nonce_too_high_retries: 0,
1827 });
1828 Ok::<_, RepositoryError>(tx)
1829 })
1830 .times(1);
1831
1832 mocks
1833 .job_producer
1834 .expect_produce_submit_transaction_job()
1835 .withf(|_, scheduled_on| scheduled_on.is_some())
1836 .times(1)
1837 .returning(|_, _| {
1838 Box::pin(async {
1839 Err(crate::jobs::JobProducerError::QueueError(
1840 "queue unavailable".to_string(),
1841 ))
1842 })
1843 });
1844
1845 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
1846 let mut tx = create_test_transaction(&relayer.id);
1847 tx.status = TransactionStatus::Sent;
1848 if let NetworkTransactionData::Stellar(ref mut data) = tx.network_data {
1849 data.signatures.push(dummy_signature());
1850 data.signed_envelope_xdr = Some(create_signed_xdr(TEST_PK, TEST_PK_2));
1851 }
1852
1853 let res = handler.submit_transaction_impl(tx).await;
1854
1855 let returned_tx = res.unwrap();
1856 assert_eq!(returned_tx.status, TransactionStatus::Sent);
1857 assert_eq!(
1858 returned_tx
1859 .metadata
1860 .as_ref()
1861 .map(|metadata| metadata.try_again_later_retries),
1862 Some(1)
1863 );
1864 }
1865
1866 #[tokio::test]
1867 async fn submit_transaction_insufficient_fee_exceeding_retry_limit_fails() {
1868 let relayer = create_test_relayer();
1869 let mut mocks = default_test_mocks();
1870 let retry_limit_reason = format!(
1871 "insufficient fee retry limit exceeded ({STELLAR_INSUFFICIENT_FEE_MAX_RETRIES})"
1872 );
1873 let retry_limit_reason_for_mock = retry_limit_reason.clone();
1874
1875 let mut response = create_send_tx_response(
1876 "ERROR",
1877 "0101010101010101010101010101010101010101010101010101010101010101",
1878 );
1879 response.error_result_xdr = Some("AAAAAAAAY/n////3AAAAAA==".to_string());
1880 mocks
1881 .provider
1882 .expect_send_transaction_with_status()
1883 .returning(move |_| {
1884 let r = response.clone();
1885 Box::pin(async move { Ok(r) })
1886 });
1887
1888 mocks
1889 .tx_repo
1890 .expect_partial_update()
1891 .withf(move |_, upd| {
1892 upd.status == Some(TransactionStatus::Failed)
1893 && upd
1894 .status_reason
1895 .as_ref()
1896 .is_some_and(|reason| reason.contains(&retry_limit_reason_for_mock))
1897 })
1898 .returning(|id, upd| {
1899 let mut tx = create_test_transaction("relayer-1");
1900 tx.id = id;
1901 tx.status = upd.status.unwrap();
1902 tx.status_reason = upd.status_reason;
1903 Ok::<_, RepositoryError>(tx)
1904 });
1905
1906 mocks
1907 .job_producer
1908 .expect_produce_send_notification_job()
1909 .times(1)
1910 .returning(|_, _| Box::pin(async { Ok(()) }));
1911
1912 mocks
1913 .tx_repo
1914 .expect_find_by_status_paginated()
1915 .returning(move |_, _, _, _| {
1916 Ok(PaginatedResult {
1917 items: vec![],
1918 total: 0,
1919 page: 1,
1920 per_page: 1,
1921 })
1922 });
1923
1924 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
1925 let mut tx = create_test_transaction(&relayer.id);
1926 tx.status = TransactionStatus::Sent;
1927 tx.metadata = Some(TransactionMetadata {
1928 insufficient_fee_retries: STELLAR_INSUFFICIENT_FEE_MAX_RETRIES,
1929 ..Default::default()
1930 });
1931 if let NetworkTransactionData::Stellar(ref mut data) = tx.network_data {
1932 data.signatures.push(dummy_signature());
1933 data.signed_envelope_xdr = Some(create_signed_xdr(TEST_PK, TEST_PK_2));
1934 }
1935
1936 let res = handler.submit_transaction_impl(tx).await;
1937
1938 let failed_tx = res.unwrap();
1939 assert_eq!(failed_tx.status, TransactionStatus::Failed);
1940 assert!(failed_tx
1941 .status_reason
1942 .as_ref()
1943 .is_some_and(|reason| reason.contains(&retry_limit_reason)));
1944 }
1945
1946 #[tokio::test]
1947 async fn submit_transaction_error_non_fee_still_fails() {
1948 let relayer = create_test_relayer();
1949 let mut mocks = default_test_mocks();
1950
1951 let mut response = create_send_tx_response(
1953 "ERROR",
1954 "0101010101010101010101010101010101010101010101010101010101010101",
1955 );
1956 response.error_result_xdr = Some("AAAAAAAAA/v////6AAAAAA==".to_string());
1957 mocks
1958 .provider
1959 .expect_send_transaction_with_status()
1960 .returning(move |_| {
1961 let r = response.clone();
1962 Box::pin(async move { Ok(r) })
1963 });
1964
1965 mocks
1967 .tx_repo
1968 .expect_partial_update()
1969 .withf(|_, upd| upd.status == Some(TransactionStatus::Failed))
1970 .returning(|id, upd| {
1971 let mut tx = create_test_transaction("relayer-1");
1972 tx.id = id;
1973 tx.status = upd.status.unwrap();
1974 Ok::<_, RepositoryError>(tx)
1975 });
1976
1977 mocks
1979 .job_producer
1980 .expect_produce_send_notification_job()
1981 .times(1)
1982 .returning(|_, _| Box::pin(async { Ok(()) }));
1983
1984 mocks
1986 .tx_repo
1987 .expect_find_by_status_paginated()
1988 .returning(move |_, _, _, _| {
1989 Ok(PaginatedResult {
1990 items: vec![],
1991 total: 0,
1992 page: 1,
1993 per_page: 1,
1994 })
1995 });
1996
1997 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
1998 let mut tx = create_test_transaction(&relayer.id);
1999 tx.status = TransactionStatus::Sent;
2000 if let NetworkTransactionData::Stellar(ref mut data) = tx.network_data {
2001 data.signatures.push(dummy_signature());
2002 data.signed_envelope_xdr = Some(create_signed_xdr(TEST_PK, TEST_PK_2));
2003 }
2004
2005 let res = handler.submit_transaction_impl(tx).await;
2006
2007 let failed_tx = res.unwrap();
2009 assert_eq!(failed_tx.status, TransactionStatus::Failed);
2010 }
2011
2012 #[tokio::test]
2013 async fn submit_transaction_concurrent_update_conflict_reloads_latest_state() {
2014 let relayer = create_test_relayer();
2017 let mut mocks = default_test_mocks();
2018
2019 let response = create_send_tx_response(
2021 "PENDING",
2022 "0101010101010101010101010101010101010101010101010101010101010101",
2023 );
2024 mocks
2025 .provider
2026 .expect_send_transaction_with_status()
2027 .returning(move |_| {
2028 let r = response.clone();
2029 Box::pin(async move { Ok(r) })
2030 });
2031
2032 mocks
2034 .tx_repo
2035 .expect_partial_update()
2036 .withf(|_, upd| upd.status == Some(TransactionStatus::Submitted))
2037 .times(1)
2038 .returning(|_, _| {
2039 Err(RepositoryError::ConcurrentUpdateConflict(
2040 "CAS mismatch".to_string(),
2041 ))
2042 });
2043
2044 let reloaded_tx = {
2046 let mut t = create_test_transaction(&relayer.id);
2047 t.status = TransactionStatus::Submitted;
2048 t
2049 };
2050 let reloaded_clone = reloaded_tx.clone();
2051 mocks
2052 .tx_repo
2053 .expect_get_by_id()
2054 .times(1)
2055 .returning(move |_| Ok(reloaded_clone.clone()));
2056
2057 mocks
2059 .job_producer
2060 .expect_produce_send_notification_job()
2061 .never();
2062 mocks
2063 .job_producer
2064 .expect_produce_transaction_request_job()
2065 .never();
2066
2067 let handler = make_stellar_tx_handler(relayer.clone(), mocks);
2068 let mut tx = create_test_transaction(&relayer.id);
2069 tx.status = TransactionStatus::Sent;
2070 if let NetworkTransactionData::Stellar(ref mut data) = tx.network_data {
2071 data.signatures.push(dummy_signature());
2072 data.signed_envelope_xdr = Some(create_signed_xdr(TEST_PK, TEST_PK_2));
2073 }
2074
2075 let res = handler.submit_transaction_impl(tx).await;
2076
2077 assert!(res.is_ok(), "CAS conflict should return Ok after reload");
2078 let returned_tx = res.unwrap();
2079 assert_eq!(returned_tx.status, TransactionStatus::Submitted);
2081 }
2082 }
2083}