openzeppelin_relayer/domain/relayer/stellar/
stellar_relayer.rs

1use crate::constants::get_stellar_sponsored_transaction_validity_duration;
2use crate::domain::relayer::evm::create_error_response;
3use crate::services::stellar_dex::StellarDexService;
4use crate::utils::{map_provider_error, sanitize_error_description};
5/// This module defines the `StellarRelayer` struct and its associated functionality for
6/// interacting with Stellar networks. The `StellarRelayer` is responsible for managing
7/// transactions, synchronizing sequence numbers, and ensuring the relayer's state is
8/// consistent with the Stellar blockchain.
9///
10/// # Components
11///
12/// - `StellarRelayer`: The main struct that encapsulates the relayer's state and operations for Stellar.
13/// - `RelayerRepoModel`: Represents the relayer's data model.
14/// - `StellarProvider`: Provides blockchain interaction capabilities, such as fetching account details.
15/// - `TransactionCounterService`: Manages the sequence number for transactions to ensure correct ordering.
16/// - `JobProducer`: Produces jobs for processing transactions and sending notifications.
17///
18/// # Error Handling
19///
20/// The module uses the `RelayerError` enum to handle various errors that can occur during
21/// operations, such as provider errors, sequence synchronization failures, and transaction failures.
22///
23/// # Usage
24///
25/// To use the `StellarRelayer`, create an instance using the `new` method, providing the necessary
26/// components. Then, call the appropriate methods to process transactions and manage the relayer's state.
27use crate::{
28    constants::{
29        transactions::PENDING_TRANSACTION_STATUSES, STELLAR_SMALLEST_UNIT_NAME,
30        STELLAR_STATUS_CHECK_INITIAL_DELAY_SECONDS,
31    },
32    domain::{
33        create_success_response, transaction::stellar::fetch_next_sequence_from_chain,
34        BalanceResponse, SignDataRequest, SignDataResponse, SignTransactionExternalResponse,
35        SignTransactionExternalResponseStellar, SignTransactionRequest, SignTypedDataRequest,
36    },
37    jobs::{JobProducerTrait, RelayerHealthCheck, TransactionRequest, TransactionStatusCheck},
38    models::{
39        produce_relayer_disabled_payload, DeletePendingTransactionsResponse, DisabledReason,
40        HealthCheckFailure, JsonRpcRequest, JsonRpcResponse, NetworkRepoModel, NetworkRpcRequest,
41        NetworkRpcResult, NetworkTransactionRequest, NetworkType, PaginationQuery,
42        RelayerNetworkPolicy, RelayerRepoModel, RelayerStatus, RelayerStellarPolicy,
43        RepositoryError, RpcErrorCodes, StellarAllowedTokensPolicy, StellarFeePaymentStrategy,
44        StellarNetwork, StellarRpcRequest, TransactionRepoModel, TransactionStatus,
45        TransactionUpdateRequest,
46    },
47    repositories::{NetworkRepository, RelayerRepository, Repository, TransactionRepository},
48    services::{
49        provider::{StellarProvider, StellarProviderTrait},
50        signer::{StellarSignTrait, StellarSigner},
51        stellar_dex::StellarDexServiceTrait,
52        TransactionCounterService, TransactionCounterServiceTrait,
53    },
54    utils::calculate_scheduled_timestamp,
55};
56use async_trait::async_trait;
57use eyre::Result;
58use futures::future::try_join_all;
59use std::sync::Arc;
60use tracing::{debug, error, info, instrument, warn};
61
62use crate::domain::relayer::stellar::xdr_utils::parse_transaction_xdr;
63use crate::domain::relayer::{Relayer, RelayerError, StellarRelayerDexTrait};
64use crate::domain::transaction::stellar::token::get_token_metadata;
65use crate::domain::transaction::stellar::StellarTransactionValidator;
66
67/// Dependencies container for `StellarRelayer` construction.
68pub struct StellarRelayerDependencies<RR, NR, TR, J, TCS>
69where
70    RR: Repository<RelayerRepoModel, String> + RelayerRepository + Send + Sync + 'static,
71    NR: NetworkRepository + Repository<NetworkRepoModel, String> + Send + Sync + 'static,
72    TR: Repository<TransactionRepoModel, String> + TransactionRepository + Send + Sync + 'static,
73    J: JobProducerTrait + Send + Sync + 'static,
74    TCS: TransactionCounterServiceTrait + Send + Sync + 'static,
75{
76    pub relayer_repository: Arc<RR>,
77    pub network_repository: Arc<NR>,
78    pub transaction_repository: Arc<TR>,
79    pub transaction_counter_service: Arc<TCS>,
80    pub job_producer: Arc<J>,
81}
82
83impl<RR, NR, TR, J, TCS> StellarRelayerDependencies<RR, NR, TR, J, TCS>
84where
85    RR: Repository<RelayerRepoModel, String> + RelayerRepository + Send + Sync + 'static,
86    NR: NetworkRepository + Repository<NetworkRepoModel, String> + Send + Sync + 'static,
87    TR: Repository<TransactionRepoModel, String> + TransactionRepository + Send + Sync + 'static,
88    J: JobProducerTrait + Send + Sync,
89    TCS: TransactionCounterServiceTrait + Send + Sync + 'static,
90{
91    /// Creates a new dependencies container for `StellarRelayer`.
92    ///
93    /// # Arguments
94    ///
95    /// * `relayer_repository` - Repository for managing relayer model persistence
96    /// * `network_repository` - Repository for accessing network configuration data (RPC URLs, chain settings)
97    /// * `transaction_repository` - Repository for storing and retrieving transaction models
98    /// * `transaction_counter_service` - Service for managing sequence numbers to ensure proper transaction ordering
99    /// * `job_producer` - Service for creating background jobs for transaction processing and notifications
100    ///
101    /// # Returns
102    ///
103    /// Returns a new `StellarRelayerDependencies` instance containing all provided dependencies.
104    pub fn new(
105        relayer_repository: Arc<RR>,
106        network_repository: Arc<NR>,
107        transaction_repository: Arc<TR>,
108        transaction_counter_service: Arc<TCS>,
109        job_producer: Arc<J>,
110    ) -> Self {
111        Self {
112            relayer_repository,
113            network_repository,
114            transaction_repository,
115            transaction_counter_service,
116            job_producer,
117        }
118    }
119}
120
121#[allow(dead_code)]
122pub struct StellarRelayer<P, RR, NR, TR, J, TCS, S, D>
123where
124    P: StellarProviderTrait + Send + Sync + 'static,
125    RR: Repository<RelayerRepoModel, String> + RelayerRepository + Send + Sync + 'static,
126    NR: NetworkRepository + Repository<NetworkRepoModel, String> + Send + Sync + 'static,
127    TR: Repository<TransactionRepoModel, String> + TransactionRepository + Send + Sync + 'static,
128    J: JobProducerTrait + Send + Sync + 'static,
129    TCS: TransactionCounterServiceTrait + Send + Sync + 'static,
130    S: StellarSignTrait + Send + Sync + 'static,
131    D: StellarDexServiceTrait + Send + Sync + 'static,
132{
133    pub(crate) relayer: RelayerRepoModel,
134    pub(crate) signer: Arc<S>,
135    pub(crate) network: StellarNetwork,
136    pub(crate) provider: P,
137    pub(crate) relayer_repository: Arc<RR>,
138    network_repository: Arc<NR>,
139    transaction_repository: Arc<TR>,
140    transaction_counter_service: Arc<TCS>,
141    pub(crate) job_producer: Arc<J>,
142    pub(crate) dex_service: Arc<D>,
143}
144
145pub type DefaultStellarRelayer<J, TR, NR, RR, TCR> = StellarRelayer<
146    StellarProvider,
147    RR,
148    NR,
149    TR,
150    J,
151    TransactionCounterService<TCR>,
152    StellarSigner,
153    StellarDexService<StellarProvider, StellarSigner>,
154>;
155
156impl<P, RR, NR, TR, J, TCS, S, D> StellarRelayer<P, RR, NR, TR, J, TCS, S, D>
157where
158    P: StellarProviderTrait + Send + Sync,
159    RR: Repository<RelayerRepoModel, String> + RelayerRepository + Send + Sync + 'static,
160    NR: NetworkRepository + Repository<NetworkRepoModel, String> + Send + Sync + 'static,
161    TR: Repository<TransactionRepoModel, String> + TransactionRepository + Send + Sync + 'static,
162    J: JobProducerTrait + Send + Sync + 'static,
163    TCS: TransactionCounterServiceTrait + Send + Sync + 'static,
164    S: StellarSignTrait + Send + Sync + 'static,
165    D: StellarDexServiceTrait + Send + Sync + 'static,
166{
167    /// Creates a new `StellarRelayer` instance.
168    ///
169    /// This constructor initializes a new Stellar relayer with the provided configuration,
170    /// provider, and dependencies. It validates the network configuration and sets up
171    /// all necessary components for transaction processing.
172    ///
173    /// # Arguments
174    ///
175    /// * `relayer` - The relayer model containing configuration like ID, address, network name, and policies
176    /// * `signer` - The Stellar signer for signing transactions
177    /// * `provider` - The Stellar provider implementation for blockchain interactions (account queries, transaction submission)
178    /// * `dependencies` - Container with all required repositories and services (see [`StellarRelayerDependencies`])
179    /// * `dex_service` - The DEX service implementation for swap operations
180    ///
181    /// # Returns
182    ///
183    /// * `Ok(StellarRelayer)` - Successfully initialized relayer ready for operation
184    /// * `Err(RelayerError)` - If initialization fails due to configuration or validation errors
185    #[allow(clippy::too_many_arguments)]
186    pub async fn new(
187        relayer: RelayerRepoModel,
188        signer: Arc<S>,
189        provider: P,
190        dependencies: StellarRelayerDependencies<RR, NR, TR, J, TCS>,
191        dex_service: Arc<D>,
192    ) -> Result<Self, RelayerError> {
193        let network_repo = dependencies
194            .network_repository
195            .get_by_name(NetworkType::Stellar, &relayer.network)
196            .await
197            .ok()
198            .flatten()
199            .ok_or_else(|| {
200                RelayerError::NetworkConfiguration(format!("Network {} not found", relayer.network))
201            })?;
202
203        let network = StellarNetwork::try_from(network_repo.clone())?;
204
205        Ok(Self {
206            relayer,
207            signer,
208            network,
209            provider,
210            relayer_repository: dependencies.relayer_repository,
211            network_repository: dependencies.network_repository,
212            transaction_repository: dependencies.transaction_repository,
213            transaction_counter_service: dependencies.transaction_counter_service,
214            job_producer: dependencies.job_producer,
215            dex_service,
216        })
217    }
218
219    #[instrument(
220        level = "debug",
221        skip(self),
222        fields(
223            request_id = ?crate::observability::request_id::get_request_id(),
224            relayer_id = %self.relayer.id,
225        )
226    )]
227    async fn sync_sequence(&self) -> Result<(), RelayerError> {
228        info!(
229            address = %self.relayer.address,
230            "syncing sequence from chain"
231        );
232
233        let next = fetch_next_sequence_from_chain(&self.provider, &self.relayer.address)
234            .await
235            .map_err(RelayerError::ProviderError)?;
236
237        // Raise the local counter to the on-chain floor *monotonically*: a blind `set` here
238        // would rewind the counter when concurrent allocations have already advanced it past
239        // the chain value, causing duplicate/gapped sequences and a TxBadSeq cascade
240        // (the live chain sequence is a floor, never the authoritative next assignment).
241        // `sync_floor` only advances when the chain is ahead.
242        let effective = self
243            .transaction_counter_service
244            .sync_floor(next)
245            .await
246            .map_err(RelayerError::from)?;
247
248        info!(
249            chain_sequence = %next,
250            effective_sequence = %effective,
251            "synced local sequence counter to chain floor"
252        );
253        Ok(())
254    }
255
256    /// Populates the allowed tokens metadata for the Stellar relayer policy.
257    ///
258    /// This method checks whether allowed tokens have been configured in the relayer's policy.
259    /// If allowed tokens are provided, it concurrently fetches token metadata for each token,
260    /// determines the token kind (Native, Classic, or Contract), and populates metadata including
261    /// decimals and canonical asset ID. The updated policy is then stored in the repository.
262    ///
263    /// If no allowed tokens are specified, it logs an informational message and returns the policy
264    /// unchanged.
265    #[instrument(
266        level = "debug",
267        skip(self),
268        fields(
269            request_id = ?crate::observability::request_id::get_request_id(),
270            relayer_id = %self.relayer.id,
271        )
272    )]
273    async fn populate_allowed_tokens_metadata(&self) -> Result<RelayerStellarPolicy, RelayerError> {
274        let mut policy = self.relayer.policies.get_stellar_policy();
275        // Check if allowed_tokens is specified; if not, return the policy unchanged.
276        let allowed_tokens = match policy.allowed_tokens.as_ref() {
277            Some(tokens) if !tokens.is_empty() => tokens,
278            _ => {
279                info!("No allowed tokens specified; skipping token metadata population.");
280                return Ok(policy);
281            }
282        };
283
284        let token_metadata_futures = allowed_tokens.iter().map(|token| {
285            let asset_id = token.asset.clone();
286            let provider = &self.provider;
287            async move {
288                let metadata = get_token_metadata(provider, &asset_id)
289                    .await
290                    .map_err(RelayerError::from)?;
291
292                Ok::<StellarAllowedTokensPolicy, RelayerError>(StellarAllowedTokensPolicy {
293                    asset: asset_id,
294                    metadata: Some(metadata),
295                    max_allowed_fee: token.max_allowed_fee,
296                    swap_config: token.swap_config.clone(),
297                })
298            }
299        });
300
301        let updated_allowed_tokens = try_join_all(token_metadata_futures).await?;
302
303        policy.allowed_tokens = Some(updated_allowed_tokens.clone());
304
305        self.relayer_repository
306            .update_policy(
307                self.relayer.id.clone(),
308                RelayerNetworkPolicy::Stellar(policy.clone()),
309            )
310            .await?;
311
312        Ok(policy)
313    }
314
315    /// Migrates fee_payment_strategy policy for older relayers that don't have it set.
316    ///
317    /// This migration is needed for relayers that were created before `fee_payment_strategy`
318    /// became a required policy. For relayers persisted in Redis storage, this ensures
319    /// backward compatibility by setting the policy to `Relayer` (the old default behavior).
320    ///
321    /// In-memory relayers don't need this migration as they are recreated from config.json
322    /// on startup, which would have the policy set if using a newer version.
323    #[instrument(
324        level = "debug",
325        skip(self),
326        fields(
327            request_id = ?crate::observability::request_id::get_request_id(),
328            relayer_id = %self.relayer.id,
329        )
330    )]
331    async fn migrate_fee_payment_strategy_if_needed(&self) -> Result<(), RelayerError> {
332        // Only migrate if using persistent storage (Redis)
333        // In-memory relayers are recreated from config.json on startup
334        if !self.relayer_repository.is_persistent_storage() {
335            debug!(
336                relayer_id = %self.relayer.id,
337                "Skipping migration: using in-memory storage"
338            );
339            return Ok(());
340        }
341
342        let policy = self.relayer.policies.get_stellar_policy();
343
344        // If fee_payment_strategy is already set, no migration needed
345        if policy.fee_payment_strategy.is_some() {
346            return Ok(());
347        }
348
349        // Migration needed: fee_payment_strategy is missing
350        info!(
351            relayer_id = %self.relayer.id,
352            "Migrating Stellar relayer: setting fee_payment_strategy to 'Relayer' (old default behavior)"
353        );
354
355        // Create updated policy with fee_payment_strategy set to Relayer
356        let mut updated_policy = policy;
357        updated_policy.fee_payment_strategy = Some(StellarFeePaymentStrategy::Relayer);
358
359        // Update the relayer in the repository
360        self.relayer_repository
361            .update_policy(
362                self.relayer.id.clone(),
363                RelayerNetworkPolicy::Stellar(updated_policy),
364            )
365            .await
366            .map_err(|e| {
367                RelayerError::PolicyConfigurationError(format!(
368                    "Failed to migrate fee_payment_strategy policy: {e}"
369                ))
370            })?;
371
372        debug!(
373            relayer_id = %self.relayer.id,
374            "Successfully migrated fee_payment_strategy policy"
375        );
376
377        Ok(())
378    }
379
380    /// Checks the relayer's XLM balance and triggers token swap if it falls below the
381    /// specified threshold. Only proceeds with swap if balance is below the configured
382    /// min_balance_threshold.
383    #[instrument(
384        level = "debug",
385        skip(self),
386        fields(
387            request_id = ?crate::observability::request_id::get_request_id(),
388            relayer_id = %self.relayer.id,
389        )
390    )]
391    async fn check_balance_and_trigger_token_swap_if_needed(&self) -> Result<(), RelayerError> {
392        let policy = self.relayer.policies.get_stellar_policy();
393
394        // Check if swap config exists
395        let swap_config = match policy.get_swap_config() {
396            Some(config) => config,
397            None => {
398                debug!(
399                    relayer_id = %self.relayer.id,
400                    "No swap configuration specified; skipping balance check"
401                );
402                return Ok(());
403            }
404        };
405
406        // Early return if no threshold is configured (mirrors Solana logic)
407        let threshold = match swap_config.min_balance_threshold {
408            Some(threshold) => threshold,
409            None => {
410                debug!(
411                    relayer_id = %self.relayer.id,
412                    "No swap min balance threshold specified; skipping validation"
413                );
414                return Ok(());
415            }
416        };
417
418        // Get balance only when threshold is configured
419        let balance_response = self.get_balance().await?;
420        let current_balance = u64::try_from(balance_response.balance).map_err(|_| {
421            RelayerError::Internal("Account balance exceeds u64 maximum value".to_string())
422        })?;
423
424        // Only trigger swap if balance is below threshold
425        if current_balance < threshold {
426            debug!(
427                relayer_id = %self.relayer.id,
428                balance = current_balance,
429                threshold = threshold,
430                "XLM balance is below threshold, triggering token swap"
431            );
432
433            let _swap_results = self
434                .handle_token_swap_request(self.relayer.id.clone())
435                .await?;
436        } else {
437            debug!(
438                relayer_id = %self.relayer.id,
439                balance = current_balance,
440                threshold = threshold,
441                "XLM balance is above threshold, no swap needed"
442            );
443        }
444
445        Ok(())
446    }
447}
448
449#[async_trait]
450impl<P, RR, NR, TR, J, TCS, S, D> Relayer for StellarRelayer<P, RR, NR, TR, J, TCS, S, D>
451where
452    P: StellarProviderTrait + Send + Sync + 'static,
453    D: StellarDexServiceTrait + Send + Sync + 'static,
454    RR: Repository<RelayerRepoModel, String> + RelayerRepository + Send + Sync + 'static,
455    NR: NetworkRepository + Repository<NetworkRepoModel, String> + Send + Sync + 'static,
456    TR: Repository<TransactionRepoModel, String> + TransactionRepository + Send + Sync + 'static,
457    J: JobProducerTrait + Send + Sync + 'static,
458    TCS: TransactionCounterServiceTrait + Send + Sync + 'static,
459    S: StellarSignTrait + Send + Sync + 'static,
460{
461    #[instrument(
462        level = "debug",
463        skip(self, network_transaction),
464        fields(
465            request_id = ?crate::observability::request_id::get_request_id(),
466            relayer_id = %self.relayer.id,
467            network_type = ?self.relayer.network_type,
468        )
469    )]
470    async fn process_transaction_request(
471        &self,
472        network_transaction: NetworkTransactionRequest,
473    ) -> Result<TransactionRepoModel, RelayerError> {
474        let network_model = self
475            .network_repository
476            .get_by_name(NetworkType::Stellar, &self.relayer.network)
477            .await?
478            .ok_or_else(|| {
479                RelayerError::NetworkConfiguration(format!(
480                    "Network {} not found",
481                    self.relayer.network
482                ))
483            })?;
484        let transaction =
485            TransactionRepoModel::try_from((&network_transaction, &self.relayer, &network_model))?;
486
487        self.transaction_repository
488            .create(transaction.clone())
489            .await
490            .map_err(|e| RepositoryError::TransactionFailure(e.to_string()))?;
491
492        // Status check FIRST - this is our safety net for monitoring.
493        // If this fails, mark transaction as failed and don't proceed.
494        // This ensures we never have an unmonitored transaction.
495        if let Err(e) = self
496            .job_producer
497            .produce_check_transaction_status_job(
498                TransactionStatusCheck::new(
499                    transaction.id.clone(),
500                    transaction.relayer_id.clone(),
501                    NetworkType::Stellar,
502                ),
503                Some(calculate_scheduled_timestamp(
504                    STELLAR_STATUS_CHECK_INITIAL_DELAY_SECONDS,
505                )),
506            )
507            .await
508        {
509            // Status queue failed - mark transaction as failed to prevent orphaned tx
510            error!(
511                relayer_id = %self.relayer.id,
512                transaction_id = %transaction.id,
513                error = %e,
514                "Status check queue push failed - marking transaction as failed"
515            );
516            if let Err(update_err) = self
517                .transaction_repository
518                .partial_update(
519                    transaction.id.clone(),
520                    TransactionUpdateRequest {
521                        status: Some(TransactionStatus::Failed),
522                        status_reason: Some("Queue unavailable".to_string()),
523                        ..Default::default()
524                    },
525                )
526                .await
527            {
528                warn!(
529                    relayer_id = %self.relayer.id,
530                    transaction_id = %transaction.id,
531                    error = %update_err,
532                    "Failed to mark transaction as failed after queue push failure"
533                );
534            }
535            return Err(e.into());
536        }
537
538        // Now safe to push transaction request.
539        // Even if this fails, status check will monitor and detect the stuck transaction.
540        self.job_producer
541            .produce_transaction_request_job(
542                TransactionRequest::new(transaction.id.clone(), transaction.relayer_id.clone()),
543                None,
544            )
545            .await?;
546
547        Ok(transaction)
548    }
549
550    #[instrument(
551        level = "debug",
552        skip(self),
553        fields(
554            request_id = ?crate::observability::request_id::get_request_id(),
555            relayer_id = %self.relayer.id,
556        )
557    )]
558    async fn get_balance(&self) -> Result<BalanceResponse, RelayerError> {
559        let account_entry = self
560            .provider
561            .get_account(&self.relayer.address)
562            .await
563            .map_err(|e| {
564                warn!(
565                    relayer_id = %self.relayer.id,
566                    address = %self.relayer.address,
567                    error = %e,
568                    "get_account failed in get_balance (called before transaction creation)"
569                );
570                // Track RPC failure metric
571                crate::metrics::API_RPC_FAILURES
572                    .with_label_values(&[
573                        self.relayer.id.as_str(),
574                        "stellar",
575                        "get_balance",
576                        "get_account_failed",
577                    ])
578                    .inc();
579                RelayerError::ProviderError(format!("Failed to fetch account for balance: {e}"))
580            })?;
581
582        Ok(BalanceResponse {
583            balance: account_entry.balance as u128,
584            unit: STELLAR_SMALLEST_UNIT_NAME.to_string(),
585        })
586    }
587
588    #[instrument(
589        level = "debug",
590        skip(self),
591        fields(
592            request_id = ?crate::observability::request_id::get_request_id(),
593            relayer_id = %self.relayer.id,
594        )
595    )]
596    async fn get_status(&self) -> Result<RelayerStatus, RelayerError> {
597        let relayer_model = &self.relayer;
598
599        let account_entry = self
600            .provider
601            .get_account(&relayer_model.address)
602            .await
603            .map_err(|e| {
604                warn!(
605                    relayer_id = %relayer_model.id,
606                    address = %relayer_model.address,
607                    error = %e,
608                    "get_account failed in get_status (called before transaction creation)"
609                );
610                // Track RPC failure metric
611                crate::metrics::API_RPC_FAILURES
612                    .with_label_values(&[
613                        relayer_model.id.as_str(),
614                        "stellar",
615                        "get_status",
616                        "get_account_failed",
617                    ])
618                    .inc();
619                RelayerError::ProviderError(format!("Failed to get account details: {e}"))
620            })?;
621
622        let sequence_number_str = account_entry.seq_num.0.to_string();
623
624        let balance_response = self.get_balance().await?;
625
626        // Use optimized count_by_status
627        let pending_transactions_count = self
628            .transaction_repository
629            .count_by_status(&relayer_model.id, PENDING_TRANSACTION_STATUSES)
630            .await
631            .map_err(RelayerError::from)?;
632
633        // Use find_by_status_paginated to get the latest confirmed transaction (newest first)
634        let last_confirmed_transaction_timestamp = self
635            .transaction_repository
636            .find_by_status_paginated(
637                &relayer_model.id,
638                &[TransactionStatus::Confirmed],
639                PaginationQuery {
640                    page: 1,
641                    per_page: 1,
642                },
643                false, // oldest_first = false means newest first
644            )
645            .await
646            .map_err(RelayerError::from)?
647            .items
648            .into_iter()
649            .next()
650            .and_then(|tx| tx.confirmed_at);
651
652        Ok(RelayerStatus::Stellar {
653            balance: balance_response.balance.to_string(),
654            pending_transactions_count,
655            last_confirmed_transaction_timestamp,
656            system_disabled: relayer_model.system_disabled,
657            paused: relayer_model.paused,
658            sequence_number: sequence_number_str,
659        })
660    }
661
662    #[instrument(
663        level = "debug",
664        skip(self),
665        fields(
666            request_id = ?crate::observability::request_id::get_request_id(),
667            relayer_id = %self.relayer.id,
668        )
669    )]
670    async fn delete_pending_transactions(
671        &self,
672    ) -> Result<DeletePendingTransactionsResponse, RelayerError> {
673        println!("Stellar delete_pending_transactions...");
674        Ok(DeletePendingTransactionsResponse {
675            queued_for_cancellation_transaction_ids: vec![],
676            failed_to_queue_transaction_ids: vec![],
677            total_processed: 0,
678        })
679    }
680
681    #[instrument(
682        level = "debug",
683        skip(self, _request),
684        fields(
685            request_id = ?crate::observability::request_id::get_request_id(),
686            relayer_id = %self.relayer.id,
687        )
688    )]
689    async fn sign_data(&self, _request: SignDataRequest) -> Result<SignDataResponse, RelayerError> {
690        Err(RelayerError::NotSupported(
691            "Signing data not supported for Stellar".to_string(),
692        ))
693    }
694
695    #[instrument(
696        level = "debug",
697        skip(self, _request),
698        fields(
699            request_id = ?crate::observability::request_id::get_request_id(),
700            relayer_id = %self.relayer.id,
701        )
702    )]
703    async fn sign_typed_data(
704        &self,
705        _request: SignTypedDataRequest,
706    ) -> Result<SignDataResponse, RelayerError> {
707        Err(RelayerError::NotSupported(
708            "Signing typed data not supported for Stellar".to_string(),
709        ))
710    }
711
712    #[instrument(
713        level = "debug",
714        skip(self, request),
715        fields(
716            request_id = ?crate::observability::request_id::get_request_id(),
717            relayer_id = %self.relayer.id,
718        )
719    )]
720    async fn rpc(
721        &self,
722        request: JsonRpcRequest<NetworkRpcRequest>,
723    ) -> Result<JsonRpcResponse<NetworkRpcResult>, RelayerError> {
724        let JsonRpcRequest { id, params, .. } = request;
725        let stellar_request = match params {
726            NetworkRpcRequest::Stellar(stellar_req) => stellar_req,
727            _ => {
728                return Ok(create_error_response(
729                    id.clone(),
730                    RpcErrorCodes::INVALID_PARAMS,
731                    "Invalid params",
732                    "Expected Stellar network request",
733                ))
734            }
735        };
736
737        // Parse method and params from the Stellar request (single unified variant)
738        let (method, params_json) = match stellar_request {
739            StellarRpcRequest::RawRpcRequest { method, params } => (method, params),
740        };
741
742        match self
743            .provider
744            .raw_request_dyn(&method, params_json, id.clone())
745            .await
746        {
747            Ok(result_value) => Ok(create_success_response(id.clone(), result_value)),
748            Err(provider_error) => {
749                // Log the full error internally for debugging
750                tracing::error!(
751                    error = %provider_error,
752                    "RPC provider error occurred"
753                );
754                let (error_code, error_message) = map_provider_error(&provider_error);
755                let sanitized_description = sanitize_error_description(&provider_error);
756                Ok(create_error_response(
757                    id.clone(),
758                    error_code,
759                    error_message,
760                    &sanitized_description,
761                ))
762            }
763        }
764    }
765
766    #[instrument(
767        level = "debug",
768        skip(self),
769        fields(
770            request_id = ?crate::observability::request_id::get_request_id(),
771            relayer_id = %self.relayer.id,
772        )
773    )]
774    async fn validate_min_balance(&self) -> Result<(), RelayerError> {
775        Ok(())
776    }
777
778    #[instrument(
779        level = "debug",
780        skip(self),
781        fields(
782            request_id = ?crate::observability::request_id::get_request_id(),
783            relayer_id = %self.relayer.id,
784        )
785    )]
786    async fn initialize_relayer(&self) -> Result<(), RelayerError> {
787        debug!("initializing Stellar relayer");
788
789        // Migration: Check if relayer needs fee_payment_strategy migration
790        // Older relayers persisted in Redis may not have this policy set.
791        // We automatically set it to "Relayer" (the old default behavior) for backward compatibility.
792        self.migrate_fee_payment_strategy_if_needed().await?;
793
794        // Populate model with allowed token metadata and update DB entry
795        // Error will be thrown if any of the tokens are not found
796        self.populate_allowed_tokens_metadata().await.map_err(|e| {
797            RelayerError::PolicyConfigurationError(format!(
798                "Error while processing allowed tokens policy: {e}"
799            ))
800        })?;
801
802        match self.check_health().await {
803            Ok(_) => {
804                // All checks passed
805                if self.relayer.system_disabled {
806                    // Silently re-enable if was disabled (startup, not recovery)
807                    self.relayer_repository
808                        .enable_relayer(self.relayer.id.clone())
809                        .await?;
810                }
811            }
812            Err(failures) => {
813                // Health checks failed
814                let reason = DisabledReason::from_health_failures(failures).unwrap_or_else(|| {
815                    DisabledReason::SequenceSyncFailed("Unknown error".to_string())
816                });
817
818                warn!(reason = %reason, "disabling relayer");
819                let updated_relayer = self
820                    .relayer_repository
821                    .disable_relayer(self.relayer.id.clone(), reason.clone())
822                    .await?;
823
824                // Send notification if configured
825                if let Some(notification_id) = &self.relayer.notification_id {
826                    self.job_producer
827                        .produce_send_notification_job(
828                            produce_relayer_disabled_payload(
829                                notification_id,
830                                &updated_relayer,
831                                &reason.safe_description(),
832                            ),
833                            None,
834                        )
835                        .await?;
836                }
837
838                // Schedule health check to try re-enabling the relayer after 10 seconds
839                self.job_producer
840                    .produce_relayer_health_check_job(
841                        RelayerHealthCheck::new(self.relayer.id.clone()),
842                        Some(calculate_scheduled_timestamp(10)),
843                    )
844                    .await?;
845            }
846        }
847        debug!(
848            "Stellar relayer initialized successfully: {}",
849            self.relayer.id
850        );
851        Ok(())
852    }
853
854    #[instrument(
855        level = "debug",
856        skip(self),
857        fields(
858            request_id = ?crate::observability::request_id::get_request_id(),
859            relayer_id = %self.relayer.id,
860        )
861    )]
862    async fn check_health(&self) -> Result<(), Vec<HealthCheckFailure>> {
863        debug!(
864            "running health checks for Stellar relayer {}",
865            self.relayer.id
866        );
867
868        let mut failures = Vec::new();
869
870        // Check sequence synchronization
871        match self.sync_sequence().await {
872            Ok(_) => {
873                debug!(
874                    "sequence sync passed for Stellar relayer {}",
875                    self.relayer.id
876                );
877            }
878            Err(e) => {
879                let reason = HealthCheckFailure::SequenceSyncFailed(e.to_string());
880                warn!("sequence sync failed: {:?}", reason);
881                failures.push(reason);
882            }
883        }
884
885        // Check balance and trigger token swap if fee_payment_strategy is User
886        // Note: Swap failures are logged but don't cause health check failures
887        // to avoid disabling the relayer due to transient swap issues
888        let policy = self.relayer.policies.get_stellar_policy();
889        if matches!(
890            policy.fee_payment_strategy,
891            Some(StellarFeePaymentStrategy::User)
892        ) {
893            debug!(
894                "checking balance and attempting token swap for user fee payment strategy relayer {}",
895                self.relayer.id
896            );
897            if let Err(e) = self.check_balance_and_trigger_token_swap_if_needed().await {
898                warn!(
899                    relayer_id = %self.relayer.id,
900                    error = %e,
901                    "Balance check or token swap failed, but not treating as health check failure"
902                );
903            } else {
904                debug!(
905                    "balance check and token swap completed for Stellar relayer {}",
906                    self.relayer.id
907                );
908            }
909        }
910
911        if failures.is_empty() {
912            debug!(
913                "all health checks passed for Stellar relayer {}",
914                self.relayer.id
915            );
916            Ok(())
917        } else {
918            warn!(
919                "health checks failed for Stellar relayer {}: {:?}",
920                self.relayer.id, failures
921            );
922            Err(failures)
923        }
924    }
925
926    #[instrument(
927        level = "debug",
928        skip(self, request),
929        fields(
930            request_id = ?crate::observability::request_id::get_request_id(),
931            relayer_id = %self.relayer.id,
932        )
933    )]
934    async fn sign_transaction(
935        &self,
936        request: &SignTransactionRequest,
937    ) -> Result<SignTransactionExternalResponse, RelayerError> {
938        let stellar_req = match request {
939            SignTransactionRequest::Stellar(req) => req,
940            _ => {
941                return Err(RelayerError::NotSupported(
942                    "Invalid request type for Stellar relayer".to_string(),
943                ))
944            }
945        };
946
947        let policy = self.relayer.policies.get_stellar_policy();
948        let user_pays_fee = matches!(
949            policy.fee_payment_strategy,
950            Some(StellarFeePaymentStrategy::User)
951        );
952
953        // For user-paid fees, validate transaction before signing
954        if user_pays_fee {
955            // Parse the transaction XDR
956            let envelope = parse_transaction_xdr(&stellar_req.unsigned_xdr, false)
957                .map_err(|e| RelayerError::ValidationError(format!("Failed to parse XDR: {e}")))?;
958
959            // Comprehensive validation for user fee payment transactions when signing
960            // This validates: transaction structure, fee payments, allowed tokens, payment amounts, and time bounds
961            StellarTransactionValidator::validate_user_fee_payment_transaction(
962                &envelope,
963                &self.relayer.address,
964                &policy,
965                &self.provider,
966                self.dex_service.as_ref(),
967                Some(get_stellar_sponsored_transaction_validity_duration()), // Enforce 1 minute max validity for signing flow
968            )
969            .await
970            .map_err(|e| {
971                RelayerError::ValidationError(format!("Failed to validate transaction: {e}"))
972            })?;
973        }
974
975        // Use the signer's sign_xdr_transaction method
976        let response = self
977            .signer
978            .sign_xdr_transaction(&stellar_req.unsigned_xdr, &self.network.passphrase)
979            .await
980            .map_err(RelayerError::SignerError)?;
981
982        // Convert DecoratedSignature to base64 string
983        let signature_bytes = &response.signature.signature.0;
984        let signature_string =
985            base64::Engine::encode(&base64::engine::general_purpose::STANDARD, signature_bytes);
986
987        Ok(SignTransactionExternalResponse::Stellar(
988            SignTransactionExternalResponseStellar {
989                signed_xdr: response.signed_xdr,
990                signature: signature_string,
991            },
992        ))
993    }
994}
995
996#[cfg(test)]
997mod tests {
998    use super::*;
999    use crate::{
1000        config::{NetworkConfigCommon, StellarNetworkConfig},
1001        constants::STELLAR_SMALLEST_UNIT_NAME,
1002        domain::{SignTransactionRequestStellar, SignXdrTransactionResponseStellar},
1003        jobs::MockJobProducerTrait,
1004        models::{
1005            NetworkConfigData, NetworkRepoModel, NetworkType, RelayerNetworkPolicy,
1006            RelayerRepoModel, RelayerStellarPolicy, RpcConfig, SignerError,
1007        },
1008        repositories::{
1009            InMemoryNetworkRepository, MockRelayerRepository, MockTransactionRepository,
1010        },
1011        services::{
1012            provider::{MockStellarProviderTrait, ProviderError},
1013            signer::MockStellarSignTrait,
1014            stellar_dex::MockStellarDexServiceTrait,
1015            MockTransactionCounterServiceTrait,
1016        },
1017    };
1018    use mockall::predicate::*;
1019    use soroban_rs::xdr::{
1020        AccountEntry, AccountEntryExt, AccountId, DecoratedSignature, PublicKey, SequenceNumber,
1021        Signature, SignatureHint, String32, Thresholds, Uint256, VecM,
1022    };
1023    use std::future::ready;
1024    use std::sync::Arc;
1025
1026    /// Helper function to create a mock DEX service for testing
1027    fn create_mock_dex_service() -> Arc<MockStellarDexServiceTrait> {
1028        let mut mock_dex = MockStellarDexServiceTrait::new();
1029        mock_dex.expect_supported_asset_types().returning(|| {
1030            use crate::services::stellar_dex::AssetType;
1031            std::collections::HashSet::from([AssetType::Native, AssetType::Classic])
1032        });
1033        Arc::new(mock_dex)
1034    }
1035
1036    /// Test context structure to manage test dependencies
1037    struct TestCtx {
1038        relayer_model: RelayerRepoModel,
1039        network_repository: Arc<InMemoryNetworkRepository>,
1040    }
1041
1042    impl Default for TestCtx {
1043        fn default() -> Self {
1044            let network_repository = Arc::new(InMemoryNetworkRepository::new());
1045
1046            let relayer_model = RelayerRepoModel {
1047                id: "test-relayer-id".to_string(),
1048                name: "Test Relayer".to_string(),
1049                network: "testnet".to_string(),
1050                paused: false,
1051                network_type: NetworkType::Stellar,
1052                signer_id: "signer-id".to_string(),
1053                policies: RelayerNetworkPolicy::Stellar(RelayerStellarPolicy::default()),
1054                address: "GAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAWHF".to_string(),
1055                notification_id: Some("notification-id".to_string()),
1056                system_disabled: false,
1057                custom_rpc_urls: None,
1058                ..Default::default()
1059            };
1060
1061            TestCtx {
1062                relayer_model,
1063                network_repository,
1064            }
1065        }
1066    }
1067
1068    impl TestCtx {
1069        async fn setup_network(&self) {
1070            let test_network = NetworkRepoModel {
1071                id: "stellar:testnet".to_string(),
1072                name: "testnet".to_string(),
1073                network_type: NetworkType::Stellar,
1074                config: NetworkConfigData::Stellar(StellarNetworkConfig {
1075                    common: NetworkConfigCommon {
1076                        network: "testnet".to_string(),
1077                        from: None,
1078                        rpc_urls: Some(vec![RpcConfig::new(
1079                            "https://horizon-testnet.stellar.org".to_string(),
1080                        )]),
1081                        explorer_urls: None,
1082                        average_blocktime_ms: Some(5000),
1083                        is_testnet: Some(true),
1084                        tags: None,
1085                    },
1086                    passphrase: Some("Test SDF Network ; September 2015".to_string()),
1087                    horizon_url: Some("https://horizon-testnet.stellar.org".to_string()),
1088                }),
1089            };
1090
1091            self.network_repository.create(test_network).await.unwrap();
1092        }
1093    }
1094
1095    #[tokio::test]
1096    async fn test_sync_sequence_success() {
1097        let ctx = TestCtx::default();
1098        ctx.setup_network().await;
1099        let relayer_model = ctx.relayer_model.clone();
1100        let mut provider = MockStellarProviderTrait::new();
1101        provider
1102            .expect_get_account()
1103            .with(eq(relayer_model.address.clone()))
1104            .returning(|_| {
1105                Box::pin(async {
1106                    Ok(AccountEntry {
1107                        account_id: AccountId(PublicKey::PublicKeyTypeEd25519(Uint256([0; 32]))),
1108                        balance: 0,
1109                        ext: AccountEntryExt::V0,
1110                        flags: 0,
1111                        home_domain: String32::default(),
1112                        inflation_dest: None,
1113                        seq_num: SequenceNumber(5),
1114                        num_sub_entries: 0,
1115                        signers: VecM::default(),
1116                        thresholds: Thresholds([0, 0, 0, 0]),
1117                    })
1118                })
1119            });
1120        let mut counter = MockTransactionCounterServiceTrait::new();
1121        counter
1122            .expect_sync_floor()
1123            .with(eq(6u64))
1124            .returning(|floor| Box::pin(async move { Ok(floor) }));
1125        let relayer_repo = MockRelayerRepository::new();
1126        let tx_repo = MockTransactionRepository::new();
1127        let job_producer = MockJobProducerTrait::new();
1128        let signer = Arc::new(MockStellarSignTrait::new());
1129        let dex_service = create_mock_dex_service();
1130
1131        let relayer = StellarRelayer::new(
1132            relayer_model.clone(),
1133            signer,
1134            provider,
1135            StellarRelayerDependencies::new(
1136                Arc::new(relayer_repo),
1137                ctx.network_repository.clone(),
1138                Arc::new(tx_repo),
1139                Arc::new(counter),
1140                Arc::new(job_producer),
1141            ),
1142            dex_service,
1143        )
1144        .await
1145        .unwrap();
1146
1147        let result = relayer.sync_sequence().await;
1148        assert!(result.is_ok());
1149    }
1150
1151    /// Verifies `sync_sequence` never rewinds the counter below sequences already
1152    /// allocated by concurrent `get_and_increment` calls. On the multi-threaded runtime,
1153    /// a periodic/startup health-check sync racing with in-flight transaction submissions
1154    /// must only ever raise the counter to the chain floor, never blindly `set` it.
1155    #[tokio::test]
1156    async fn test_sync_sequence_does_not_rewind_counter() {
1157        let ctx = TestCtx::default();
1158        ctx.setup_network().await;
1159        let relayer_model = ctx.relayer_model.clone();
1160        let mut provider = MockStellarProviderTrait::new();
1161        // Chain reports on-chain sequence 5, so next usable sequence is 6.
1162        provider
1163            .expect_get_account()
1164            .with(eq(relayer_model.address.clone()))
1165            .returning(|_| {
1166                Box::pin(async {
1167                    Ok(AccountEntry {
1168                        account_id: AccountId(PublicKey::PublicKeyTypeEd25519(Uint256([0; 32]))),
1169                        balance: 0,
1170                        ext: AccountEntryExt::V0,
1171                        flags: 0,
1172                        home_domain: String32::default(),
1173                        inflation_dest: None,
1174                        seq_num: SequenceNumber(5),
1175                        num_sub_entries: 0,
1176                        signers: VecM::default(),
1177                        thresholds: Thresholds([0, 0, 0, 0]),
1178                    })
1179                })
1180            });
1181
1182        let mut counter = MockTransactionCounterServiceTrait::new();
1183        // Counter is already at 15 (advanced past the chain floor of 6 by concurrent
1184        // allocations). `sync_floor` must return the higher existing value, not rewind.
1185        counter
1186            .expect_sync_floor()
1187            .with(eq(6u64))
1188            .returning(|floor| Box::pin(async move { Ok(std::cmp::max(floor, 15)) }));
1189
1190        let relayer_repo = MockRelayerRepository::new();
1191        let tx_repo = MockTransactionRepository::new();
1192        let job_producer = MockJobProducerTrait::new();
1193        let signer = Arc::new(MockStellarSignTrait::new());
1194        let dex_service = create_mock_dex_service();
1195
1196        let relayer = StellarRelayer::new(
1197            relayer_model.clone(),
1198            signer,
1199            provider,
1200            StellarRelayerDependencies::new(
1201                Arc::new(relayer_repo),
1202                ctx.network_repository.clone(),
1203                Arc::new(tx_repo),
1204                Arc::new(counter),
1205                Arc::new(job_producer),
1206            ),
1207            dex_service,
1208        )
1209        .await
1210        .unwrap();
1211
1212        // sync_sequence succeeds and, because it goes through sync_floor rather than a
1213        // blind set, does not clobber the counter's higher, already-allocated value.
1214        let result = relayer.sync_sequence().await;
1215        assert!(result.is_ok());
1216    }
1217
1218    #[tokio::test]
1219    async fn test_sync_sequence_provider_error() {
1220        let ctx = TestCtx::default();
1221        ctx.setup_network().await;
1222        let relayer_model = ctx.relayer_model.clone();
1223        let mut provider = MockStellarProviderTrait::new();
1224        provider
1225            .expect_get_account()
1226            .with(eq(relayer_model.address.clone()))
1227            .returning(|_| Box::pin(async { Err(ProviderError::Other("fail".to_string())) }));
1228        let counter = MockTransactionCounterServiceTrait::new();
1229        let relayer_repo = MockRelayerRepository::new();
1230        let tx_repo = MockTransactionRepository::new();
1231        let job_producer = MockJobProducerTrait::new();
1232        let signer = Arc::new(MockStellarSignTrait::new());
1233        let dex_service = create_mock_dex_service();
1234
1235        let relayer = StellarRelayer::new(
1236            relayer_model.clone(),
1237            signer,
1238            provider,
1239            StellarRelayerDependencies::new(
1240                Arc::new(relayer_repo),
1241                ctx.network_repository.clone(),
1242                Arc::new(tx_repo),
1243                Arc::new(counter),
1244                Arc::new(job_producer),
1245            ),
1246            dex_service,
1247        )
1248        .await
1249        .unwrap();
1250
1251        let result = relayer.sync_sequence().await;
1252        assert!(matches!(result, Err(RelayerError::ProviderError(_))));
1253    }
1254
1255    #[tokio::test]
1256    async fn test_get_status_success_stellar() {
1257        let ctx = TestCtx::default();
1258        ctx.setup_network().await;
1259        let relayer_model = ctx.relayer_model.clone();
1260        let mut provider_mock = MockStellarProviderTrait::new();
1261        let mut tx_repo_mock = MockTransactionRepository::new();
1262        let relayer_repo_mock = MockRelayerRepository::new();
1263        let job_producer_mock = MockJobProducerTrait::new();
1264        let counter_mock = MockTransactionCounterServiceTrait::new();
1265
1266        provider_mock.expect_get_account().times(2).returning(|_| {
1267            Box::pin(ready(Ok(AccountEntry {
1268                account_id: AccountId(PublicKey::PublicKeyTypeEd25519(Uint256([0; 32]))),
1269                balance: 10000000,
1270                seq_num: SequenceNumber(12345),
1271                ext: AccountEntryExt::V0,
1272                flags: 0,
1273                home_domain: String32::default(),
1274                inflation_dest: None,
1275                num_sub_entries: 0,
1276                signers: VecM::default(),
1277                thresholds: Thresholds([0, 0, 0, 0]),
1278            })))
1279        });
1280
1281        // Mock count_by_status for pending transactions count
1282        tx_repo_mock
1283            .expect_count_by_status()
1284            .withf(|relayer_id, statuses| {
1285                relayer_id == "test-relayer-id"
1286                    && statuses
1287                        == [
1288                            TransactionStatus::Pending,
1289                            TransactionStatus::Sent,
1290                            TransactionStatus::Submitted,
1291                        ]
1292            })
1293            .returning(|_, _| Ok(0u64))
1294            .once();
1295
1296        // Mock find_by_status_paginated for latest confirmed transaction
1297        let confirmed_tx = TransactionRepoModel {
1298            id: "tx1_stellar".to_string(),
1299            relayer_id: relayer_model.id.clone(),
1300            status: TransactionStatus::Confirmed,
1301            confirmed_at: Some("2023-02-01T12:00:00Z".to_string()),
1302            ..TransactionRepoModel::default()
1303        };
1304        let relayer_id_clone = relayer_model.id.clone();
1305        tx_repo_mock
1306            .expect_find_by_status_paginated()
1307            .withf(move |relayer_id, statuses, query, oldest_first| {
1308                *relayer_id == relayer_id_clone
1309                    && statuses == [TransactionStatus::Confirmed]
1310                    && query.page == 1
1311                    && query.per_page == 1
1312                    && !(*oldest_first)
1313            })
1314            .returning(move |_, _, _, _| {
1315                Ok(crate::repositories::PaginatedResult {
1316                    items: vec![confirmed_tx.clone()],
1317                    total: 1,
1318                    page: 1,
1319                    per_page: 1,
1320                })
1321            })
1322            .once();
1323        let signer = Arc::new(MockStellarSignTrait::new());
1324        let dex_service = create_mock_dex_service();
1325
1326        let stellar_relayer = StellarRelayer::new(
1327            relayer_model.clone(),
1328            signer,
1329            provider_mock,
1330            StellarRelayerDependencies::new(
1331                Arc::new(relayer_repo_mock),
1332                ctx.network_repository.clone(),
1333                Arc::new(tx_repo_mock),
1334                Arc::new(counter_mock),
1335                Arc::new(job_producer_mock),
1336            ),
1337            dex_service,
1338        )
1339        .await
1340        .unwrap();
1341
1342        let status = stellar_relayer.get_status().await.unwrap();
1343
1344        match status {
1345            RelayerStatus::Stellar {
1346                balance,
1347                pending_transactions_count,
1348                last_confirmed_transaction_timestamp,
1349                system_disabled,
1350                paused,
1351                sequence_number,
1352            } => {
1353                assert_eq!(balance, "10000000");
1354                assert_eq!(pending_transactions_count, 0);
1355                assert_eq!(
1356                    last_confirmed_transaction_timestamp,
1357                    Some("2023-02-01T12:00:00Z".to_string())
1358                );
1359                assert_eq!(system_disabled, relayer_model.system_disabled);
1360                assert_eq!(paused, relayer_model.paused);
1361                assert_eq!(sequence_number, "12345");
1362            }
1363            _ => panic!("Expected Stellar RelayerStatus"),
1364        }
1365    }
1366
1367    #[tokio::test]
1368    async fn test_get_status_stellar_provider_error() {
1369        let ctx = TestCtx::default();
1370        ctx.setup_network().await;
1371        let relayer_model = ctx.relayer_model.clone();
1372        let mut provider_mock = MockStellarProviderTrait::new();
1373        let tx_repo_mock = MockTransactionRepository::new();
1374        let relayer_repo_mock = MockRelayerRepository::new();
1375        let job_producer_mock = MockJobProducerTrait::new();
1376        let counter_mock = MockTransactionCounterServiceTrait::new();
1377
1378        provider_mock
1379            .expect_get_account()
1380            .with(eq(relayer_model.address.clone()))
1381            .returning(|_| {
1382                Box::pin(async { Err(ProviderError::Other("Stellar provider down".to_string())) })
1383            });
1384        let signer = Arc::new(MockStellarSignTrait::new());
1385        let dex_service = create_mock_dex_service();
1386
1387        let stellar_relayer = StellarRelayer::new(
1388            relayer_model.clone(),
1389            signer,
1390            provider_mock,
1391            StellarRelayerDependencies::new(
1392                Arc::new(relayer_repo_mock),
1393                ctx.network_repository.clone(),
1394                Arc::new(tx_repo_mock),
1395                Arc::new(counter_mock),
1396                Arc::new(job_producer_mock),
1397            ),
1398            dex_service,
1399        )
1400        .await
1401        .unwrap();
1402
1403        let result = stellar_relayer.get_status().await;
1404        assert!(result.is_err());
1405        match result.err().unwrap() {
1406            RelayerError::ProviderError(msg) => {
1407                assert!(msg.contains("Failed to get account details"))
1408            }
1409            _ => panic!("Expected ProviderError for get_account failure"),
1410        }
1411    }
1412
1413    #[tokio::test]
1414    async fn test_get_balance_success() {
1415        let ctx = TestCtx::default();
1416        ctx.setup_network().await;
1417        let relayer_model = ctx.relayer_model.clone();
1418        let mut provider = MockStellarProviderTrait::new();
1419        let expected_balance = 100_000_000i64; // 10 XLM in stroops
1420
1421        provider
1422            .expect_get_account()
1423            .with(eq(relayer_model.address.clone()))
1424            .returning(move |_| {
1425                Box::pin(async move {
1426                    Ok(AccountEntry {
1427                        account_id: AccountId(PublicKey::PublicKeyTypeEd25519(Uint256([0; 32]))),
1428                        balance: expected_balance,
1429                        ext: AccountEntryExt::V0,
1430                        flags: 0,
1431                        home_domain: String32::default(),
1432                        inflation_dest: None,
1433                        seq_num: SequenceNumber(5),
1434                        num_sub_entries: 0,
1435                        signers: VecM::default(),
1436                        thresholds: Thresholds([0, 0, 0, 0]),
1437                    })
1438                })
1439            });
1440
1441        let relayer_repo = Arc::new(MockRelayerRepository::new());
1442        let tx_repo = Arc::new(MockTransactionRepository::new());
1443        let job_producer = Arc::new(MockJobProducerTrait::new());
1444        let counter = Arc::new(MockTransactionCounterServiceTrait::new());
1445        let signer = Arc::new(MockStellarSignTrait::new());
1446        let dex_service = create_mock_dex_service();
1447
1448        let relayer = StellarRelayer::new(
1449            relayer_model,
1450            signer,
1451            provider,
1452            StellarRelayerDependencies::new(
1453                relayer_repo,
1454                ctx.network_repository.clone(),
1455                tx_repo,
1456                counter,
1457                job_producer,
1458            ),
1459            dex_service,
1460        )
1461        .await
1462        .unwrap();
1463
1464        let result = relayer.get_balance().await;
1465        assert!(result.is_ok());
1466        let balance_response = result.unwrap();
1467        assert_eq!(balance_response.balance, expected_balance as u128);
1468        assert_eq!(balance_response.unit, STELLAR_SMALLEST_UNIT_NAME);
1469    }
1470
1471    #[tokio::test]
1472    async fn test_get_balance_provider_error() {
1473        let ctx = TestCtx::default();
1474        ctx.setup_network().await;
1475        let relayer_model = ctx.relayer_model.clone();
1476        let mut provider = MockStellarProviderTrait::new();
1477
1478        provider
1479            .expect_get_account()
1480            .with(eq(relayer_model.address.clone()))
1481            .returning(|_| {
1482                Box::pin(async { Err(ProviderError::Other("provider failed".to_string())) })
1483            });
1484
1485        let relayer_repo = Arc::new(MockRelayerRepository::new());
1486        let tx_repo = Arc::new(MockTransactionRepository::new());
1487        let job_producer = Arc::new(MockJobProducerTrait::new());
1488        let counter = Arc::new(MockTransactionCounterServiceTrait::new());
1489        let signer = Arc::new(MockStellarSignTrait::new());
1490        let dex_service = create_mock_dex_service();
1491
1492        let relayer = StellarRelayer::new(
1493            relayer_model,
1494            signer,
1495            provider,
1496            StellarRelayerDependencies::new(
1497                relayer_repo,
1498                ctx.network_repository.clone(),
1499                tx_repo,
1500                counter,
1501                job_producer,
1502            ),
1503            dex_service,
1504        )
1505        .await
1506        .unwrap();
1507
1508        let result = relayer.get_balance().await;
1509        assert!(result.is_err());
1510        match result.err().unwrap() {
1511            RelayerError::ProviderError(msg) => {
1512                assert!(msg.contains("Failed to fetch account for balance"));
1513            }
1514            _ => panic!("Unexpected error type"),
1515        }
1516    }
1517
1518    #[tokio::test]
1519    async fn test_sign_transaction_success() {
1520        let ctx = TestCtx::default();
1521        ctx.setup_network().await;
1522        let relayer_model = ctx.relayer_model.clone();
1523        let provider = MockStellarProviderTrait::new();
1524        let mut signer = MockStellarSignTrait::new();
1525
1526        let unsigned_xdr = "AAAAAgAAAAD///8AAAAAAAAAAQAAAAAAAAACAAAAAQAAAAAAAAAB";
1527        let expected_signed_xdr =
1528            "AAAAAgAAAAD///8AAAAAAAABAAAAAAAAAAIAAAABAAAAAAAAAAEAAAABAAAAA...";
1529        let expected_signature = DecoratedSignature {
1530            hint: SignatureHint([1, 2, 3, 4]),
1531            signature: Signature([5u8; 64].try_into().unwrap()),
1532        };
1533        let expected_signature_for_closure = expected_signature.clone();
1534
1535        signer
1536            .expect_sign_xdr_transaction()
1537            .with(eq(unsigned_xdr), eq("Test SDF Network ; September 2015"))
1538            .returning(move |_, _| {
1539                Ok(SignXdrTransactionResponseStellar {
1540                    signed_xdr: expected_signed_xdr.to_string(),
1541                    signature: expected_signature_for_closure.clone(),
1542                })
1543            });
1544
1545        let relayer_repo = Arc::new(MockRelayerRepository::new());
1546        let tx_repo = Arc::new(MockTransactionRepository::new());
1547        let job_producer = Arc::new(MockJobProducerTrait::new());
1548        let counter = Arc::new(MockTransactionCounterServiceTrait::new());
1549        let dex_service = create_mock_dex_service();
1550
1551        let relayer = StellarRelayer::new(
1552            relayer_model,
1553            Arc::new(signer),
1554            provider,
1555            StellarRelayerDependencies::new(
1556                relayer_repo,
1557                ctx.network_repository.clone(),
1558                tx_repo,
1559                counter,
1560                job_producer,
1561            ),
1562            dex_service,
1563        )
1564        .await
1565        .unwrap();
1566
1567        let request = SignTransactionRequest::Stellar(SignTransactionRequestStellar {
1568            unsigned_xdr: unsigned_xdr.to_string(),
1569        });
1570        let result = relayer.sign_transaction(&request).await;
1571        assert!(result.is_ok());
1572
1573        match result.unwrap() {
1574            SignTransactionExternalResponse::Stellar(response) => {
1575                assert_eq!(response.signed_xdr, expected_signed_xdr);
1576                // Compare the base64 encoded signature
1577                let expected_signature_base64 = base64::Engine::encode(
1578                    &base64::engine::general_purpose::STANDARD,
1579                    &expected_signature.signature.0,
1580                );
1581                assert_eq!(response.signature, expected_signature_base64);
1582            }
1583            _ => panic!("Expected Stellar response"),
1584        }
1585    }
1586
1587    #[tokio::test]
1588    async fn test_sign_transaction_signer_error() {
1589        let ctx = TestCtx::default();
1590        ctx.setup_network().await;
1591        let relayer_model = ctx.relayer_model.clone();
1592        let provider = MockStellarProviderTrait::new();
1593        let mut signer = MockStellarSignTrait::new();
1594
1595        let unsigned_xdr = "INVALID_XDR";
1596
1597        signer
1598            .expect_sign_xdr_transaction()
1599            .with(eq(unsigned_xdr), eq("Test SDF Network ; September 2015"))
1600            .returning(|_, _| Err(SignerError::SigningError("Invalid XDR format".to_string())));
1601
1602        let relayer_repo = Arc::new(MockRelayerRepository::new());
1603        let tx_repo = Arc::new(MockTransactionRepository::new());
1604        let job_producer = Arc::new(MockJobProducerTrait::new());
1605        let counter = Arc::new(MockTransactionCounterServiceTrait::new());
1606        let dex_service = create_mock_dex_service();
1607
1608        let relayer = StellarRelayer::new(
1609            relayer_model,
1610            Arc::new(signer),
1611            provider,
1612            StellarRelayerDependencies::new(
1613                relayer_repo,
1614                ctx.network_repository.clone(),
1615                tx_repo,
1616                counter,
1617                job_producer,
1618            ),
1619            dex_service,
1620        )
1621        .await
1622        .unwrap();
1623
1624        let request = SignTransactionRequest::Stellar(SignTransactionRequestStellar {
1625            unsigned_xdr: unsigned_xdr.to_string(),
1626        });
1627        let result = relayer.sign_transaction(&request).await;
1628        assert!(result.is_err());
1629
1630        match result.err().unwrap() {
1631            RelayerError::SignerError(err) => match err {
1632                SignerError::SigningError(msg) => {
1633                    assert_eq!(msg, "Invalid XDR format");
1634                }
1635                _ => panic!("Expected SigningError"),
1636            },
1637            _ => panic!("Expected RelayerError::SignerError"),
1638        }
1639    }
1640
1641    #[tokio::test]
1642    async fn test_sign_transaction_with_different_network_passphrase() {
1643        let ctx = TestCtx::default();
1644        // Create a custom network with a different passphrase
1645        let custom_network = NetworkRepoModel {
1646            id: "stellar:mainnet".to_string(),
1647            name: "mainnet".to_string(),
1648            network_type: NetworkType::Stellar,
1649            config: NetworkConfigData::Stellar(StellarNetworkConfig {
1650                common: NetworkConfigCommon {
1651                    network: "mainnet".to_string(),
1652                    from: None,
1653                    rpc_urls: Some(vec![RpcConfig::new(
1654                        "https://horizon.stellar.org".to_string(),
1655                    )]),
1656                    explorer_urls: None,
1657                    average_blocktime_ms: Some(5000),
1658                    is_testnet: Some(false),
1659                    tags: None,
1660                },
1661                passphrase: Some("Public Global Stellar Network ; September 2015".to_string()),
1662                horizon_url: Some("https://horizon.stellar.org".to_string()),
1663            }),
1664        };
1665        ctx.network_repository.create(custom_network).await.unwrap();
1666
1667        let mut relayer_model = ctx.relayer_model.clone();
1668        relayer_model.network = "mainnet".to_string();
1669
1670        let provider = MockStellarProviderTrait::new();
1671        let mut signer = MockStellarSignTrait::new();
1672
1673        let unsigned_xdr = "AAAAAgAAAAD///8AAAAAAAAAAQAAAAAAAAACAAAAAQAAAAAAAAAB";
1674        let expected_signature = DecoratedSignature {
1675            hint: SignatureHint([10, 20, 30, 40]),
1676            signature: Signature([15u8; 64].try_into().unwrap()),
1677        };
1678        let expected_signature_for_closure = expected_signature.clone();
1679
1680        signer
1681            .expect_sign_xdr_transaction()
1682            .with(
1683                eq(unsigned_xdr),
1684                eq("Public Global Stellar Network ; September 2015"),
1685            )
1686            .returning(move |_, _| {
1687                Ok(SignXdrTransactionResponseStellar {
1688                    signed_xdr: "mainnet_signed_xdr".to_string(),
1689                    signature: expected_signature_for_closure.clone(),
1690                })
1691            });
1692
1693        let relayer_repo = Arc::new(MockRelayerRepository::new());
1694        let tx_repo = Arc::new(MockTransactionRepository::new());
1695        let job_producer = Arc::new(MockJobProducerTrait::new());
1696        let counter = Arc::new(MockTransactionCounterServiceTrait::new());
1697        let dex_service = create_mock_dex_service();
1698
1699        let relayer = StellarRelayer::new(
1700            relayer_model,
1701            Arc::new(signer),
1702            provider,
1703            StellarRelayerDependencies::new(
1704                relayer_repo,
1705                ctx.network_repository.clone(),
1706                tx_repo,
1707                counter,
1708                job_producer,
1709            ),
1710            dex_service,
1711        )
1712        .await
1713        .unwrap();
1714
1715        let request = SignTransactionRequest::Stellar(SignTransactionRequestStellar {
1716            unsigned_xdr: unsigned_xdr.to_string(),
1717        });
1718        let result = relayer.sign_transaction(&request).await;
1719        assert!(result.is_ok());
1720
1721        match result.unwrap() {
1722            SignTransactionExternalResponse::Stellar(response) => {
1723                assert_eq!(response.signed_xdr, "mainnet_signed_xdr");
1724                // Convert expected signature to base64 for comparison (just the signature bytes, not the whole struct)
1725                let expected_signature_string = base64::Engine::encode(
1726                    &base64::engine::general_purpose::STANDARD,
1727                    &expected_signature.signature.0,
1728                );
1729                assert_eq!(response.signature, expected_signature_string);
1730            }
1731            _ => panic!("Expected Stellar response"),
1732        }
1733    }
1734
1735    #[tokio::test]
1736    async fn test_initialize_relayer_disables_when_validation_fails() {
1737        let ctx = TestCtx::default();
1738        ctx.setup_network().await;
1739        let mut relayer_model = ctx.relayer_model.clone();
1740        relayer_model.system_disabled = false; // Start as enabled
1741        relayer_model.notification_id = Some("test-notification-id".to_string());
1742
1743        let mut provider = MockStellarProviderTrait::new();
1744        let mut relayer_repo = MockRelayerRepository::new();
1745        let mut job_producer = MockJobProducerTrait::new();
1746
1747        relayer_repo
1748            .expect_is_persistent_storage()
1749            .returning(|| false);
1750
1751        // Mock validation failure - sequence sync fails
1752        provider
1753            .expect_get_account()
1754            .returning(|_| Box::pin(ready(Err(ProviderError::Other("RPC error".to_string())))));
1755
1756        // Mock disable_relayer call
1757        let mut disabled_relayer = relayer_model.clone();
1758        disabled_relayer.system_disabled = true;
1759        relayer_repo
1760            .expect_disable_relayer()
1761            .withf(|id, reason| {
1762                id == "test-relayer-id"
1763                    && matches!(reason, crate::models::DisabledReason::SequenceSyncFailed(_))
1764            })
1765            .returning(move |_, _| Ok(disabled_relayer.clone()));
1766
1767        // Mock notification job production
1768        job_producer
1769            .expect_produce_send_notification_job()
1770            .returning(|_, _| Box::pin(async { Ok(()) }));
1771
1772        // Mock health check job scheduling
1773        job_producer
1774            .expect_produce_relayer_health_check_job()
1775            .returning(|_, _| Box::pin(async { Ok(()) }));
1776
1777        let tx_repo = MockTransactionRepository::new();
1778        let counter = MockTransactionCounterServiceTrait::new();
1779        let signer = Arc::new(MockStellarSignTrait::new());
1780        let dex_service = create_mock_dex_service();
1781
1782        let relayer = StellarRelayer::new(
1783            relayer_model.clone(),
1784            signer,
1785            provider,
1786            StellarRelayerDependencies::new(
1787                Arc::new(relayer_repo),
1788                ctx.network_repository.clone(),
1789                Arc::new(tx_repo),
1790                Arc::new(counter),
1791                Arc::new(job_producer),
1792            ),
1793            dex_service,
1794        )
1795        .await
1796        .unwrap();
1797
1798        let result = relayer.initialize_relayer().await;
1799        assert!(result.is_ok());
1800    }
1801
1802    #[tokio::test]
1803    async fn test_initialize_relayer_enables_when_validation_passes_and_was_disabled() {
1804        let ctx = TestCtx::default();
1805        ctx.setup_network().await;
1806        let mut relayer_model = ctx.relayer_model.clone();
1807        relayer_model.system_disabled = true; // Start as disabled
1808
1809        let mut provider = MockStellarProviderTrait::new();
1810        let mut relayer_repo = MockRelayerRepository::new();
1811
1812        relayer_repo
1813            .expect_is_persistent_storage()
1814            .returning(|| false);
1815
1816        // Mock successful validations - sequence sync succeeds
1817        provider.expect_get_account().returning(|_| {
1818            Box::pin(ready(Ok(AccountEntry {
1819                account_id: AccountId(PublicKey::PublicKeyTypeEd25519(Uint256([0; 32]))),
1820                balance: 1000000000, // 100 XLM
1821                seq_num: SequenceNumber(1),
1822                num_sub_entries: 0,
1823                inflation_dest: None,
1824                flags: 0,
1825                home_domain: String32::default(),
1826                thresholds: Thresholds([0; 4]),
1827                signers: VecM::default(),
1828                ext: AccountEntryExt::V0,
1829            })))
1830        });
1831
1832        // Mock enable_relayer call
1833        let mut enabled_relayer = relayer_model.clone();
1834        enabled_relayer.system_disabled = false;
1835        relayer_repo
1836            .expect_enable_relayer()
1837            .with(eq("test-relayer-id".to_string()))
1838            .returning(move |_| Ok(enabled_relayer.clone()));
1839
1840        let tx_repo = MockTransactionRepository::new();
1841        let mut counter = MockTransactionCounterServiceTrait::new();
1842        counter
1843            .expect_sync_floor()
1844            .returning(|floor| Box::pin(async move { Ok(floor) }));
1845        let signer = Arc::new(MockStellarSignTrait::new());
1846        let dex_service = create_mock_dex_service();
1847        let job_producer = MockJobProducerTrait::new();
1848
1849        let relayer = StellarRelayer::new(
1850            relayer_model.clone(),
1851            signer,
1852            provider,
1853            StellarRelayerDependencies::new(
1854                Arc::new(relayer_repo),
1855                ctx.network_repository.clone(),
1856                Arc::new(tx_repo),
1857                Arc::new(counter),
1858                Arc::new(job_producer),
1859            ),
1860            dex_service,
1861        )
1862        .await
1863        .unwrap();
1864
1865        let result = relayer.initialize_relayer().await;
1866        assert!(result.is_ok());
1867    }
1868
1869    #[tokio::test]
1870    async fn test_initialize_relayer_no_action_when_enabled_and_validation_passes() {
1871        let ctx = TestCtx::default();
1872        ctx.setup_network().await;
1873        let mut relayer_model = ctx.relayer_model.clone();
1874        relayer_model.system_disabled = false; // Start as enabled
1875
1876        let mut provider = MockStellarProviderTrait::new();
1877
1878        // Mock successful validations - sequence sync succeeds
1879        provider.expect_get_account().returning(|_| {
1880            Box::pin(ready(Ok(AccountEntry {
1881                account_id: AccountId(PublicKey::PublicKeyTypeEd25519(Uint256([0; 32]))),
1882                balance: 1000000000, // 100 XLM
1883                seq_num: SequenceNumber(1),
1884                num_sub_entries: 0,
1885                inflation_dest: None,
1886                flags: 0,
1887                home_domain: String32::default(),
1888                thresholds: Thresholds([0; 4]),
1889                signers: VecM::default(),
1890                ext: AccountEntryExt::V0,
1891            })))
1892        });
1893
1894        // No repository calls should be made since relayer is already enabled
1895
1896        let tx_repo = MockTransactionRepository::new();
1897        let mut counter = MockTransactionCounterServiceTrait::new();
1898        counter
1899            .expect_sync_floor()
1900            .returning(|floor| Box::pin(async move { Ok(floor) }));
1901        let signer = Arc::new(MockStellarSignTrait::new());
1902        let dex_service = create_mock_dex_service();
1903        let job_producer = MockJobProducerTrait::new();
1904        let mut relayer_repo = MockRelayerRepository::new();
1905
1906        relayer_repo
1907            .expect_is_persistent_storage()
1908            .returning(|| false);
1909
1910        let relayer = StellarRelayer::new(
1911            relayer_model.clone(),
1912            signer,
1913            provider,
1914            StellarRelayerDependencies::new(
1915                Arc::new(relayer_repo),
1916                ctx.network_repository.clone(),
1917                Arc::new(tx_repo),
1918                Arc::new(counter),
1919                Arc::new(job_producer),
1920            ),
1921            dex_service,
1922        )
1923        .await
1924        .unwrap();
1925
1926        let result = relayer.initialize_relayer().await;
1927        assert!(result.is_ok());
1928    }
1929
1930    #[tokio::test]
1931    async fn test_initialize_relayer_sends_notification_when_disabled() {
1932        let ctx = TestCtx::default();
1933        ctx.setup_network().await;
1934        let mut relayer_model = ctx.relayer_model.clone();
1935        relayer_model.system_disabled = false; // Start as enabled
1936        relayer_model.notification_id = Some("test-notification-id".to_string());
1937
1938        let mut provider = MockStellarProviderTrait::new();
1939        let mut relayer_repo = MockRelayerRepository::new();
1940        let mut job_producer = MockJobProducerTrait::new();
1941
1942        relayer_repo
1943            .expect_is_persistent_storage()
1944            .returning(|| false);
1945
1946        // Mock validation failure - sequence sync fails
1947        provider.expect_get_account().returning(|_| {
1948            Box::pin(ready(Err(ProviderError::Other(
1949                "Sequence sync failed".to_string(),
1950            ))))
1951        });
1952
1953        // Mock disable_relayer call
1954        let mut disabled_relayer = relayer_model.clone();
1955        disabled_relayer.system_disabled = true;
1956        relayer_repo
1957            .expect_disable_relayer()
1958            .withf(|id, reason| {
1959                id == "test-relayer-id"
1960                    && matches!(reason, crate::models::DisabledReason::SequenceSyncFailed(_))
1961            })
1962            .returning(move |_, _| Ok(disabled_relayer.clone()));
1963
1964        // Mock notification job production - verify it's called
1965        job_producer
1966            .expect_produce_send_notification_job()
1967            .returning(|_, _| Box::pin(async { Ok(()) }));
1968
1969        // Mock health check job scheduling
1970        job_producer
1971            .expect_produce_relayer_health_check_job()
1972            .returning(|_, _| Box::pin(async { Ok(()) }));
1973
1974        let tx_repo = MockTransactionRepository::new();
1975        let counter = MockTransactionCounterServiceTrait::new();
1976        let signer = Arc::new(MockStellarSignTrait::new());
1977        let dex_service = create_mock_dex_service();
1978
1979        let relayer = StellarRelayer::new(
1980            relayer_model.clone(),
1981            signer,
1982            provider,
1983            StellarRelayerDependencies::new(
1984                Arc::new(relayer_repo),
1985                ctx.network_repository.clone(),
1986                Arc::new(tx_repo),
1987                Arc::new(counter),
1988                Arc::new(job_producer),
1989            ),
1990            dex_service,
1991        )
1992        .await
1993        .unwrap();
1994
1995        let result = relayer.initialize_relayer().await;
1996        assert!(result.is_ok());
1997    }
1998
1999    #[tokio::test]
2000    async fn test_initialize_relayer_no_notification_when_no_notification_id() {
2001        let ctx = TestCtx::default();
2002        ctx.setup_network().await;
2003        let mut relayer_model = ctx.relayer_model.clone();
2004        relayer_model.system_disabled = false; // Start as enabled
2005        relayer_model.notification_id = None; // No notification ID
2006
2007        let mut provider = MockStellarProviderTrait::new();
2008        let mut relayer_repo = MockRelayerRepository::new();
2009        relayer_repo
2010            .expect_is_persistent_storage()
2011            .returning(|| false);
2012
2013        // Mock validation failure - sequence sync fails
2014        provider.expect_get_account().returning(|_| {
2015            Box::pin(ready(Err(ProviderError::Other(
2016                "Sequence sync failed".to_string(),
2017            ))))
2018        });
2019
2020        // Mock disable_relayer call
2021        let mut disabled_relayer = relayer_model.clone();
2022        disabled_relayer.system_disabled = true;
2023        relayer_repo
2024            .expect_disable_relayer()
2025            .withf(|id, reason| {
2026                id == "test-relayer-id"
2027                    && matches!(reason, crate::models::DisabledReason::SequenceSyncFailed(_))
2028            })
2029            .returning(move |_, _| Ok(disabled_relayer.clone()));
2030
2031        // No notification job should be produced since notification_id is None
2032        // But health check job should still be scheduled
2033        let mut job_producer = MockJobProducerTrait::new();
2034        job_producer
2035            .expect_produce_relayer_health_check_job()
2036            .returning(|_, _| Box::pin(async { Ok(()) }));
2037
2038        let tx_repo = MockTransactionRepository::new();
2039        let counter = MockTransactionCounterServiceTrait::new();
2040        let signer = Arc::new(MockStellarSignTrait::new());
2041        let dex_service = create_mock_dex_service();
2042
2043        let relayer = StellarRelayer::new(
2044            relayer_model.clone(),
2045            signer,
2046            provider,
2047            StellarRelayerDependencies::new(
2048                Arc::new(relayer_repo),
2049                ctx.network_repository.clone(),
2050                Arc::new(tx_repo),
2051                Arc::new(counter),
2052                Arc::new(job_producer),
2053            ),
2054            dex_service,
2055        )
2056        .await
2057        .unwrap();
2058
2059        let result = relayer.initialize_relayer().await;
2060        assert!(result.is_ok());
2061    }
2062
2063    mod process_transaction_request_tests {
2064        use super::*;
2065        use crate::constants::STELLAR_STATUS_CHECK_INITIAL_DELAY_SECONDS;
2066        use crate::models::{
2067            NetworkTransactionRequest, NetworkType, StellarTransactionRequest, TransactionStatus,
2068        };
2069        use chrono::Utc;
2070
2071        // Helper function to create a valid test transaction request
2072        fn create_test_transaction_request() -> NetworkTransactionRequest {
2073            NetworkTransactionRequest::Stellar(StellarTransactionRequest {
2074                source_account: None,
2075                network: "testnet".to_string(),
2076                operations: None,
2077                memo: None,
2078                valid_until: None,
2079                transaction_xdr: Some("AAAAAgAAAACige4lTdwSB/sto4SniEdJ2kOa2X65s5bqkd40J4DjSwAAAAEAAHAkAAAADwAAAAAAAAAAAAAAAQAAAAAAAAABAAAAAKKB7iVN3BIH+y2jhKeIR0naQ5rZfrmzluqR3jQngONLAAAAAAAAAAAAD0JAAAAAAAAAAAA=".to_string()),
2080                fee_bump: None,
2081                max_fee: None,
2082                signed_auth_entry: None,
2083            })
2084        }
2085
2086        #[tokio::test]
2087        async fn test_process_transaction_request_calls_job_producer_methods() {
2088            let ctx = TestCtx::default();
2089            ctx.setup_network().await;
2090            let relayer_model = ctx.relayer_model.clone();
2091
2092            let provider = MockStellarProviderTrait::new();
2093            let signer = Arc::new(MockStellarSignTrait::new());
2094            let dex_service = create_mock_dex_service();
2095
2096            // Create a test transaction request
2097            let tx_request = create_test_transaction_request();
2098
2099            // Mock transaction repository - we expect it to create a transaction
2100            let mut tx_repo = MockTransactionRepository::new();
2101            tx_repo.expect_create().returning(|t| Ok(t.clone()));
2102
2103            // Mock job producer to verify both methods are called
2104            let mut job_producer = MockJobProducerTrait::new();
2105
2106            // Verify produce_transaction_request_job is called
2107            job_producer
2108                .expect_produce_transaction_request_job()
2109                .withf(|req, delay| {
2110                    !req.transaction_id.is_empty() && !req.relayer_id.is_empty() && delay.is_none()
2111                })
2112                .times(1)
2113                .returning(|_, _| Box::pin(async { Ok(()) }));
2114
2115            // Verify produce_check_transaction_status_job is called with correct parameters
2116            job_producer
2117                .expect_produce_check_transaction_status_job()
2118                .withf(|check, delay| {
2119                    !check.transaction_id.is_empty()
2120                        && !check.relayer_id.is_empty()
2121                        && check.network_type == Some(NetworkType::Stellar)
2122                        && delay.is_some()
2123                })
2124                .times(1)
2125                .returning(|_, _| Box::pin(async { Ok(()) }));
2126
2127            let relayer_repo = Arc::new(MockRelayerRepository::new());
2128            let counter = MockTransactionCounterServiceTrait::new();
2129
2130            let relayer = StellarRelayer::new(
2131                relayer_model,
2132                signer,
2133                provider,
2134                StellarRelayerDependencies::new(
2135                    relayer_repo,
2136                    ctx.network_repository.clone(),
2137                    Arc::new(tx_repo),
2138                    Arc::new(counter),
2139                    Arc::new(job_producer),
2140                ),
2141                dex_service,
2142            )
2143            .await
2144            .unwrap();
2145
2146            let result = relayer.process_transaction_request(tx_request).await;
2147            if let Err(e) = &result {
2148                panic!("process_transaction_request failed: {e}");
2149            }
2150            assert!(result.is_ok());
2151        }
2152
2153        #[tokio::test]
2154        async fn test_process_transaction_request_with_scheduled_delay() {
2155            let ctx = TestCtx::default();
2156            ctx.setup_network().await;
2157            let relayer_model = ctx.relayer_model.clone();
2158
2159            let provider = MockStellarProviderTrait::new();
2160            let signer = Arc::new(MockStellarSignTrait::new());
2161            let dex_service = create_mock_dex_service();
2162
2163            let tx_request = create_test_transaction_request();
2164
2165            let mut tx_repo = MockTransactionRepository::new();
2166            tx_repo.expect_create().returning(|t| Ok(t.clone()));
2167
2168            let mut job_producer = MockJobProducerTrait::new();
2169
2170            job_producer
2171                .expect_produce_transaction_request_job()
2172                .returning(|_, _| Box::pin(async { Ok(()) }));
2173
2174            // Verify that the status check is scheduled with the initial delay
2175            job_producer
2176                .expect_produce_check_transaction_status_job()
2177                .withf(|_, delay| {
2178                    // Should have a delay timestamp
2179                    if let Some(scheduled_at) = delay {
2180                        // The scheduled time should be approximately STELLAR_STATUS_CHECK_INITIAL_DELAY_SECONDS from now
2181                        let now = Utc::now().timestamp();
2182                        let diff = scheduled_at - now;
2183                        // Allow some tolerance (within 2 seconds)
2184                        ((STELLAR_STATUS_CHECK_INITIAL_DELAY_SECONDS - 2)
2185                            ..=(STELLAR_STATUS_CHECK_INITIAL_DELAY_SECONDS + 2))
2186                            .contains(&diff)
2187                    } else {
2188                        false
2189                    }
2190                })
2191                .times(1)
2192                .returning(|_, _| Box::pin(async { Ok(()) }));
2193
2194            let relayer_repo = Arc::new(MockRelayerRepository::new());
2195            let counter = MockTransactionCounterServiceTrait::new();
2196
2197            let relayer = StellarRelayer::new(
2198                relayer_model,
2199                signer,
2200                provider,
2201                StellarRelayerDependencies::new(
2202                    relayer_repo,
2203                    ctx.network_repository.clone(),
2204                    Arc::new(tx_repo),
2205                    Arc::new(counter),
2206                    Arc::new(job_producer),
2207                ),
2208                dex_service,
2209            )
2210            .await
2211            .unwrap();
2212
2213            let result = relayer.process_transaction_request(tx_request).await;
2214            assert!(result.is_ok());
2215        }
2216
2217        #[tokio::test]
2218        async fn test_process_transaction_request_repository_failure() {
2219            let ctx = TestCtx::default();
2220            ctx.setup_network().await;
2221            let relayer_model = ctx.relayer_model.clone();
2222
2223            let provider = MockStellarProviderTrait::new();
2224            let signer = Arc::new(MockStellarSignTrait::new());
2225            let dex_service = create_mock_dex_service();
2226
2227            let tx_request = create_test_transaction_request();
2228
2229            // Mock repository failure
2230            let mut tx_repo = MockTransactionRepository::new();
2231            tx_repo.expect_create().returning(|_| {
2232                Err(RepositoryError::TransactionFailure(
2233                    "Database connection failed".to_string(),
2234                ))
2235            });
2236
2237            // Job producer should NOT be called when repository fails
2238            let job_producer = MockJobProducerTrait::new();
2239
2240            let relayer_repo = Arc::new(MockRelayerRepository::new());
2241            let counter = MockTransactionCounterServiceTrait::new();
2242
2243            let relayer = StellarRelayer::new(
2244                relayer_model,
2245                signer,
2246                provider,
2247                StellarRelayerDependencies::new(
2248                    relayer_repo,
2249                    ctx.network_repository.clone(),
2250                    Arc::new(tx_repo),
2251                    Arc::new(counter),
2252                    Arc::new(job_producer),
2253                ),
2254                dex_service,
2255            )
2256            .await
2257            .unwrap();
2258
2259            let result = relayer.process_transaction_request(tx_request).await;
2260            assert!(result.is_err());
2261            // RepositoryError is converted to RelayerError::NetworkConfiguration
2262            let err_msg = result.err().unwrap().to_string();
2263            assert!(
2264                err_msg.contains("Database connection failed"),
2265                "Error was: {err_msg}"
2266            );
2267        }
2268
2269        #[tokio::test]
2270        async fn test_process_transaction_request_job_producer_request_failure() {
2271            let ctx = TestCtx::default();
2272            ctx.setup_network().await;
2273            let relayer_model = ctx.relayer_model.clone();
2274
2275            let provider = MockStellarProviderTrait::new();
2276            let signer = Arc::new(MockStellarSignTrait::new());
2277            let dex_service = create_mock_dex_service();
2278
2279            let tx_request = create_test_transaction_request();
2280
2281            let mut tx_repo = MockTransactionRepository::new();
2282            tx_repo.expect_create().returning(|t| Ok(t.clone()));
2283
2284            let mut job_producer = MockJobProducerTrait::new();
2285
2286            // Status check is called FIRST and succeeds (safety net)
2287            job_producer
2288                .expect_produce_check_transaction_status_job()
2289                .returning(|_, _| Box::pin(async { Ok(()) }));
2290
2291            // Transaction request fails AFTER status check succeeds
2292            // This is safe because status check will monitor the stuck transaction
2293            job_producer
2294                .expect_produce_transaction_request_job()
2295                .returning(|_, _| {
2296                    Box::pin(async {
2297                        Err(crate::jobs::JobProducerError::QueueError(
2298                            "Queue is full".to_string(),
2299                        ))
2300                    })
2301                });
2302
2303            let relayer_repo = Arc::new(MockRelayerRepository::new());
2304            let counter = MockTransactionCounterServiceTrait::new();
2305
2306            let relayer = StellarRelayer::new(
2307                relayer_model,
2308                signer,
2309                provider,
2310                StellarRelayerDependencies::new(
2311                    relayer_repo,
2312                    ctx.network_repository.clone(),
2313                    Arc::new(tx_repo),
2314                    Arc::new(counter),
2315                    Arc::new(job_producer),
2316                ),
2317                dex_service,
2318            )
2319            .await
2320            .unwrap();
2321
2322            let result = relayer.process_transaction_request(tx_request).await;
2323            assert!(result.is_err());
2324        }
2325
2326        #[tokio::test]
2327        async fn test_process_transaction_request_job_producer_status_check_failure() {
2328            let ctx = TestCtx::default();
2329            ctx.setup_network().await;
2330            let relayer_model = ctx.relayer_model.clone();
2331
2332            let provider = MockStellarProviderTrait::new();
2333            let signer = Arc::new(MockStellarSignTrait::new());
2334            let dex_service = create_mock_dex_service();
2335
2336            let tx_request = create_test_transaction_request();
2337
2338            let mut tx_repo = MockTransactionRepository::new();
2339            tx_repo.expect_create().returning(|t| Ok(t.clone()));
2340            // When status check fails, transaction is marked as failed
2341            tx_repo
2342                .expect_partial_update()
2343                .returning(|_, _| Ok(TransactionRepoModel::default()));
2344
2345            let mut job_producer = MockJobProducerTrait::new();
2346
2347            // Status check is called FIRST and fails
2348            // This prevents orphaned transactions without monitoring
2349            job_producer
2350                .expect_produce_check_transaction_status_job()
2351                .returning(|_, _| {
2352                    Box::pin(async {
2353                        Err(crate::jobs::JobProducerError::QueueError(
2354                            "Failed to queue job".to_string(),
2355                        ))
2356                    })
2357                });
2358
2359            // Transaction request should NOT be called when status check fails
2360            // (no expectation set = test fails if called)
2361
2362            let relayer_repo = Arc::new(MockRelayerRepository::new());
2363            let counter = MockTransactionCounterServiceTrait::new();
2364
2365            let relayer = StellarRelayer::new(
2366                relayer_model,
2367                signer,
2368                provider,
2369                StellarRelayerDependencies::new(
2370                    relayer_repo,
2371                    ctx.network_repository.clone(),
2372                    Arc::new(tx_repo),
2373                    Arc::new(counter),
2374                    Arc::new(job_producer),
2375                ),
2376                dex_service,
2377            )
2378            .await
2379            .unwrap();
2380
2381            let result = relayer.process_transaction_request(tx_request).await;
2382            assert!(result.is_err());
2383        }
2384
2385        #[tokio::test]
2386        async fn test_process_transaction_request_status_check_failure_marks_tx_failed() {
2387            // Verify that when status check queue fails, the transaction is marked
2388            // as Failed with "Queue unavailable" reason
2389            let ctx = TestCtx::default();
2390            ctx.setup_network().await;
2391            let relayer_model = ctx.relayer_model.clone();
2392
2393            let provider = MockStellarProviderTrait::new();
2394            let signer = Arc::new(MockStellarSignTrait::new());
2395            let dex_service = create_mock_dex_service();
2396
2397            let tx_request = create_test_transaction_request();
2398
2399            let mut tx_repo = MockTransactionRepository::new();
2400            tx_repo.expect_create().returning(|t| Ok(t.clone()));
2401
2402            // Verify partial_update is called with correct status and reason
2403            tx_repo
2404                .expect_partial_update()
2405                .withf(|_tx_id, update| {
2406                    update.status == Some(TransactionStatus::Failed)
2407                        && update.status_reason == Some("Queue unavailable".to_string())
2408                })
2409                .returning(|_, _| Ok(TransactionRepoModel::default()));
2410
2411            let mut job_producer = MockJobProducerTrait::new();
2412            job_producer
2413                .expect_produce_check_transaction_status_job()
2414                .returning(|_, _| {
2415                    Box::pin(async {
2416                        Err(crate::jobs::JobProducerError::QueueError(
2417                            "Redis timeout".to_string(),
2418                        ))
2419                    })
2420                });
2421
2422            let relayer_repo = Arc::new(MockRelayerRepository::new());
2423            let counter = MockTransactionCounterServiceTrait::new();
2424
2425            let relayer = StellarRelayer::new(
2426                relayer_model,
2427                signer,
2428                provider,
2429                StellarRelayerDependencies::new(
2430                    relayer_repo,
2431                    ctx.network_repository.clone(),
2432                    Arc::new(tx_repo),
2433                    Arc::new(counter),
2434                    Arc::new(job_producer),
2435                ),
2436                dex_service,
2437            )
2438            .await
2439            .unwrap();
2440
2441            let result = relayer.process_transaction_request(tx_request).await;
2442            assert!(result.is_err());
2443            // The mock verification (withf) ensures partial_update was called correctly
2444        }
2445
2446        #[tokio::test]
2447        async fn test_process_transaction_request_preserves_transaction_data() {
2448            let ctx = TestCtx::default();
2449            ctx.setup_network().await;
2450            let relayer_model = ctx.relayer_model.clone();
2451
2452            let provider = MockStellarProviderTrait::new();
2453            let signer = Arc::new(MockStellarSignTrait::new());
2454            let dex_service = create_mock_dex_service();
2455
2456            let tx_request = create_test_transaction_request();
2457
2458            let mut tx_repo = MockTransactionRepository::new();
2459            tx_repo.expect_create().returning(|t| Ok(t.clone()));
2460
2461            let mut job_producer = MockJobProducerTrait::new();
2462            job_producer
2463                .expect_produce_transaction_request_job()
2464                .returning(|_, _| Box::pin(async { Ok(()) }));
2465            job_producer
2466                .expect_produce_check_transaction_status_job()
2467                .returning(|_, _| Box::pin(async { Ok(()) }));
2468
2469            let relayer_repo = Arc::new(MockRelayerRepository::new());
2470            let counter = MockTransactionCounterServiceTrait::new();
2471
2472            let relayer = StellarRelayer::new(
2473                relayer_model.clone(),
2474                signer,
2475                provider,
2476                StellarRelayerDependencies::new(
2477                    relayer_repo,
2478                    ctx.network_repository.clone(),
2479                    Arc::new(tx_repo),
2480                    Arc::new(counter),
2481                    Arc::new(job_producer),
2482                ),
2483                dex_service,
2484            )
2485            .await
2486            .unwrap();
2487
2488            let result = relayer.process_transaction_request(tx_request).await;
2489            assert!(result.is_ok());
2490
2491            let returned_tx = result.unwrap();
2492            assert_eq!(returned_tx.relayer_id, relayer_model.id);
2493            assert_eq!(returned_tx.network_type, NetworkType::Stellar);
2494            assert_eq!(returned_tx.status, TransactionStatus::Pending);
2495        }
2496    }
2497
2498    // Tests for populate_allowed_tokens_metadata
2499    mod populate_allowed_tokens_metadata_tests {
2500        use super::*;
2501        use crate::models::StellarTokenKind;
2502
2503        #[tokio::test]
2504        async fn test_populate_allowed_tokens_metadata_no_tokens() {
2505            let ctx = TestCtx::default();
2506            ctx.setup_network().await;
2507            let relayer_model = ctx.relayer_model.clone();
2508
2509            let provider = MockStellarProviderTrait::new();
2510            let signer = Arc::new(MockStellarSignTrait::new());
2511            let dex_service = create_mock_dex_service();
2512
2513            let mut relayer_repo = MockRelayerRepository::new();
2514            // Should not be called since no tokens
2515            relayer_repo.expect_update_policy().times(0);
2516
2517            let tx_repo = MockTransactionRepository::new();
2518            let job_producer = MockJobProducerTrait::new();
2519            let counter = MockTransactionCounterServiceTrait::new();
2520
2521            let relayer = StellarRelayer::new(
2522                relayer_model.clone(),
2523                signer,
2524                provider,
2525                StellarRelayerDependencies::new(
2526                    Arc::new(relayer_repo),
2527                    ctx.network_repository.clone(),
2528                    Arc::new(tx_repo),
2529                    Arc::new(counter),
2530                    Arc::new(job_producer),
2531                ),
2532                dex_service,
2533            )
2534            .await
2535            .unwrap();
2536
2537            let result = relayer.populate_allowed_tokens_metadata().await;
2538            assert!(result.is_ok());
2539        }
2540
2541        #[tokio::test]
2542        async fn test_populate_allowed_tokens_metadata_empty_tokens() {
2543            let ctx = TestCtx::default();
2544            ctx.setup_network().await;
2545            let mut relayer_model = ctx.relayer_model.clone();
2546
2547            // Set up empty allowed tokens
2548            let mut policy = RelayerStellarPolicy::default();
2549            policy.allowed_tokens = Some(vec![]);
2550            relayer_model.policies = RelayerNetworkPolicy::Stellar(policy);
2551
2552            let provider = MockStellarProviderTrait::new();
2553            let signer = Arc::new(MockStellarSignTrait::new());
2554            let dex_service = create_mock_dex_service();
2555
2556            let mut relayer_repo = MockRelayerRepository::new();
2557            // Should not be called since tokens list is empty
2558            relayer_repo.expect_update_policy().times(0);
2559
2560            let tx_repo = MockTransactionRepository::new();
2561            let job_producer = MockJobProducerTrait::new();
2562            let counter = MockTransactionCounterServiceTrait::new();
2563
2564            let relayer = StellarRelayer::new(
2565                relayer_model.clone(),
2566                signer,
2567                provider,
2568                StellarRelayerDependencies::new(
2569                    Arc::new(relayer_repo),
2570                    ctx.network_repository.clone(),
2571                    Arc::new(tx_repo),
2572                    Arc::new(counter),
2573                    Arc::new(job_producer),
2574                ),
2575                dex_service,
2576            )
2577            .await
2578            .unwrap();
2579
2580            let result = relayer.populate_allowed_tokens_metadata().await;
2581            assert!(result.is_ok());
2582        }
2583
2584        #[tokio::test]
2585        async fn test_populate_allowed_tokens_metadata_classic_asset_success() {
2586            let ctx = TestCtx::default();
2587            ctx.setup_network().await;
2588            let mut relayer_model = ctx.relayer_model.clone();
2589
2590            // Set up allowed tokens with a classic asset (USDC)
2591            let mut policy = RelayerStellarPolicy::default();
2592            policy.allowed_tokens = Some(vec![crate::models::StellarAllowedTokensPolicy {
2593                asset: "USDC:GBBD47IF6LWK7P7MDEVSCWR7DPUWV3NY3DTQEVFL4NAT4AQH3ZLLFLA5".to_string(),
2594                metadata: None,
2595                max_allowed_fee: None,
2596                swap_config: None,
2597            }]);
2598            relayer_model.policies = RelayerNetworkPolicy::Stellar(policy);
2599
2600            let provider = MockStellarProviderTrait::new();
2601            let signer = Arc::new(MockStellarSignTrait::new());
2602            let dex_service = create_mock_dex_service();
2603
2604            let mut relayer_repo = MockRelayerRepository::new();
2605            relayer_repo
2606                .expect_update_policy()
2607                .times(1)
2608                .returning(|_, _| Ok(RelayerRepoModel::default()));
2609
2610            let tx_repo = MockTransactionRepository::new();
2611            let job_producer = MockJobProducerTrait::new();
2612            let counter = MockTransactionCounterServiceTrait::new();
2613
2614            let relayer = StellarRelayer::new(
2615                relayer_model.clone(),
2616                signer,
2617                provider,
2618                StellarRelayerDependencies::new(
2619                    Arc::new(relayer_repo),
2620                    ctx.network_repository.clone(),
2621                    Arc::new(tx_repo),
2622                    Arc::new(counter),
2623                    Arc::new(job_producer),
2624                ),
2625                dex_service,
2626            )
2627            .await
2628            .unwrap();
2629
2630            let result = relayer.populate_allowed_tokens_metadata().await;
2631            assert!(result.is_ok());
2632
2633            let updated_policy = result.unwrap();
2634            assert!(updated_policy.allowed_tokens.is_some());
2635
2636            let tokens = updated_policy.allowed_tokens.unwrap();
2637            assert_eq!(tokens.len(), 1);
2638
2639            // Verify metadata was populated
2640            let token = &tokens[0];
2641            assert!(token.metadata.is_some());
2642
2643            let metadata = token.metadata.as_ref().unwrap();
2644            assert_eq!(metadata.decimals, 7); // Default Stellar decimals
2645            assert_eq!(
2646                metadata.canonical_asset_id,
2647                "USDC:GBBD47IF6LWK7P7MDEVSCWR7DPUWV3NY3DTQEVFL4NAT4AQH3ZLLFLA5"
2648            );
2649
2650            // Verify it's a classic asset
2651            match &metadata.kind {
2652                StellarTokenKind::Classic { code, issuer } => {
2653                    assert_eq!(code, "USDC");
2654                    assert_eq!(
2655                        issuer,
2656                        "GBBD47IF6LWK7P7MDEVSCWR7DPUWV3NY3DTQEVFL4NAT4AQH3ZLLFLA5"
2657                    );
2658                }
2659                _ => panic!("Expected Classic token kind"),
2660            }
2661        }
2662
2663        #[tokio::test]
2664        async fn test_populate_allowed_tokens_metadata_multiple_tokens() {
2665            let ctx = TestCtx::default();
2666            ctx.setup_network().await;
2667            let mut relayer_model = ctx.relayer_model.clone();
2668
2669            // Set up multiple allowed tokens
2670            let mut policy = RelayerStellarPolicy::default();
2671            policy.allowed_tokens = Some(vec![
2672                crate::models::StellarAllowedTokensPolicy {
2673                    asset: "USDC:GBBD47IF6LWK7P7MDEVSCWR7DPUWV3NY3DTQEVFL4NAT4AQH3ZLLFLA5"
2674                        .to_string(),
2675                    metadata: None,
2676                    max_allowed_fee: None,
2677                    swap_config: None,
2678                },
2679                crate::models::StellarAllowedTokensPolicy {
2680                    asset: "AQUA:GAHPYWLK6YRN7CVYZOO4H3VDRZ7PVF5UJGLZCSPAEIKJE2XSWF5LAGER"
2681                        .to_string(),
2682                    metadata: None,
2683                    max_allowed_fee: Some(1000000),
2684                    swap_config: None,
2685                },
2686            ]);
2687            relayer_model.policies = RelayerNetworkPolicy::Stellar(policy);
2688
2689            let provider = MockStellarProviderTrait::new();
2690            let signer = Arc::new(MockStellarSignTrait::new());
2691            let dex_service = create_mock_dex_service();
2692
2693            let mut relayer_repo = MockRelayerRepository::new();
2694            relayer_repo
2695                .expect_update_policy()
2696                .times(1)
2697                .returning(|_, _| Ok(RelayerRepoModel::default()));
2698
2699            let tx_repo = MockTransactionRepository::new();
2700            let job_producer = MockJobProducerTrait::new();
2701            let counter = MockTransactionCounterServiceTrait::new();
2702
2703            let relayer = StellarRelayer::new(
2704                relayer_model.clone(),
2705                signer,
2706                provider,
2707                StellarRelayerDependencies::new(
2708                    Arc::new(relayer_repo),
2709                    ctx.network_repository.clone(),
2710                    Arc::new(tx_repo),
2711                    Arc::new(counter),
2712                    Arc::new(job_producer),
2713                ),
2714                dex_service,
2715            )
2716            .await
2717            .unwrap();
2718
2719            let result = relayer.populate_allowed_tokens_metadata().await;
2720            assert!(result.is_ok());
2721
2722            let updated_policy = result.unwrap();
2723            let tokens = updated_policy.allowed_tokens.unwrap();
2724            assert_eq!(tokens.len(), 2);
2725
2726            // Verify both tokens have metadata
2727            assert!(tokens[0].metadata.is_some());
2728            assert!(tokens[1].metadata.is_some());
2729
2730            // Verify first token (USDC)
2731            let usdc_metadata = tokens[0].metadata.as_ref().unwrap();
2732            match &usdc_metadata.kind {
2733                StellarTokenKind::Classic { code, .. } => {
2734                    assert_eq!(code, "USDC");
2735                }
2736                _ => panic!("Expected Classic token kind for USDC"),
2737            }
2738
2739            // Verify second token (AQUA)
2740            let aqua_metadata = tokens[1].metadata.as_ref().unwrap();
2741            match &aqua_metadata.kind {
2742                StellarTokenKind::Classic { code, .. } => {
2743                    assert_eq!(code, "AQUA");
2744                }
2745                _ => panic!("Expected Classic token kind for AQUA"),
2746            }
2747
2748            // Verify max_allowed_fee is preserved
2749            assert_eq!(tokens[1].max_allowed_fee, Some(1000000));
2750        }
2751
2752        #[tokio::test]
2753        async fn test_populate_allowed_tokens_metadata_invalid_asset() {
2754            let ctx = TestCtx::default();
2755            ctx.setup_network().await;
2756            let mut relayer_model = ctx.relayer_model.clone();
2757
2758            // Set up allowed tokens with invalid asset format
2759            let mut policy = RelayerStellarPolicy::default();
2760            policy.allowed_tokens = Some(vec![crate::models::StellarAllowedTokensPolicy {
2761                asset: "INVALID_FORMAT".to_string(), // Missing issuer
2762                metadata: None,
2763                max_allowed_fee: None,
2764                swap_config: None,
2765            }]);
2766            relayer_model.policies = RelayerNetworkPolicy::Stellar(policy);
2767
2768            let provider = MockStellarProviderTrait::new();
2769            let signer = Arc::new(MockStellarSignTrait::new());
2770            let dex_service = create_mock_dex_service();
2771
2772            let relayer_repo = MockRelayerRepository::new();
2773            let tx_repo = MockTransactionRepository::new();
2774            let job_producer = MockJobProducerTrait::new();
2775            let counter = MockTransactionCounterServiceTrait::new();
2776
2777            let relayer = StellarRelayer::new(
2778                relayer_model.clone(),
2779                signer,
2780                provider,
2781                StellarRelayerDependencies::new(
2782                    Arc::new(relayer_repo),
2783                    ctx.network_repository.clone(),
2784                    Arc::new(tx_repo),
2785                    Arc::new(counter),
2786                    Arc::new(job_producer),
2787                ),
2788                dex_service,
2789            )
2790            .await
2791            .unwrap();
2792
2793            let result = relayer.populate_allowed_tokens_metadata().await;
2794            assert!(result.is_err());
2795        }
2796    }
2797
2798    // Tests for migrate_fee_payment_strategy_if_needed
2799    mod migrate_fee_payment_strategy_tests {
2800        use super::*;
2801
2802        #[tokio::test]
2803        async fn test_migrate_fee_payment_strategy_in_memory_storage() {
2804            let ctx = TestCtx::default();
2805            ctx.setup_network().await;
2806            let relayer_model = ctx.relayer_model.clone();
2807
2808            let provider = MockStellarProviderTrait::new();
2809            let signer = Arc::new(MockStellarSignTrait::new());
2810            let dex_service = create_mock_dex_service();
2811
2812            let mut relayer_repo = MockRelayerRepository::new();
2813            // Mock in-memory storage
2814            relayer_repo
2815                .expect_is_persistent_storage()
2816                .returning(|| false);
2817            // Should not call update_policy for in-memory storage
2818            relayer_repo.expect_update_policy().times(0);
2819
2820            let tx_repo = MockTransactionRepository::new();
2821            let job_producer = MockJobProducerTrait::new();
2822            let counter = MockTransactionCounterServiceTrait::new();
2823
2824            let relayer = StellarRelayer::new(
2825                relayer_model.clone(),
2826                signer,
2827                provider,
2828                StellarRelayerDependencies::new(
2829                    Arc::new(relayer_repo),
2830                    ctx.network_repository.clone(),
2831                    Arc::new(tx_repo),
2832                    Arc::new(counter),
2833                    Arc::new(job_producer),
2834                ),
2835                dex_service,
2836            )
2837            .await
2838            .unwrap();
2839
2840            let result = relayer.migrate_fee_payment_strategy_if_needed().await;
2841            assert!(result.is_ok());
2842        }
2843
2844        #[tokio::test]
2845        async fn test_migrate_fee_payment_strategy_already_set() {
2846            let ctx = TestCtx::default();
2847            ctx.setup_network().await;
2848            let mut relayer_model = ctx.relayer_model.clone();
2849
2850            // Set fee_payment_strategy
2851            let mut policy = RelayerStellarPolicy::default();
2852            policy.fee_payment_strategy = Some(StellarFeePaymentStrategy::User);
2853            relayer_model.policies = RelayerNetworkPolicy::Stellar(policy);
2854
2855            let provider = MockStellarProviderTrait::new();
2856            let signer = Arc::new(MockStellarSignTrait::new());
2857            let dex_service = create_mock_dex_service();
2858
2859            let mut relayer_repo = MockRelayerRepository::new();
2860            relayer_repo
2861                .expect_is_persistent_storage()
2862                .returning(|| true);
2863            // Should not call update_policy since already set
2864            relayer_repo.expect_update_policy().times(0);
2865
2866            let tx_repo = MockTransactionRepository::new();
2867            let job_producer = MockJobProducerTrait::new();
2868            let counter = MockTransactionCounterServiceTrait::new();
2869
2870            let relayer = StellarRelayer::new(
2871                relayer_model.clone(),
2872                signer,
2873                provider,
2874                StellarRelayerDependencies::new(
2875                    Arc::new(relayer_repo),
2876                    ctx.network_repository.clone(),
2877                    Arc::new(tx_repo),
2878                    Arc::new(counter),
2879                    Arc::new(job_producer),
2880                ),
2881                dex_service,
2882            )
2883            .await
2884            .unwrap();
2885
2886            let result = relayer.migrate_fee_payment_strategy_if_needed().await;
2887            assert!(result.is_ok());
2888        }
2889
2890        #[tokio::test]
2891        async fn test_migrate_fee_payment_strategy_migration_needed() {
2892            let ctx = TestCtx::default();
2893            ctx.setup_network().await;
2894            let relayer_model = ctx.relayer_model.clone();
2895
2896            let provider = MockStellarProviderTrait::new();
2897            let signer = Arc::new(MockStellarSignTrait::new());
2898            let dex_service = create_mock_dex_service();
2899
2900            let mut relayer_repo = MockRelayerRepository::new();
2901            relayer_repo
2902                .expect_is_persistent_storage()
2903                .returning(|| true);
2904            relayer_repo
2905                .expect_update_policy()
2906                .times(1)
2907                .returning(|_, policy| {
2908                    // Verify the policy is set to Relayer
2909                    if let RelayerNetworkPolicy::Stellar(stellar_policy) = &policy {
2910                        assert_eq!(
2911                            stellar_policy.fee_payment_strategy,
2912                            Some(StellarFeePaymentStrategy::Relayer)
2913                        );
2914                    }
2915                    Ok(RelayerRepoModel::default())
2916                });
2917
2918            let tx_repo = MockTransactionRepository::new();
2919            let job_producer = MockJobProducerTrait::new();
2920            let counter = MockTransactionCounterServiceTrait::new();
2921
2922            let relayer = StellarRelayer::new(
2923                relayer_model.clone(),
2924                signer,
2925                provider,
2926                StellarRelayerDependencies::new(
2927                    Arc::new(relayer_repo),
2928                    ctx.network_repository.clone(),
2929                    Arc::new(tx_repo),
2930                    Arc::new(counter),
2931                    Arc::new(job_producer),
2932                ),
2933                dex_service,
2934            )
2935            .await
2936            .unwrap();
2937
2938            let result = relayer.migrate_fee_payment_strategy_if_needed().await;
2939            assert!(result.is_ok());
2940        }
2941
2942        #[tokio::test]
2943        async fn test_migrate_fee_payment_strategy_update_fails() {
2944            let ctx = TestCtx::default();
2945            ctx.setup_network().await;
2946            let relayer_model = ctx.relayer_model.clone();
2947
2948            let provider = MockStellarProviderTrait::new();
2949            let signer = Arc::new(MockStellarSignTrait::new());
2950            let dex_service = create_mock_dex_service();
2951
2952            let mut relayer_repo = MockRelayerRepository::new();
2953            relayer_repo
2954                .expect_is_persistent_storage()
2955                .returning(|| true);
2956            relayer_repo
2957                .expect_update_policy()
2958                .times(1)
2959                .returning(|_, _| {
2960                    Err(RepositoryError::TransactionFailure(
2961                        "Database error".to_string(),
2962                    ))
2963                });
2964
2965            let tx_repo = MockTransactionRepository::new();
2966            let job_producer = MockJobProducerTrait::new();
2967            let counter = MockTransactionCounterServiceTrait::new();
2968
2969            let relayer = StellarRelayer::new(
2970                relayer_model.clone(),
2971                signer,
2972                provider,
2973                StellarRelayerDependencies::new(
2974                    Arc::new(relayer_repo),
2975                    ctx.network_repository.clone(),
2976                    Arc::new(tx_repo),
2977                    Arc::new(counter),
2978                    Arc::new(job_producer),
2979                ),
2980                dex_service,
2981            )
2982            .await
2983            .unwrap();
2984
2985            let result = relayer.migrate_fee_payment_strategy_if_needed().await;
2986            assert!(result.is_err());
2987            assert!(matches!(
2988                result.unwrap_err(),
2989                RelayerError::PolicyConfigurationError(_)
2990            ));
2991        }
2992    }
2993
2994    // Tests for check_balance_and_trigger_token_swap_if_needed
2995    mod check_balance_and_trigger_token_swap_tests {
2996        use super::*;
2997        use crate::models::RelayerStellarSwapConfig;
2998
2999        #[tokio::test]
3000        async fn test_check_balance_no_swap_config() {
3001            let ctx = TestCtx::default();
3002            ctx.setup_network().await;
3003            let relayer_model = ctx.relayer_model.clone();
3004
3005            let provider = MockStellarProviderTrait::new();
3006            let signer = Arc::new(MockStellarSignTrait::new());
3007            let dex_service = create_mock_dex_service();
3008
3009            let relayer_repo = MockRelayerRepository::new();
3010            let tx_repo = MockTransactionRepository::new();
3011            let job_producer = MockJobProducerTrait::new();
3012            let counter = MockTransactionCounterServiceTrait::new();
3013
3014            let relayer = StellarRelayer::new(
3015                relayer_model.clone(),
3016                signer,
3017                provider,
3018                StellarRelayerDependencies::new(
3019                    Arc::new(relayer_repo),
3020                    ctx.network_repository.clone(),
3021                    Arc::new(tx_repo),
3022                    Arc::new(counter),
3023                    Arc::new(job_producer),
3024                ),
3025                dex_service,
3026            )
3027            .await
3028            .unwrap();
3029
3030            let result = relayer
3031                .check_balance_and_trigger_token_swap_if_needed()
3032                .await;
3033            assert!(result.is_ok());
3034        }
3035
3036        #[tokio::test]
3037        async fn test_check_balance_no_threshold() {
3038            let ctx = TestCtx::default();
3039            ctx.setup_network().await;
3040            let mut relayer_model = ctx.relayer_model.clone();
3041
3042            // Set up swap config without threshold
3043            let mut policy = RelayerStellarPolicy::default();
3044            policy.swap_config = Some(RelayerStellarSwapConfig {
3045                strategies: vec![],
3046                min_balance_threshold: None,
3047                cron_schedule: None,
3048            });
3049            relayer_model.policies = RelayerNetworkPolicy::Stellar(policy);
3050
3051            let provider = MockStellarProviderTrait::new();
3052            let signer = Arc::new(MockStellarSignTrait::new());
3053            let dex_service = create_mock_dex_service();
3054
3055            let relayer_repo = MockRelayerRepository::new();
3056            let tx_repo = MockTransactionRepository::new();
3057            let job_producer = MockJobProducerTrait::new();
3058            let counter = MockTransactionCounterServiceTrait::new();
3059
3060            let relayer = StellarRelayer::new(
3061                relayer_model.clone(),
3062                signer,
3063                provider,
3064                StellarRelayerDependencies::new(
3065                    Arc::new(relayer_repo),
3066                    ctx.network_repository.clone(),
3067                    Arc::new(tx_repo),
3068                    Arc::new(counter),
3069                    Arc::new(job_producer),
3070                ),
3071                dex_service,
3072            )
3073            .await
3074            .unwrap();
3075
3076            let result = relayer
3077                .check_balance_and_trigger_token_swap_if_needed()
3078                .await;
3079            assert!(result.is_ok());
3080        }
3081
3082        #[tokio::test]
3083        async fn test_check_balance_above_threshold() {
3084            let ctx = TestCtx::default();
3085            ctx.setup_network().await;
3086            let mut relayer_model = ctx.relayer_model.clone();
3087
3088            // Set up swap config with threshold
3089            let mut policy = RelayerStellarPolicy::default();
3090            policy.swap_config = Some(RelayerStellarSwapConfig {
3091                strategies: vec![],
3092                min_balance_threshold: Some(1000000), // 1 XLM
3093                cron_schedule: None,
3094            });
3095            relayer_model.policies = RelayerNetworkPolicy::Stellar(policy);
3096
3097            let mut provider = MockStellarProviderTrait::new();
3098            // Mock get_account to return balance above threshold
3099            provider.expect_get_account().returning(|_| {
3100                Box::pin(async {
3101                    Ok(AccountEntry {
3102                        account_id: AccountId(PublicKey::PublicKeyTypeEd25519(Uint256([0; 32]))),
3103                        balance: 10000000, // 10 XLM (above threshold)
3104                        ext: AccountEntryExt::V0,
3105                        flags: 0,
3106                        home_domain: String32::default(),
3107                        inflation_dest: None,
3108                        seq_num: SequenceNumber(5),
3109                        num_sub_entries: 0,
3110                        signers: VecM::default(),
3111                        thresholds: Thresholds([0, 0, 0, 0]),
3112                    })
3113                })
3114            });
3115
3116            let signer = Arc::new(MockStellarSignTrait::new());
3117            let dex_service = create_mock_dex_service();
3118
3119            let relayer_repo = MockRelayerRepository::new();
3120            let tx_repo = MockTransactionRepository::new();
3121            let job_producer = MockJobProducerTrait::new();
3122            let counter = MockTransactionCounterServiceTrait::new();
3123
3124            let relayer = StellarRelayer::new(
3125                relayer_model.clone(),
3126                signer,
3127                provider,
3128                StellarRelayerDependencies::new(
3129                    Arc::new(relayer_repo),
3130                    ctx.network_repository.clone(),
3131                    Arc::new(tx_repo),
3132                    Arc::new(counter),
3133                    Arc::new(job_producer),
3134                ),
3135                dex_service,
3136            )
3137            .await
3138            .unwrap();
3139
3140            let result = relayer
3141                .check_balance_and_trigger_token_swap_if_needed()
3142                .await;
3143            assert!(result.is_ok());
3144        }
3145
3146        #[tokio::test]
3147        async fn test_check_balance_provider_error() {
3148            let ctx = TestCtx::default();
3149            ctx.setup_network().await;
3150            let mut relayer_model = ctx.relayer_model.clone();
3151
3152            // Set up swap config with threshold
3153            let mut policy = RelayerStellarPolicy::default();
3154            policy.swap_config = Some(RelayerStellarSwapConfig {
3155                strategies: vec![],
3156                min_balance_threshold: Some(1000000),
3157                cron_schedule: None,
3158            });
3159            relayer_model.policies = RelayerNetworkPolicy::Stellar(policy);
3160
3161            let mut provider = MockStellarProviderTrait::new();
3162            provider.expect_get_account().returning(|_| {
3163                Box::pin(async { Err(ProviderError::Other("Network error".to_string())) })
3164            });
3165
3166            let signer = Arc::new(MockStellarSignTrait::new());
3167            let dex_service = create_mock_dex_service();
3168
3169            let relayer_repo = MockRelayerRepository::new();
3170            let tx_repo = MockTransactionRepository::new();
3171            let job_producer = MockJobProducerTrait::new();
3172            let counter = MockTransactionCounterServiceTrait::new();
3173
3174            let relayer = StellarRelayer::new(
3175                relayer_model.clone(),
3176                signer,
3177                provider,
3178                StellarRelayerDependencies::new(
3179                    Arc::new(relayer_repo),
3180                    ctx.network_repository.clone(),
3181                    Arc::new(tx_repo),
3182                    Arc::new(counter),
3183                    Arc::new(job_producer),
3184                ),
3185                dex_service,
3186            )
3187            .await
3188            .unwrap();
3189
3190            let result = relayer
3191                .check_balance_and_trigger_token_swap_if_needed()
3192                .await;
3193            assert!(result.is_err());
3194        }
3195    }
3196
3197    // Tests for check_health
3198    mod check_health_tests {
3199        use super::*;
3200        use crate::models::RelayerStellarSwapConfig;
3201
3202        #[tokio::test]
3203        async fn test_check_health_success() {
3204            let ctx = TestCtx::default();
3205            ctx.setup_network().await;
3206            let relayer_model = ctx.relayer_model.clone();
3207
3208            let mut provider = MockStellarProviderTrait::new();
3209            provider.expect_get_account().returning(|_| {
3210                Box::pin(async {
3211                    Ok(AccountEntry {
3212                        account_id: AccountId(PublicKey::PublicKeyTypeEd25519(Uint256([0; 32]))),
3213                        balance: 10000000,
3214                        ext: AccountEntryExt::V0,
3215                        flags: 0,
3216                        home_domain: String32::default(),
3217                        inflation_dest: None,
3218                        seq_num: SequenceNumber(5),
3219                        num_sub_entries: 0,
3220                        signers: VecM::default(),
3221                        thresholds: Thresholds([0, 0, 0, 0]),
3222                    })
3223                })
3224            });
3225
3226            let signer = Arc::new(MockStellarSignTrait::new());
3227            let dex_service = create_mock_dex_service();
3228
3229            let relayer_repo = MockRelayerRepository::new();
3230            let tx_repo = MockTransactionRepository::new();
3231            let job_producer = MockJobProducerTrait::new();
3232
3233            let mut counter = MockTransactionCounterServiceTrait::new();
3234            counter
3235                .expect_sync_floor()
3236                .returning(|floor| Box::pin(async move { Ok(floor) }));
3237
3238            let relayer = StellarRelayer::new(
3239                relayer_model.clone(),
3240                signer,
3241                provider,
3242                StellarRelayerDependencies::new(
3243                    Arc::new(relayer_repo),
3244                    ctx.network_repository.clone(),
3245                    Arc::new(tx_repo),
3246                    Arc::new(counter),
3247                    Arc::new(job_producer),
3248                ),
3249                dex_service,
3250            )
3251            .await
3252            .unwrap();
3253
3254            let result = relayer.check_health().await;
3255            assert!(result.is_ok());
3256        }
3257
3258        #[tokio::test]
3259        async fn test_check_health_sequence_sync_fails() {
3260            let ctx = TestCtx::default();
3261            ctx.setup_network().await;
3262            let relayer_model = ctx.relayer_model.clone();
3263
3264            let mut provider = MockStellarProviderTrait::new();
3265            provider.expect_get_account().returning(|_| {
3266                Box::pin(async { Err(ProviderError::Other("Network error".to_string())) })
3267            });
3268
3269            let signer = Arc::new(MockStellarSignTrait::new());
3270            let dex_service = create_mock_dex_service();
3271
3272            let relayer_repo = MockRelayerRepository::new();
3273            let tx_repo = MockTransactionRepository::new();
3274            let job_producer = MockJobProducerTrait::new();
3275            let counter = MockTransactionCounterServiceTrait::new();
3276
3277            let relayer = StellarRelayer::new(
3278                relayer_model.clone(),
3279                signer,
3280                provider,
3281                StellarRelayerDependencies::new(
3282                    Arc::new(relayer_repo),
3283                    ctx.network_repository.clone(),
3284                    Arc::new(tx_repo),
3285                    Arc::new(counter),
3286                    Arc::new(job_producer),
3287                ),
3288                dex_service,
3289            )
3290            .await
3291            .unwrap();
3292
3293            let result = relayer.check_health().await;
3294            assert!(result.is_err());
3295            let failures = result.unwrap_err();
3296            assert_eq!(failures.len(), 1);
3297            assert!(matches!(
3298                failures[0],
3299                HealthCheckFailure::SequenceSyncFailed(_)
3300            ));
3301        }
3302
3303        #[tokio::test]
3304        async fn test_check_health_with_user_fee_strategy() {
3305            let ctx = TestCtx::default();
3306            ctx.setup_network().await;
3307            let mut relayer_model = ctx.relayer_model.clone();
3308
3309            // Set up user fee payment strategy
3310            let mut policy = RelayerStellarPolicy::default();
3311            policy.fee_payment_strategy = Some(StellarFeePaymentStrategy::User);
3312            policy.swap_config = Some(RelayerStellarSwapConfig {
3313                strategies: vec![],
3314                min_balance_threshold: Some(1000000),
3315                cron_schedule: None,
3316            });
3317            relayer_model.policies = RelayerNetworkPolicy::Stellar(policy);
3318
3319            let mut provider = MockStellarProviderTrait::new();
3320            provider.expect_get_account().returning(|_| {
3321                Box::pin(async {
3322                    Ok(AccountEntry {
3323                        account_id: AccountId(PublicKey::PublicKeyTypeEd25519(Uint256([0; 32]))),
3324                        balance: 10000000, // Above threshold
3325                        ext: AccountEntryExt::V0,
3326                        flags: 0,
3327                        home_domain: String32::default(),
3328                        inflation_dest: None,
3329                        seq_num: SequenceNumber(5),
3330                        num_sub_entries: 0,
3331                        signers: VecM::default(),
3332                        thresholds: Thresholds([0, 0, 0, 0]),
3333                    })
3334                })
3335            });
3336
3337            let signer = Arc::new(MockStellarSignTrait::new());
3338            let dex_service = create_mock_dex_service();
3339
3340            let relayer_repo = MockRelayerRepository::new();
3341            let tx_repo = MockTransactionRepository::new();
3342            let job_producer = MockJobProducerTrait::new();
3343
3344            let mut counter = MockTransactionCounterServiceTrait::new();
3345            counter
3346                .expect_sync_floor()
3347                .returning(|floor| Box::pin(async move { Ok(floor) }));
3348
3349            let relayer = StellarRelayer::new(
3350                relayer_model.clone(),
3351                signer,
3352                provider,
3353                StellarRelayerDependencies::new(
3354                    Arc::new(relayer_repo),
3355                    ctx.network_repository.clone(),
3356                    Arc::new(tx_repo),
3357                    Arc::new(counter),
3358                    Arc::new(job_producer),
3359                ),
3360                dex_service,
3361            )
3362            .await
3363            .unwrap();
3364
3365            let result = relayer.check_health().await;
3366            // Should pass even with user fee strategy
3367            assert!(result.is_ok());
3368        }
3369    }
3370
3371    // Tests for RPC method
3372    mod rpc_tests {
3373        use super::*;
3374        use crate::models::{JsonRpcId, StellarRpcRequest};
3375
3376        #[tokio::test]
3377        async fn test_rpc_invalid_network_request() {
3378            let ctx = TestCtx::default();
3379            ctx.setup_network().await;
3380            let relayer_model = ctx.relayer_model.clone();
3381
3382            let provider = MockStellarProviderTrait::new();
3383            let signer = Arc::new(MockStellarSignTrait::new());
3384            let dex_service = create_mock_dex_service();
3385
3386            let relayer_repo = MockRelayerRepository::new();
3387            let tx_repo = MockTransactionRepository::new();
3388            let job_producer = MockJobProducerTrait::new();
3389            let counter = MockTransactionCounterServiceTrait::new();
3390
3391            let relayer = StellarRelayer::new(
3392                relayer_model.clone(),
3393                signer,
3394                provider,
3395                StellarRelayerDependencies::new(
3396                    Arc::new(relayer_repo),
3397                    ctx.network_repository.clone(),
3398                    Arc::new(tx_repo),
3399                    Arc::new(counter),
3400                    Arc::new(job_producer),
3401                ),
3402                dex_service,
3403            )
3404            .await
3405            .unwrap();
3406
3407            // Create a request with wrong network type
3408            let request = JsonRpcRequest {
3409                jsonrpc: "2.0".to_string(),
3410                id: Some(JsonRpcId::Number(1)),
3411                params: NetworkRpcRequest::Evm(crate::models::EvmRpcRequest::RawRpcRequest {
3412                    method: "eth_blockNumber".to_string(),
3413                    params: serde_json::Value::Null,
3414                }),
3415            };
3416
3417            let result = relayer.rpc(request).await;
3418            assert!(result.is_ok());
3419            let response = result.unwrap();
3420            // Should return an error response for invalid network type
3421            assert!(response.error.is_some());
3422        }
3423
3424        #[tokio::test]
3425        async fn test_rpc_provider_error() {
3426            let ctx = TestCtx::default();
3427            ctx.setup_network().await;
3428            let relayer_model = ctx.relayer_model.clone();
3429
3430            let mut provider = MockStellarProviderTrait::new();
3431            provider.expect_raw_request_dyn().returning(|_, _, _| {
3432                Box::pin(async { Err(ProviderError::Other("RPC error".to_string())) })
3433            });
3434
3435            let signer = Arc::new(MockStellarSignTrait::new());
3436            let dex_service = create_mock_dex_service();
3437
3438            let relayer_repo = MockRelayerRepository::new();
3439            let tx_repo = MockTransactionRepository::new();
3440            let job_producer = MockJobProducerTrait::new();
3441            let counter = MockTransactionCounterServiceTrait::new();
3442
3443            let relayer = StellarRelayer::new(
3444                relayer_model.clone(),
3445                signer,
3446                provider,
3447                StellarRelayerDependencies::new(
3448                    Arc::new(relayer_repo),
3449                    ctx.network_repository.clone(),
3450                    Arc::new(tx_repo),
3451                    Arc::new(counter),
3452                    Arc::new(job_producer),
3453                ),
3454                dex_service,
3455            )
3456            .await
3457            .unwrap();
3458
3459            let request = JsonRpcRequest {
3460                jsonrpc: "2.0".to_string(),
3461                id: Some(JsonRpcId::Number(1)),
3462                params: NetworkRpcRequest::Stellar(StellarRpcRequest::RawRpcRequest {
3463                    method: "getHealth".to_string(),
3464                    params: serde_json::Value::Null,
3465                }),
3466            };
3467
3468            let result = relayer.rpc(request).await;
3469            assert!(result.is_ok());
3470            let response = result.unwrap();
3471            // Should return an error response for provider error
3472            assert!(response.error.is_some());
3473        }
3474
3475        #[tokio::test]
3476        async fn test_rpc_success() {
3477            let ctx = TestCtx::default();
3478            ctx.setup_network().await;
3479            let relayer_model = ctx.relayer_model.clone();
3480
3481            let mut provider = MockStellarProviderTrait::new();
3482            provider.expect_raw_request_dyn().returning(|_, _, _| {
3483                Box::pin(async { Ok(serde_json::json!({"status": "healthy"})) })
3484            });
3485
3486            let signer = Arc::new(MockStellarSignTrait::new());
3487            let dex_service = create_mock_dex_service();
3488
3489            let relayer_repo = MockRelayerRepository::new();
3490            let tx_repo = MockTransactionRepository::new();
3491            let job_producer = MockJobProducerTrait::new();
3492            let counter = MockTransactionCounterServiceTrait::new();
3493
3494            let relayer = StellarRelayer::new(
3495                relayer_model.clone(),
3496                signer,
3497                provider,
3498                StellarRelayerDependencies::new(
3499                    Arc::new(relayer_repo),
3500                    ctx.network_repository.clone(),
3501                    Arc::new(tx_repo),
3502                    Arc::new(counter),
3503                    Arc::new(job_producer),
3504                ),
3505                dex_service,
3506            )
3507            .await
3508            .unwrap();
3509
3510            let request = JsonRpcRequest {
3511                jsonrpc: "2.0".to_string(),
3512                id: Some(JsonRpcId::Number(1)),
3513                params: NetworkRpcRequest::Stellar(StellarRpcRequest::RawRpcRequest {
3514                    method: "getHealth".to_string(),
3515                    params: serde_json::Value::Null,
3516                }),
3517            };
3518
3519            let result = relayer.rpc(request).await;
3520            assert!(result.is_ok());
3521            let response = result.unwrap();
3522            assert!(response.error.is_none());
3523            assert!(response.result.is_some());
3524        }
3525    }
3526}