openzeppelin_relayer/domain/transaction/stellar/
submit.rs

1//! This module contains the submission-related functionality for Stellar transactions.
2//! It includes methods for submitting transactions with robust error handling,
3//! ensuring proper transaction state management on failure.
4
5use 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    /// Main submission method with robust error handling.
43    /// Unlike prepare, submit doesn't claim lanes but still needs proper error handling.
44    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        // Defensive check: if transaction is in a final state or unexpected state, don't retry
56        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        // Check if transaction has expired before attempting submission
67        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        // Call core submission logic with error handling
80        match self.submit_core(tx.clone()).await {
81            Ok(submitted_tx) => Ok(submitted_tx),
82            Err(error) => {
83                // Handle submission failure - mark as failed and send notification
84                self.handle_submit_failure(tx, error).await
85            }
86        }
87    }
88
89    /// Core submission logic - pure business logic without error handling concerns.
90    ///
91    /// Uses `send_transaction_with_status` to get full status information from the RPC.
92    /// Handles status codes:
93    /// - PENDING: Transaction accepted for processing
94    /// - DUPLICATE: Transaction already submitted (treat as success)
95    /// - TRY_AGAIN_LATER: Network congested but tx is valid; update sent_at and schedule a
96    ///   fast delayed resubmit (5s × attempt) for the first 3 rejections
97    /// - ERROR: Transaction validation failed, except insufficient fee schedules the same fast
98    ///   delayed resubmit through the insufficient-fee retry cap (6)
99    ///
100    /// The status checker remains the fallback rescue. Each fast attempt refreshes sent_at before
101    /// enqueueing, so the ladder's ≥10s time-since-last-submit gate stays closed while the
102    /// 5s-spaced fast attempts are active.
103    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        // Use send_transaction_with_status to get full status information
113        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        // Handle status codes from the RPC response
125        match response.status.as_str() {
126            "PENDING" | "DUPLICATE" => {
127                // Success - transaction is accepted or already exists
128                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                // Send notification for newly submitted transaction
158                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                // Network is temporarily congested — the transaction is valid but the
171                // node's queue is full. Atomically update sent_at and increment
172                // try_again_later_retries so the status checker's backoff gate measures
173                // time since this attempt. Return Ok to keep the transaction alive;
174                // the fast path handles the bounded early retries, and the status
175                // checker remains the fallback rescue.
176                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                // The repo returns the transaction unchanged if it reached a final state
182                // between this submission and the record call (e.g. a racing job finalized
183                // it). Nothing left to rescue — skip metrics and scheduling.
184                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                // Only push on first encounter (dedup: won't fire on retry 2, 3, etc.)
200                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                // Transaction validation failed
248                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                // Insufficient fee is a transient condition (network fee spike).
254                // Treat like TRY_AGAIN_LATER: update sent_at and let the status
255                // checker retry with exponential backoff.
256                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                    // Only push on first encounter (dedup: won't fire on retry 2, 3, etc.)
264                    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                    // Atomically sets `sent_at` and increments Stellar insufficient-fee retries.
280                    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                    // The repo returns the transaction unchanged if it reached a final state
289                    // between this submission and the record call (e.g. a racing job finalized
290                    // it). Nothing left to rescue — skip metrics and scheduling.
291                    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                    // This can only happen when a concurrent submission increments the counter
307                    // between our read and record. That winner already scheduled the retry; skip
308                    // this enqueue so the next rejection terminates normally at the pre-record cap.
309                    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                // Unknown status - treat as error
363                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    /// Handles submission failures with comprehensive cleanup and error reporting.
380    /// For bad sequence errors, resets the transaction and re-enqueues it for retry.
381    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        // CAS conflict in the submission path only occurs after the RPC
397        // already accepted the transaction (PENDING status update raced).
398        // The on-chain state is valid; reload the latest DB state and return
399        // Ok — the status checker will reconcile on its next poll.
400        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            // For bad sequence errors, sync sequence from chain first
415            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            // Reset the transaction to pending state
444            // Status check will handle resubmission when it detects a pending transaction without hash
445            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 success since we've reset the transaction
458                    // Status check job (scheduled with delay) will detect pending without hash
459                    // and schedule a recovery job to go through the pipeline again
460                    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                    // Fall through to normal failure handling
470                }
471            }
472        }
473
474        // For non-bad-sequence errors or if reset failed, mark as failed
475        // Step 1: Mark transaction as Failed with detailed reason
476        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                // Finalization failed — propagate error so the queue retries
494                // and the next attempt will either finalize or hit is_final_state
495                return Err(error);
496            }
497        };
498
499        // Attempt to enqueue next pending transaction or release lane
500        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        // Transaction successfully marked as failed — return Ok to avoid
517        // a pointless queue retry (the defensive is_final_state check at the
518        // top of submit_transaction_impl would short-circuit anyway).
519        Ok(failed_tx)
520    }
521
522    /// Resubmit transaction - delegates to submit_transaction_impl
523    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    /// Helper to create a SendTransactionResponse with given status
541    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            // provider returns PENDING status
565            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            // expect partial update to Submitted
578            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            // Expect notification
590            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; // Must be Sent for idempotent submit
600            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                // Valid XDR
604            }
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            // Provider fails with non-bad-sequence error
616            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            // Mock finalize_transaction_state for failure handling
624            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            // Mock notification for failed transaction
636            mocks
637                .job_producer
638                .expect_produce_send_notification_job()
639                .times(1)
640                .returning(|_, _| Box::pin(async { Ok(()) }));
641
642            // Mock find_by_status_paginated for enqueue_next_pending_transaction
643            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                }); // No pending transactions
654
655            let handler = make_stellar_tx_handler(relayer.clone(), mocks);
656            let mut tx = create_test_transaction(&relayer.id);
657            tx.status = TransactionStatus::Sent; // Must be Sent for idempotent submit
658            if let NetworkTransactionData::Stellar(ref mut data) = tx.network_data {
659                data.signatures.push(dummy_signature());
660                data.sequence_number = Some(42); // Set sequence number
661                data.signed_envelope_xdr = Some("test-xdr".to_string()); // Required for submission
662            }
663
664            let res = handler.submit_transaction_impl(tx).await;
665
666            // Transaction is marked as failed and returned as Ok (no queue retry needed)
667            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            // Provider returns PENDING status
677            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            // Repository fails on first update (submission)
690            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            // Mock finalize_transaction_state for failure handling
697            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            // Mock notification for failed transaction
709            mocks
710                .job_producer
711                .expect_produce_send_notification_job()
712                .times(1)
713                .returning(|_, _| Box::pin(async { Ok(()) }));
714
715            // Mock find_by_status_paginated for enqueue_next_pending_transaction
716            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                }); // No pending transactions
727
728            let handler = make_stellar_tx_handler(relayer.clone(), mocks);
729            let mut tx = create_test_transaction(&relayer.id);
730            tx.status = TransactionStatus::Sent; // Must be Sent for idempotent submit
731            if let NetworkTransactionData::Stellar(ref mut data) = tx.network_data {
732                data.signatures.push(dummy_signature());
733                data.sequence_number = Some(42); // Set sequence number
734                data.signed_envelope_xdr = Some("test-xdr".to_string()); // Required for submission
735            }
736
737            let res = handler.submit_transaction_impl(tx).await;
738
739            // Even though provider succeeded and repo failed on Submitted update,
740            // the failure handler marks the tx as Failed and returns Ok
741            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            // Create a transaction with signed_envelope_xdr set
751            let mut tx = create_test_transaction(&relayer.id);
752            tx.status = TransactionStatus::Sent; // Must be Sent for idempotent submit
753            if let NetworkTransactionData::Stellar(ref mut data) = tx.network_data {
754                data.signatures.push(dummy_signature());
755                // Build and store the signed envelope XDR
756                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            // Provider should receive the envelope decoded from signed_envelope_xdr
764            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            // Update to Submitted
777            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            // Expect notification
789            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            // provider returns PENDING status
807            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            // expect partial update to Submitted
820            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            // Expect notification
832            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; // Must be Sent for idempotent submit
842            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                // Valid XDR
846            }
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            // Provider fails with non-bad-sequence error
858            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            // No sync expected for non-bad-sequence errors
866
867            // Mock finalize_transaction_state for failure handling
868            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            // Mock notification for failed transaction
880            mocks
881                .job_producer
882                .expect_produce_send_notification_job()
883                .times(1)
884                .returning(|_, _| Box::pin(async { Ok(()) }));
885
886            // Mock find_by_status to return a pending transaction
887            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            // Mock produce_transaction_request_job for the next pending transaction
913            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; // Must be Sent for idempotent submit
923            if let NetworkTransactionData::Stellar(ref mut data) = tx.network_data {
924                data.signatures.push(dummy_signature());
925                data.sequence_number = Some(42); // Set sequence number
926                data.signed_envelope_xdr = Some("test-xdr".to_string()); // Required for submission
927            }
928
929            let res = handler.submit_transaction_impl(tx).await;
930
931            // Transaction marked as failed and next transaction enqueued
932            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            // Mock provider to return bad sequence error
942            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            // Mock get_account for sync_sequence_from_chain
954            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            // Mock counter sync_floor for sync_sequence_from_chain
981            mocks
982                .counter
983                .expect_sync_floor()
984                .times(1)
985                .returning(|_, _, floor| Box::pin(async move { Ok(floor) }));
986
987            // Mock partial_update for reset_transaction_for_retry - should reset to Pending
988            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            // Note: Status check will handle resubmission when it detects a pending transaction without hash
1004            // We don't schedule the job here - it will be scheduled by status check when the transaction is old enough
1005
1006            let handler = make_stellar_tx_handler(relayer.clone(), mocks);
1007            let mut tx = create_test_transaction(&relayer.id);
1008            tx.status = TransactionStatus::Sent; // Must be Sent for idempotent submit
1009            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                // Valid XDR
1014            }
1015
1016            let result = handler.submit_transaction_impl(tx).await;
1017
1018            // Should return Ok since we're handling the retry
1019            assert!(result.is_ok());
1020            let reset_tx = result.unwrap();
1021            assert_eq!(reset_tx.status, TransactionStatus::Pending);
1022
1023            // Verify stellar data was reset
1024            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            // Provider returns DUPLICATE status
1040            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            // expect partial update to Submitted
1053            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            // Provider returns TRY_AGAIN_LATER status
1083            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            // Transaction stays in Sent and gets a delayed fast resubmit.
1131            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            // submission returns TRY_AGAIN_LATER, transaction remains Sent.
1140            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            // status check sees stale Sent tx and re-enqueues submit job.
1196            // Both created_at and sent_at must exceed the base resubmit interval
1197            // for the backoff logic to trigger. created_at is set earlier than sent_at
1198            // to match real-world invariants (transaction is created before being sent).
1199            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            // Provider returns TRY_AGAIN_LATER status
1231            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; // Already submitted (resubmission path)
1271            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            // Should succeed without marking as failed and get a delayed fast resubmit.
1279            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            // Provider returns ERROR status with error XDR
1403            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            // Mock finalize_transaction_state for failure handling
1417            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            // Mock notification for failed transaction
1429            mocks
1430                .job_producer
1431                .expect_produce_send_notification_job()
1432                .times(1)
1433                .returning(|_, _| Box::pin(async { Ok(()) }));
1434
1435            // Mock find_by_status_paginated for enqueue_next_pending_transaction
1436            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            // Transaction marked as failed — no error propagated
1459            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            // Provider returns ERROR status with insufficient fee XDR
1469            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            // insufficient-fee retry updates sent_at and retry metadata atomically
1483            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            // Transaction stays alive — status checker will retry
1520            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            // Provider returns ERROR status with a non-fee error
1952            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            // Mock finalize_transaction_state for failure handling
1966            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            // Mock notification for failed transaction
1978            mocks
1979                .job_producer
1980                .expect_produce_send_notification_job()
1981                .times(1)
1982                .returning(|_, _| Box::pin(async { Ok(()) }));
1983
1984            // Mock find_by_status_paginated for enqueue_next_pending_transaction
1985            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            // Non-fee ERROR still marks as failed
2008            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            // When partial_update fails with ConcurrentUpdateConflict during submission,
2015            // the handler should reload the latest state via get_by_id and return Ok.
2016            let relayer = create_test_relayer();
2017            let mut mocks = default_test_mocks();
2018
2019            // Provider returns PENDING — submission to RPC succeeded
2020            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            // partial_update (Submitted) fails with CAS conflict
2033            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            // After conflict, handler reloads via get_by_id
2045            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            // No failure handling (notifications, next-pending) should occur
2058            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            // Reloaded state reflects the concurrent writer's update
2080            assert_eq!(returned_tx.status, TransactionStatus::Submitted);
2081        }
2082    }
2083}