openzeppelin_relayer/domain/transaction/stellar/
stellar_transaction.rs

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