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};
5use 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
67pub 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 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 #[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 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 #[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 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 #[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 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 policy.fee_payment_strategy.is_some() {
346 return Ok(());
347 }
348
349 info!(
351 relayer_id = %self.relayer.id,
352 "Migrating Stellar relayer: setting fee_payment_strategy to 'Relayer' (old default behavior)"
353 );
354
355 let mut updated_policy = policy;
357 updated_policy.fee_payment_strategy = Some(StellarFeePaymentStrategy::Relayer);
358
359 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 #[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 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 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 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 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 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 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 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 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 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 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 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, )
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 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 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 self.migrate_fee_payment_strategy_if_needed().await?;
793
794 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 if self.relayer.system_disabled {
806 self.relayer_repository
808 .enable_relayer(self.relayer.id.clone())
809 .await?;
810 }
811 }
812 Err(failures) => {
813 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 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 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 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 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 if user_pays_fee {
955 let envelope = parse_transaction_xdr(&stellar_req.unsigned_xdr, false)
957 .map_err(|e| RelayerError::ValidationError(format!("Failed to parse XDR: {e}")))?;
958
959 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()), )
969 .await
970 .map_err(|e| {
971 RelayerError::ValidationError(format!("Failed to validate transaction: {e}"))
972 })?;
973 }
974
975 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 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 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 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 #[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 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
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 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 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 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; 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 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 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 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; 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 provider
1753 .expect_get_account()
1754 .returning(|_| Box::pin(ready(Err(ProviderError::Other("RPC error".to_string())))));
1755
1756 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 job_producer
1769 .expect_produce_send_notification_job()
1770 .returning(|_, _| Box::pin(async { Ok(()) }));
1771
1772 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; 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 provider.expect_get_account().returning(|_| {
1818 Box::pin(ready(Ok(AccountEntry {
1819 account_id: AccountId(PublicKey::PublicKeyTypeEd25519(Uint256([0; 32]))),
1820 balance: 1000000000, 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 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; let mut provider = MockStellarProviderTrait::new();
1877
1878 provider.expect_get_account().returning(|_| {
1880 Box::pin(ready(Ok(AccountEntry {
1881 account_id: AccountId(PublicKey::PublicKeyTypeEd25519(Uint256([0; 32]))),
1882 balance: 1000000000, 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 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; 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 provider.expect_get_account().returning(|_| {
1948 Box::pin(ready(Err(ProviderError::Other(
1949 "Sequence sync failed".to_string(),
1950 ))))
1951 });
1952
1953 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 job_producer
1966 .expect_produce_send_notification_job()
1967 .returning(|_, _| Box::pin(async { Ok(()) }));
1968
1969 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; relayer_model.notification_id = None; 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 provider.expect_get_account().returning(|_| {
2015 Box::pin(ready(Err(ProviderError::Other(
2016 "Sequence sync failed".to_string(),
2017 ))))
2018 });
2019
2020 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 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 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 let tx_request = create_test_transaction_request();
2098
2099 let mut tx_repo = MockTransactionRepository::new();
2101 tx_repo.expect_create().returning(|t| Ok(t.clone()));
2102
2103 let mut job_producer = MockJobProducerTrait::new();
2105
2106 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 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 job_producer
2176 .expect_produce_check_transaction_status_job()
2177 .withf(|_, delay| {
2178 if let Some(scheduled_at) = delay {
2180 let now = Utc::now().timestamp();
2182 let diff = scheduled_at - now;
2183 ((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 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 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 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 job_producer
2288 .expect_produce_check_transaction_status_job()
2289 .returning(|_, _| Box::pin(async { Ok(()) }));
2290
2291 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 tx_repo
2342 .expect_partial_update()
2343 .returning(|_, _| Ok(TransactionRepoModel::default()));
2344
2345 let mut job_producer = MockJobProducerTrait::new();
2346
2347 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 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 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 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 }
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 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 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 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 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 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 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); assert_eq!(
2646 metadata.canonical_asset_id,
2647 "USDC:GBBD47IF6LWK7P7MDEVSCWR7DPUWV3NY3DTQEVFL4NAT4AQH3ZLLFLA5"
2648 );
2649
2650 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 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 assert!(tokens[0].metadata.is_some());
2728 assert!(tokens[1].metadata.is_some());
2729
2730 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 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 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 let mut policy = RelayerStellarPolicy::default();
2760 policy.allowed_tokens = Some(vec![crate::models::StellarAllowedTokensPolicy {
2761 asset: "INVALID_FORMAT".to_string(), 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 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 relayer_repo
2815 .expect_is_persistent_storage()
2816 .returning(|| false);
2817 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 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 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 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 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 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 let mut policy = RelayerStellarPolicy::default();
3090 policy.swap_config = Some(RelayerStellarSwapConfig {
3091 strategies: vec![],
3092 min_balance_threshold: Some(1000000), cron_schedule: None,
3094 });
3095 relayer_model.policies = RelayerNetworkPolicy::Stellar(policy);
3096
3097 let mut provider = MockStellarProviderTrait::new();
3098 provider.expect_get_account().returning(|_| {
3100 Box::pin(async {
3101 Ok(AccountEntry {
3102 account_id: AccountId(PublicKey::PublicKeyTypeEd25519(Uint256([0; 32]))),
3103 balance: 10000000, 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 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 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 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, 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 assert!(result.is_ok());
3368 }
3369 }
3370
3371 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 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 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 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}