openzeppelin_relayer/repositories/transaction/
mod.rs

1//! Transaction Repository Module
2//!
3//! This module provides the transaction repository layer for the OpenZeppelin Relayer service.
4//! It implements the Repository pattern to abstract transaction data persistence operations,
5//! supporting both in-memory and Redis-backed storage implementations.
6//!
7//! ## Features
8//!
9//! - **CRUD Operations**: Create, read, update, and delete transactions
10//! - **Specialized Queries**: Find transactions by relayer ID, status, and nonce
11//! - **Pagination Support**: Efficient paginated listing of transactions
12//! - **Status Management**: Update transaction status and timestamps
13//! - **Partial Updates**: Support for partial transaction updates
14//! - **Network Data**: Manage transaction network-specific data
15//!
16//! ## Repository Implementations
17//!
18//! - [`InMemoryTransactionRepository`]: Fast in-memory storage for testing/development
19//! - [`RedisTransactionRepository`]: Redis-backed storage for production environments
20//!
21mod transaction_in_memory;
22mod transaction_redis;
23
24pub use transaction_in_memory::*;
25pub use transaction_redis::*;
26
27use crate::{
28    models::{
29        NetworkTransactionData, TransactionRepoModel, TransactionStatus, TransactionUpdateRequest,
30    },
31    repositories::{BatchDeleteResult, TransactionDeleteRequest, *},
32    utils::RedisConnections,
33};
34use async_trait::async_trait;
35use eyre::Result;
36use std::sync::Arc;
37
38/// A trait defining transaction repository operations
39#[async_trait]
40pub trait TransactionRepository: Repository<TransactionRepoModel, String> {
41    /// Returns underlying storage Redis connections when available.
42    ///
43    /// In-memory implementations return `None`.
44    fn connection_info(&self) -> Option<(Arc<RedisConnections>, String)> {
45        None
46    }
47
48    /// Find transactions by relayer ID with pagination
49    async fn find_by_relayer_id(
50        &self,
51        relayer_id: &str,
52        query: PaginationQuery,
53    ) -> Result<PaginatedResult<TransactionRepoModel>, RepositoryError>;
54
55    /// Find transactions by relayer ID and status(es).
56    ///
57    /// Results are sorted by created_at descending (newest first).
58    async fn find_by_status(
59        &self,
60        relayer_id: &str,
61        statuses: &[TransactionStatus],
62    ) -> Result<Vec<TransactionRepoModel>, RepositoryError>;
63
64    /// Find transactions by relayer ID and status(es) with pagination.
65    ///
66    /// Results are sorted by timestamp:
67    /// - For Confirmed transactions: sorted by confirmed_at (on-chain confirmation order)
68    /// - For all other statuses: sorted by created_at (queue/processing order)
69    ///
70    /// The `oldest_first` parameter controls sort direction:
71    /// - `false` (default): newest first (descending) - for displaying recent transactions
72    /// - `true`: oldest first (ascending) - for FIFO queue processing
73    ///
74    /// For multi-status queries, transactions are merged and sorted using the same rules,
75    /// ensuring consistent ordering across different statuses.
76    async fn find_by_status_paginated(
77        &self,
78        relayer_id: &str,
79        statuses: &[TransactionStatus],
80        query: PaginationQuery,
81        oldest_first: bool,
82    ) -> Result<PaginatedResult<TransactionRepoModel>, RepositoryError>;
83
84    /// Like [`find_by_status_paginated`], but when `exclude_canceled` is true,
85    /// transactions flagged `is_canceled` are removed BEFORE counting and pagination,
86    /// so the returned page and `total` stay consistent (no short/empty pages).
87    ///
88    /// Used by the listing API to hide cancelled-in-progress transactions from
89    /// active-status queries; internal callers use the unfiltered variant so they
90    /// still see the still-active replacement NOOP.
91    async fn find_by_status_paginated_filtered(
92        &self,
93        relayer_id: &str,
94        statuses: &[TransactionStatus],
95        query: PaginationQuery,
96        oldest_first: bool,
97        exclude_canceled: bool,
98    ) -> Result<PaginatedResult<TransactionRepoModel>, RepositoryError>;
99
100    /// Find a transaction by relayer ID and nonce
101    async fn find_by_nonce(
102        &self,
103        relayer_id: &str,
104        nonce: u64,
105    ) -> Result<Option<TransactionRepoModel>, RepositoryError>;
106
107    /// Returns the transaction status for each nonce in `[from_nonce, to_nonce)`.
108    ///
109    /// For each nonce, returns `Some(status)` if a transaction exists at that slot,
110    /// or `None` if the slot is empty. Implementations should batch I/O where possible
111    /// (e.g., Redis MGET) to minimize round trips.
112    async fn get_nonce_occupancy(
113        &self,
114        relayer_id: &str,
115        from_nonce: u64,
116        to_nonce: u64,
117    ) -> Result<Vec<(u64, Option<TransactionStatus>)>, RepositoryError>;
118
119    /// Update the status of a transaction
120    async fn update_status(
121        &self,
122        tx_id: String,
123        status: TransactionStatus,
124    ) -> Result<TransactionRepoModel, RepositoryError>;
125
126    /// Partially update a transaction
127    async fn partial_update(
128        &self,
129        tx_id: String,
130        update: TransactionUpdateRequest,
131    ) -> Result<TransactionRepoModel, RepositoryError>;
132
133    /// Update the network data of a transaction
134    async fn update_network_data(
135        &self,
136        tx_id: String,
137        network_data: NetworkTransactionData,
138    ) -> Result<TransactionRepoModel, RepositoryError>;
139
140    /// Set the sent_at timestamp of a transaction
141    async fn set_sent_at(
142        &self,
143        tx_id: String,
144        sent_at: String,
145    ) -> Result<TransactionRepoModel, RepositoryError>;
146
147    /// Atomically increments status-check failure counters using the latest stored metadata.
148    async fn increment_status_check_failures(
149        &self,
150        tx_id: String,
151    ) -> Result<TransactionRepoModel, RepositoryError>;
152
153    /// Atomically resets consecutive status-check failures to zero while preserving other counters.
154    async fn reset_status_check_consecutive_failures(
155        &self,
156        tx_id: String,
157    ) -> Result<TransactionRepoModel, RepositoryError>;
158
159    /// Atomically sets `sent_at` and increments Stellar insufficient-fee retries.
160    async fn record_stellar_insufficient_fee_retry(
161        &self,
162        tx_id: String,
163        sent_at: String,
164    ) -> Result<TransactionRepoModel, RepositoryError>;
165
166    /// Atomically sets `sent_at` and increments Stellar try-again-later retries.
167    async fn record_stellar_try_again_later_retry(
168        &self,
169        tx_id: String,
170        sent_at: String,
171    ) -> Result<TransactionRepoModel, RepositoryError>;
172
173    /// Set the confirmed_at timestamp of a transaction
174    async fn set_confirmed_at(
175        &self,
176        tx_id: String,
177        confirmed_at: String,
178    ) -> Result<TransactionRepoModel, RepositoryError>;
179
180    /// Count transactions by status(es) without fetching full transaction data.
181    /// This is an optimized O(1) operation in Redis using ZCARD.
182    async fn count_by_status(
183        &self,
184        relayer_id: &str,
185        statuses: &[TransactionStatus],
186    ) -> Result<u64, RepositoryError>;
187
188    /// Delete multiple transactions by their IDs in a single batch operation.
189    ///
190    /// This is more efficient than calling `delete_by_id` multiple times as it
191    /// reduces the number of round-trips to the storage backend.
192    ///
193    /// Note: This method requires fetching transaction data first to clean up indexes.
194    /// If you already have transaction data, use `delete_by_requests` instead for
195    /// better performance.
196    ///
197    /// # Arguments
198    /// * `ids` - List of transaction IDs to delete
199    ///
200    /// # Returns
201    /// * `BatchDeleteResult` containing the count of successful deletions and any failures
202    async fn delete_by_ids(&self, ids: Vec<String>) -> Result<BatchDeleteResult, RepositoryError>;
203
204    /// Delete multiple transactions using pre-extracted data.
205    ///
206    /// This is the most efficient batch delete method as it doesn't require
207    /// re-fetching transaction data. Use this when you already have the transaction
208    /// data (e.g., from a previous query).
209    ///
210    /// # Arguments
211    /// * `requests` - List of delete requests containing transaction data needed for cleanup
212    ///
213    /// # Returns
214    /// * `BatchDeleteResult` containing the count of successful deletions and any failures
215    async fn delete_by_requests(
216        &self,
217        requests: Vec<TransactionDeleteRequest>,
218    ) -> Result<BatchDeleteResult, RepositoryError>;
219}
220
221#[cfg(test)]
222mockall::mock! {
223  pub TransactionRepository {}
224
225  #[async_trait]
226  impl Repository<TransactionRepoModel, String> for TransactionRepository {
227      async fn create(&self, entity: TransactionRepoModel) -> Result<TransactionRepoModel, RepositoryError>;
228      async fn get_by_id(&self, id: String) -> Result<TransactionRepoModel, RepositoryError>;
229      async fn list_all(&self) -> Result<Vec<TransactionRepoModel>, RepositoryError>;
230      async fn list_paginated(&self, query: PaginationQuery) -> Result<PaginatedResult<TransactionRepoModel>, RepositoryError>;
231      async fn update(&self, id: String, entity: TransactionRepoModel) -> Result<TransactionRepoModel, RepositoryError>;
232      async fn delete_by_id(&self, id: String) -> Result<(), RepositoryError>;
233      async fn count(&self) -> Result<usize, RepositoryError>;
234      async fn has_entries(&self) -> Result<bool, RepositoryError>;
235      async fn drop_all_entries(&self) -> Result<(), RepositoryError>;
236  }
237
238  #[async_trait]
239  impl TransactionRepository for TransactionRepository {
240      fn connection_info(&self) -> Option<(Arc<RedisConnections>, String)>;
241      async fn find_by_relayer_id(&self, relayer_id: &str, query: PaginationQuery) -> Result<PaginatedResult<TransactionRepoModel>, RepositoryError>;
242      async fn find_by_status(&self, relayer_id: &str, statuses: &[TransactionStatus]) -> Result<Vec<TransactionRepoModel>, RepositoryError>;
243      async fn find_by_status_paginated(&self, relayer_id: &str, statuses: &[TransactionStatus], query: PaginationQuery, oldest_first: bool) -> Result<PaginatedResult<TransactionRepoModel>, RepositoryError>;
244      async fn find_by_status_paginated_filtered(&self, relayer_id: &str, statuses: &[TransactionStatus], query: PaginationQuery, oldest_first: bool, exclude_canceled: bool) -> Result<PaginatedResult<TransactionRepoModel>, RepositoryError>;
245      async fn find_by_nonce(&self, relayer_id: &str, nonce: u64) -> Result<Option<TransactionRepoModel>, RepositoryError>;
246      async fn get_nonce_occupancy(&self, relayer_id: &str, from_nonce: u64, to_nonce: u64) -> Result<Vec<(u64, Option<TransactionStatus>)>, RepositoryError>;
247      async fn update_status(&self, tx_id: String, status: TransactionStatus) -> Result<TransactionRepoModel, RepositoryError>;
248      async fn partial_update(&self, tx_id: String, update: TransactionUpdateRequest) -> Result<TransactionRepoModel, RepositoryError>;
249      async fn update_network_data(&self, tx_id: String, network_data: NetworkTransactionData) -> Result<TransactionRepoModel, RepositoryError>;
250      async fn set_sent_at(&self, tx_id: String, sent_at: String) -> Result<TransactionRepoModel, RepositoryError>;
251      async fn increment_status_check_failures(&self, tx_id: String) -> Result<TransactionRepoModel, RepositoryError>;
252      async fn reset_status_check_consecutive_failures(&self, tx_id: String) -> Result<TransactionRepoModel, RepositoryError>;
253      async fn record_stellar_insufficient_fee_retry(&self, tx_id: String, sent_at: String) -> Result<TransactionRepoModel, RepositoryError>;
254      async fn record_stellar_try_again_later_retry(&self, tx_id: String, sent_at: String) -> Result<TransactionRepoModel, RepositoryError>;
255      async fn set_confirmed_at(&self, tx_id: String, confirmed_at: String) -> Result<TransactionRepoModel, RepositoryError>;
256      async fn count_by_status(&self, relayer_id: &str, statuses: &[TransactionStatus]) -> Result<u64, RepositoryError>;
257      async fn delete_by_ids(&self, ids: Vec<String>) -> Result<BatchDeleteResult, RepositoryError>;
258      async fn delete_by_requests(&self, requests: Vec<TransactionDeleteRequest>) -> Result<BatchDeleteResult, RepositoryError>;
259  }
260}
261
262/// Enum wrapper for different transaction repository implementations
263#[derive(Debug, Clone)]
264pub enum TransactionRepositoryStorage {
265    InMemory(InMemoryTransactionRepository),
266    Redis(RedisTransactionRepository),
267}
268
269impl TransactionRepositoryStorage {
270    pub fn new_in_memory() -> Self {
271        Self::InMemory(InMemoryTransactionRepository::new())
272    }
273    pub fn new_redis(
274        connections: Arc<RedisConnections>,
275        key_prefix: String,
276    ) -> Result<Self, RepositoryError> {
277        Ok(Self::Redis(RedisTransactionRepository::new(
278            connections,
279            key_prefix,
280        )?))
281    }
282
283    /// Returns underlying Redis connections if this is a persistent storage backend.
284    ///
285    /// This is useful for operations that need direct storage access, such as
286    /// distributed locking and health checks.
287    ///
288    /// # Returns
289    /// * `Some((connections, key_prefix))` - If using persistent Redis storage
290    /// * `None` - If using in-memory storage
291    pub fn connection_info(&self) -> Option<(Arc<RedisConnections>, &str)> {
292        match self {
293            TransactionRepositoryStorage::InMemory(_) => None,
294            TransactionRepositoryStorage::Redis(repo) => {
295                Some((repo.connections.clone(), &repo.key_prefix))
296            }
297        }
298    }
299
300    /// Returns key prefix used by persistent storage backends.
301    pub fn key_prefix(&self) -> Option<&str> {
302        match self {
303            TransactionRepositoryStorage::InMemory(_) => None,
304            TransactionRepositoryStorage::Redis(repo) => Some(&repo.key_prefix),
305        }
306    }
307}
308
309#[async_trait]
310impl TransactionRepository for TransactionRepositoryStorage {
311    fn connection_info(&self) -> Option<(Arc<RedisConnections>, String)> {
312        TransactionRepositoryStorage::connection_info(self)
313            .map(|(connections, key_prefix)| (connections, key_prefix.to_string()))
314    }
315
316    async fn find_by_relayer_id(
317        &self,
318        relayer_id: &str,
319        query: PaginationQuery,
320    ) -> Result<PaginatedResult<TransactionRepoModel>, RepositoryError> {
321        match self {
322            TransactionRepositoryStorage::InMemory(repo) => {
323                repo.find_by_relayer_id(relayer_id, query).await
324            }
325            TransactionRepositoryStorage::Redis(repo) => {
326                repo.find_by_relayer_id(relayer_id, query).await
327            }
328        }
329    }
330
331    async fn find_by_status(
332        &self,
333        relayer_id: &str,
334        statuses: &[TransactionStatus],
335    ) -> Result<Vec<TransactionRepoModel>, RepositoryError> {
336        match self {
337            TransactionRepositoryStorage::InMemory(repo) => {
338                repo.find_by_status(relayer_id, statuses).await
339            }
340            TransactionRepositoryStorage::Redis(repo) => {
341                repo.find_by_status(relayer_id, statuses).await
342            }
343        }
344    }
345
346    async fn find_by_status_paginated(
347        &self,
348        relayer_id: &str,
349        statuses: &[TransactionStatus],
350        query: PaginationQuery,
351        oldest_first: bool,
352    ) -> Result<PaginatedResult<TransactionRepoModel>, RepositoryError> {
353        match self {
354            TransactionRepositoryStorage::InMemory(repo) => {
355                repo.find_by_status_paginated(relayer_id, statuses, query, oldest_first)
356                    .await
357            }
358            TransactionRepositoryStorage::Redis(repo) => {
359                repo.find_by_status_paginated(relayer_id, statuses, query, oldest_first)
360                    .await
361            }
362        }
363    }
364
365    async fn find_by_status_paginated_filtered(
366        &self,
367        relayer_id: &str,
368        statuses: &[TransactionStatus],
369        query: PaginationQuery,
370        oldest_first: bool,
371        exclude_canceled: bool,
372    ) -> Result<PaginatedResult<TransactionRepoModel>, RepositoryError> {
373        match self {
374            TransactionRepositoryStorage::InMemory(repo) => {
375                repo.find_by_status_paginated_filtered(
376                    relayer_id,
377                    statuses,
378                    query,
379                    oldest_first,
380                    exclude_canceled,
381                )
382                .await
383            }
384            TransactionRepositoryStorage::Redis(repo) => {
385                repo.find_by_status_paginated_filtered(
386                    relayer_id,
387                    statuses,
388                    query,
389                    oldest_first,
390                    exclude_canceled,
391                )
392                .await
393            }
394        }
395    }
396
397    async fn find_by_nonce(
398        &self,
399        relayer_id: &str,
400        nonce: u64,
401    ) -> Result<Option<TransactionRepoModel>, RepositoryError> {
402        match self {
403            TransactionRepositoryStorage::InMemory(repo) => {
404                repo.find_by_nonce(relayer_id, nonce).await
405            }
406            TransactionRepositoryStorage::Redis(repo) => {
407                repo.find_by_nonce(relayer_id, nonce).await
408            }
409        }
410    }
411
412    async fn get_nonce_occupancy(
413        &self,
414        relayer_id: &str,
415        from_nonce: u64,
416        to_nonce: u64,
417    ) -> Result<Vec<(u64, Option<TransactionStatus>)>, RepositoryError> {
418        match self {
419            TransactionRepositoryStorage::InMemory(repo) => {
420                repo.get_nonce_occupancy(relayer_id, from_nonce, to_nonce)
421                    .await
422            }
423            TransactionRepositoryStorage::Redis(repo) => {
424                repo.get_nonce_occupancy(relayer_id, from_nonce, to_nonce)
425                    .await
426            }
427        }
428    }
429
430    async fn update_status(
431        &self,
432        tx_id: String,
433        status: TransactionStatus,
434    ) -> Result<TransactionRepoModel, RepositoryError> {
435        match self {
436            TransactionRepositoryStorage::InMemory(repo) => repo.update_status(tx_id, status).await,
437            TransactionRepositoryStorage::Redis(repo) => repo.update_status(tx_id, status).await,
438        }
439    }
440
441    async fn partial_update(
442        &self,
443        tx_id: String,
444        update: TransactionUpdateRequest,
445    ) -> Result<TransactionRepoModel, RepositoryError> {
446        match self {
447            TransactionRepositoryStorage::InMemory(repo) => {
448                repo.partial_update(tx_id, update).await
449            }
450            TransactionRepositoryStorage::Redis(repo) => repo.partial_update(tx_id, update).await,
451        }
452    }
453
454    async fn update_network_data(
455        &self,
456        tx_id: String,
457        network_data: NetworkTransactionData,
458    ) -> Result<TransactionRepoModel, RepositoryError> {
459        match self {
460            TransactionRepositoryStorage::InMemory(repo) => {
461                repo.update_network_data(tx_id, network_data).await
462            }
463            TransactionRepositoryStorage::Redis(repo) => {
464                repo.update_network_data(tx_id, network_data).await
465            }
466        }
467    }
468
469    async fn set_sent_at(
470        &self,
471        tx_id: String,
472        sent_at: String,
473    ) -> Result<TransactionRepoModel, RepositoryError> {
474        match self {
475            TransactionRepositoryStorage::InMemory(repo) => repo.set_sent_at(tx_id, sent_at).await,
476            TransactionRepositoryStorage::Redis(repo) => repo.set_sent_at(tx_id, sent_at).await,
477        }
478    }
479
480    async fn increment_status_check_failures(
481        &self,
482        tx_id: String,
483    ) -> Result<TransactionRepoModel, RepositoryError> {
484        match self {
485            TransactionRepositoryStorage::InMemory(repo) => {
486                repo.increment_status_check_failures(tx_id).await
487            }
488            TransactionRepositoryStorage::Redis(repo) => {
489                repo.increment_status_check_failures(tx_id).await
490            }
491        }
492    }
493
494    async fn reset_status_check_consecutive_failures(
495        &self,
496        tx_id: String,
497    ) -> Result<TransactionRepoModel, RepositoryError> {
498        match self {
499            TransactionRepositoryStorage::InMemory(repo) => {
500                repo.reset_status_check_consecutive_failures(tx_id).await
501            }
502            TransactionRepositoryStorage::Redis(repo) => {
503                repo.reset_status_check_consecutive_failures(tx_id).await
504            }
505        }
506    }
507
508    async fn record_stellar_insufficient_fee_retry(
509        &self,
510        tx_id: String,
511        sent_at: String,
512    ) -> Result<TransactionRepoModel, RepositoryError> {
513        match self {
514            TransactionRepositoryStorage::InMemory(repo) => {
515                repo.record_stellar_insufficient_fee_retry(tx_id, sent_at)
516                    .await
517            }
518            TransactionRepositoryStorage::Redis(repo) => {
519                repo.record_stellar_insufficient_fee_retry(tx_id, sent_at)
520                    .await
521            }
522        }
523    }
524
525    async fn record_stellar_try_again_later_retry(
526        &self,
527        tx_id: String,
528        sent_at: String,
529    ) -> Result<TransactionRepoModel, RepositoryError> {
530        match self {
531            TransactionRepositoryStorage::InMemory(repo) => {
532                repo.record_stellar_try_again_later_retry(tx_id, sent_at)
533                    .await
534            }
535            TransactionRepositoryStorage::Redis(repo) => {
536                repo.record_stellar_try_again_later_retry(tx_id, sent_at)
537                    .await
538            }
539        }
540    }
541
542    async fn set_confirmed_at(
543        &self,
544        tx_id: String,
545        confirmed_at: String,
546    ) -> Result<TransactionRepoModel, RepositoryError> {
547        match self {
548            TransactionRepositoryStorage::InMemory(repo) => {
549                repo.set_confirmed_at(tx_id, confirmed_at).await
550            }
551            TransactionRepositoryStorage::Redis(repo) => {
552                repo.set_confirmed_at(tx_id, confirmed_at).await
553            }
554        }
555    }
556
557    async fn count_by_status(
558        &self,
559        relayer_id: &str,
560        statuses: &[TransactionStatus],
561    ) -> Result<u64, RepositoryError> {
562        match self {
563            TransactionRepositoryStorage::InMemory(repo) => {
564                repo.count_by_status(relayer_id, statuses).await
565            }
566            TransactionRepositoryStorage::Redis(repo) => {
567                repo.count_by_status(relayer_id, statuses).await
568            }
569        }
570    }
571
572    async fn delete_by_ids(&self, ids: Vec<String>) -> Result<BatchDeleteResult, RepositoryError> {
573        match self {
574            TransactionRepositoryStorage::InMemory(repo) => repo.delete_by_ids(ids).await,
575            TransactionRepositoryStorage::Redis(repo) => repo.delete_by_ids(ids).await,
576        }
577    }
578
579    async fn delete_by_requests(
580        &self,
581        requests: Vec<TransactionDeleteRequest>,
582    ) -> Result<BatchDeleteResult, RepositoryError> {
583        match self {
584            TransactionRepositoryStorage::InMemory(repo) => repo.delete_by_requests(requests).await,
585            TransactionRepositoryStorage::Redis(repo) => repo.delete_by_requests(requests).await,
586        }
587    }
588}
589
590#[async_trait]
591impl Repository<TransactionRepoModel, String> for TransactionRepositoryStorage {
592    async fn create(
593        &self,
594        entity: TransactionRepoModel,
595    ) -> Result<TransactionRepoModel, RepositoryError> {
596        match self {
597            TransactionRepositoryStorage::InMemory(repo) => repo.create(entity).await,
598            TransactionRepositoryStorage::Redis(repo) => repo.create(entity).await,
599        }
600    }
601
602    async fn get_by_id(&self, id: String) -> Result<TransactionRepoModel, RepositoryError> {
603        match self {
604            TransactionRepositoryStorage::InMemory(repo) => repo.get_by_id(id).await,
605            TransactionRepositoryStorage::Redis(repo) => repo.get_by_id(id).await,
606        }
607    }
608
609    async fn list_all(&self) -> Result<Vec<TransactionRepoModel>, RepositoryError> {
610        match self {
611            TransactionRepositoryStorage::InMemory(repo) => repo.list_all().await,
612            TransactionRepositoryStorage::Redis(repo) => repo.list_all().await,
613        }
614    }
615
616    async fn list_paginated(
617        &self,
618        query: PaginationQuery,
619    ) -> Result<PaginatedResult<TransactionRepoModel>, RepositoryError> {
620        match self {
621            TransactionRepositoryStorage::InMemory(repo) => repo.list_paginated(query).await,
622            TransactionRepositoryStorage::Redis(repo) => repo.list_paginated(query).await,
623        }
624    }
625
626    async fn update(
627        &self,
628        id: String,
629        entity: TransactionRepoModel,
630    ) -> Result<TransactionRepoModel, RepositoryError> {
631        match self {
632            TransactionRepositoryStorage::InMemory(repo) => repo.update(id, entity).await,
633            TransactionRepositoryStorage::Redis(repo) => repo.update(id, entity).await,
634        }
635    }
636
637    async fn delete_by_id(&self, id: String) -> Result<(), RepositoryError> {
638        match self {
639            TransactionRepositoryStorage::InMemory(repo) => repo.delete_by_id(id).await,
640            TransactionRepositoryStorage::Redis(repo) => repo.delete_by_id(id).await,
641        }
642    }
643
644    async fn count(&self) -> Result<usize, RepositoryError> {
645        match self {
646            TransactionRepositoryStorage::InMemory(repo) => repo.count().await,
647            TransactionRepositoryStorage::Redis(repo) => repo.count().await,
648        }
649    }
650
651    async fn has_entries(&self) -> Result<bool, RepositoryError> {
652        match self {
653            TransactionRepositoryStorage::InMemory(repo) => repo.has_entries().await,
654            TransactionRepositoryStorage::Redis(repo) => repo.has_entries().await,
655        }
656    }
657
658    async fn drop_all_entries(&self) -> Result<(), RepositoryError> {
659        match self {
660            TransactionRepositoryStorage::InMemory(repo) => repo.drop_all_entries().await,
661            TransactionRepositoryStorage::Redis(repo) => repo.drop_all_entries().await,
662        }
663    }
664}
665
666#[cfg(test)]
667mod tests {
668    use chrono::Utc;
669    use color_eyre::Result;
670    use deadpool_redis::{Config, Runtime};
671
672    use super::*;
673    use crate::models::{
674        EvmTransactionData, NetworkTransactionData, TransactionStatus, TransactionUpdateRequest,
675    };
676    use crate::repositories::PaginationQuery;
677    use crate::utils::mocks::mockutils::create_mock_transaction;
678
679    fn create_test_transaction(id: &str, relayer_id: &str) -> TransactionRepoModel {
680        let mut transaction = create_mock_transaction();
681        transaction.id = id.to_string();
682        transaction.relayer_id = relayer_id.to_string();
683        transaction
684    }
685
686    fn create_test_transaction_with_status(
687        id: &str,
688        relayer_id: &str,
689        status: TransactionStatus,
690    ) -> TransactionRepoModel {
691        let mut transaction = create_test_transaction(id, relayer_id);
692        transaction.status = status;
693        transaction
694    }
695
696    fn create_test_transaction_with_nonce(
697        id: &str,
698        relayer_id: &str,
699        nonce: u64,
700    ) -> TransactionRepoModel {
701        let mut transaction = create_test_transaction(id, relayer_id);
702        if let NetworkTransactionData::Evm(ref mut evm_data) = transaction.network_data {
703            evm_data.nonce = Some(nonce);
704        }
705        transaction
706    }
707
708    fn create_test_update_request() -> TransactionUpdateRequest {
709        TransactionUpdateRequest {
710            status: Some(TransactionStatus::Sent),
711            status_reason: Some("Test reason".to_string()),
712            sent_at: Some(Utc::now().to_string()),
713            confirmed_at: None,
714            network_data: None,
715            priced_at: None,
716            hashes: Some(vec!["test_hash".to_string()]),
717            noop_count: None,
718            is_canceled: None,
719            delete_at: None,
720            metadata: None,
721        }
722    }
723
724    #[tokio::test]
725    async fn test_new_in_memory() {
726        let storage = TransactionRepositoryStorage::new_in_memory();
727
728        match storage {
729            TransactionRepositoryStorage::InMemory(_) => {
730                // Success - verify it's the InMemory variant
731            }
732            TransactionRepositoryStorage::Redis(_) => {
733                panic!("Expected InMemory variant, got Redis");
734            }
735        }
736    }
737
738    #[tokio::test]
739    async fn test_connection_info_returns_none_for_in_memory() {
740        let storage = TransactionRepositoryStorage::new_in_memory();
741
742        // In-memory storage should return None for connection_info
743        assert!(storage.connection_info().is_none());
744    }
745
746    #[tokio::test]
747    #[ignore = "Requires active Redis instance"]
748    async fn test_connection_info_returns_some_for_redis() -> Result<()> {
749        let redis_url = std::env::var("REDIS_TEST_URL")
750            .unwrap_or_else(|_| "redis://127.0.0.1:6379".to_string());
751        let cfg = Config::from_url(&redis_url);
752        let pool = Arc::new(
753            cfg.builder()
754                .map_err(|e| eyre::eyre!("Failed to create Redis pool builder: {}", e))?
755                .max_size(16)
756                .runtime(Runtime::Tokio1)
757                .build()
758                .map_err(|e| eyre::eyre!("Failed to build Redis pool: {}", e))?,
759        );
760        let connections = Arc::new(RedisConnections::new_single_pool(pool.clone()));
761        let key_prefix = "test_prefix".to_string();
762
763        let storage = TransactionRepositoryStorage::new_redis(connections, key_prefix.clone())?;
764
765        let (returned_connection, returned_prefix) = storage
766            .connection_info()
767            .expect("Expected Redis connection info");
768
769        assert!(Arc::ptr_eq(&pool, returned_connection.primary()));
770        assert_eq!(returned_prefix, key_prefix);
771
772        Ok(())
773    }
774
775    #[tokio::test]
776    async fn test_create_in_memory() -> Result<()> {
777        let storage = TransactionRepositoryStorage::new_in_memory();
778        let transaction = create_test_transaction("test-tx", "test-relayer");
779
780        let created = storage.create(transaction.clone()).await?;
781        assert_eq!(created.id, transaction.id);
782        assert_eq!(created.relayer_id, transaction.relayer_id);
783        assert_eq!(created.status, transaction.status);
784
785        Ok(())
786    }
787
788    #[tokio::test]
789    async fn test_get_by_id_in_memory() -> Result<()> {
790        let storage = TransactionRepositoryStorage::new_in_memory();
791        let transaction = create_test_transaction("test-tx", "test-relayer");
792
793        // Create transaction first
794        storage.create(transaction.clone()).await?;
795
796        // Get by ID
797        let retrieved = storage.get_by_id("test-tx".to_string()).await?;
798        assert_eq!(retrieved.id, transaction.id);
799        assert_eq!(retrieved.relayer_id, transaction.relayer_id);
800        assert_eq!(retrieved.status, transaction.status);
801
802        Ok(())
803    }
804
805    #[tokio::test]
806    async fn test_get_by_id_not_found_in_memory() -> Result<()> {
807        let storage = TransactionRepositoryStorage::new_in_memory();
808
809        let result = storage.get_by_id("non-existent".to_string()).await;
810        assert!(result.is_err());
811
812        Ok(())
813    }
814
815    #[tokio::test]
816    async fn test_list_all_in_memory() -> Result<()> {
817        let storage = TransactionRepositoryStorage::new_in_memory();
818
819        // Initially empty
820        let transactions = storage.list_all().await?;
821        assert!(transactions.is_empty());
822
823        // Add transactions
824        let tx1 = create_test_transaction("tx-1", "relayer-1");
825        let tx2 = create_test_transaction("tx-2", "relayer-2");
826
827        storage.create(tx1.clone()).await?;
828        storage.create(tx2.clone()).await?;
829
830        let all_transactions = storage.list_all().await?;
831        assert_eq!(all_transactions.len(), 2);
832
833        let ids: Vec<&str> = all_transactions.iter().map(|t| t.id.as_str()).collect();
834        assert!(ids.contains(&"tx-1"));
835        assert!(ids.contains(&"tx-2"));
836
837        Ok(())
838    }
839
840    #[tokio::test]
841    async fn test_list_paginated_in_memory() -> Result<()> {
842        let storage = TransactionRepositoryStorage::new_in_memory();
843
844        // Add test transactions
845        for i in 1..=5 {
846            let tx = create_test_transaction(&format!("tx-{i}"), "test-relayer");
847            storage.create(tx).await?;
848        }
849
850        // Test pagination
851        let query = PaginationQuery {
852            page: 1,
853            per_page: 2,
854        };
855        let page = storage.list_paginated(query).await?;
856
857        assert_eq!(page.items.len(), 2);
858        assert_eq!(page.total, 5);
859        assert_eq!(page.page, 1);
860        assert_eq!(page.per_page, 2);
861
862        // Test second page
863        let query2 = PaginationQuery {
864            page: 2,
865            per_page: 2,
866        };
867        let page2 = storage.list_paginated(query2).await?;
868
869        assert_eq!(page2.items.len(), 2);
870        assert_eq!(page2.total, 5);
871        assert_eq!(page2.page, 2);
872        assert_eq!(page2.per_page, 2);
873
874        Ok(())
875    }
876
877    #[tokio::test]
878    async fn test_update_in_memory() -> Result<()> {
879        let storage = TransactionRepositoryStorage::new_in_memory();
880        let transaction = create_test_transaction("test-tx", "test-relayer");
881
882        // Create transaction first
883        storage.create(transaction.clone()).await?;
884
885        // Update it
886        let mut updated_transaction = transaction.clone();
887        updated_transaction.status = TransactionStatus::Sent;
888        updated_transaction.status_reason = Some("Updated reason".to_string());
889
890        let result = storage
891            .update("test-tx".to_string(), updated_transaction.clone())
892            .await?;
893        assert_eq!(result.id, "test-tx");
894        assert_eq!(result.status, TransactionStatus::Sent);
895        assert_eq!(result.status_reason, Some("Updated reason".to_string()));
896
897        // Verify the update persisted
898        let retrieved = storage.get_by_id("test-tx".to_string()).await?;
899        assert_eq!(retrieved.status, TransactionStatus::Sent);
900        assert_eq!(retrieved.status_reason, Some("Updated reason".to_string()));
901
902        Ok(())
903    }
904
905    #[tokio::test]
906    async fn test_update_not_found_in_memory() -> Result<()> {
907        let storage = TransactionRepositoryStorage::new_in_memory();
908        let transaction = create_test_transaction("non-existent", "test-relayer");
909
910        let result = storage
911            .update("non-existent".to_string(), transaction)
912            .await;
913        assert!(result.is_err());
914
915        Ok(())
916    }
917
918    #[tokio::test]
919    async fn test_delete_by_id_in_memory() -> Result<()> {
920        let storage = TransactionRepositoryStorage::new_in_memory();
921        let transaction = create_test_transaction("test-tx", "test-relayer");
922
923        // Create transaction first
924        storage.create(transaction.clone()).await?;
925
926        // Verify it exists
927        let retrieved = storage.get_by_id("test-tx".to_string()).await?;
928        assert_eq!(retrieved.id, "test-tx");
929
930        // Delete it
931        storage.delete_by_id("test-tx".to_string()).await?;
932
933        // Verify it's gone
934        let result = storage.get_by_id("test-tx".to_string()).await;
935        assert!(result.is_err());
936
937        Ok(())
938    }
939
940    #[tokio::test]
941    async fn test_delete_by_id_not_found_in_memory() -> Result<()> {
942        let storage = TransactionRepositoryStorage::new_in_memory();
943
944        let result = storage.delete_by_id("non-existent".to_string()).await;
945        assert!(result.is_err());
946
947        Ok(())
948    }
949
950    #[tokio::test]
951    async fn test_count_in_memory() -> Result<()> {
952        let storage = TransactionRepositoryStorage::new_in_memory();
953
954        // Initially empty
955        let count = storage.count().await?;
956        assert_eq!(count, 0);
957
958        // Add transactions
959        let tx1 = create_test_transaction("tx-1", "relayer-1");
960        let tx2 = create_test_transaction("tx-2", "relayer-2");
961
962        storage.create(tx1).await?;
963        let count_after_one = storage.count().await?;
964        assert_eq!(count_after_one, 1);
965
966        storage.create(tx2).await?;
967        let count_after_two = storage.count().await?;
968        assert_eq!(count_after_two, 2);
969
970        // Delete one
971        storage.delete_by_id("tx-1".to_string()).await?;
972        let count_after_delete = storage.count().await?;
973        assert_eq!(count_after_delete, 1);
974
975        Ok(())
976    }
977
978    #[tokio::test]
979    async fn test_has_entries_in_memory() -> Result<()> {
980        let storage = TransactionRepositoryStorage::new_in_memory();
981
982        // Initially empty
983        let has_entries = storage.has_entries().await?;
984        assert!(!has_entries);
985
986        // Add transaction
987        let transaction = create_test_transaction("test-tx", "test-relayer");
988        storage.create(transaction).await?;
989
990        let has_entries_after_create = storage.has_entries().await?;
991        assert!(has_entries_after_create);
992
993        // Delete transaction
994        storage.delete_by_id("test-tx".to_string()).await?;
995
996        let has_entries_after_delete = storage.has_entries().await?;
997        assert!(!has_entries_after_delete);
998
999        Ok(())
1000    }
1001
1002    #[tokio::test]
1003    async fn test_drop_all_entries_in_memory() -> Result<()> {
1004        let storage = TransactionRepositoryStorage::new_in_memory();
1005
1006        // Add multiple transactions
1007        for i in 1..=5 {
1008            let tx = create_test_transaction(&format!("tx-{i}"), "test-relayer");
1009            storage.create(tx).await?;
1010        }
1011
1012        // Verify they exist
1013        let count_before = storage.count().await?;
1014        assert_eq!(count_before, 5);
1015
1016        let has_entries_before = storage.has_entries().await?;
1017        assert!(has_entries_before);
1018
1019        // Drop all entries
1020        storage.drop_all_entries().await?;
1021
1022        // Verify they're gone
1023        let count_after = storage.count().await?;
1024        assert_eq!(count_after, 0);
1025
1026        let has_entries_after = storage.has_entries().await?;
1027        assert!(!has_entries_after);
1028
1029        let all_transactions = storage.list_all().await?;
1030        assert!(all_transactions.is_empty());
1031
1032        Ok(())
1033    }
1034
1035    #[tokio::test]
1036    async fn test_find_by_relayer_id_in_memory() -> Result<()> {
1037        let storage = TransactionRepositoryStorage::new_in_memory();
1038
1039        // Add transactions for different relayers
1040        let tx1 = create_test_transaction("tx-1", "relayer-1");
1041        let tx2 = create_test_transaction("tx-2", "relayer-1");
1042        let tx3 = create_test_transaction("tx-3", "relayer-2");
1043
1044        storage.create(tx1).await?;
1045        storage.create(tx2).await?;
1046        storage.create(tx3).await?;
1047
1048        // Find by relayer ID
1049        let query = PaginationQuery {
1050            page: 1,
1051            per_page: 10,
1052        };
1053        let result = storage.find_by_relayer_id("relayer-1", query).await?;
1054
1055        assert_eq!(result.items.len(), 2);
1056        assert_eq!(result.total, 2);
1057
1058        // Verify all transactions belong to relayer-1
1059        for tx in result.items {
1060            assert_eq!(tx.relayer_id, "relayer-1");
1061        }
1062
1063        Ok(())
1064    }
1065
1066    #[tokio::test]
1067    async fn test_find_by_status_in_memory() -> Result<()> {
1068        let storage = TransactionRepositoryStorage::new_in_memory();
1069
1070        // Add transactions with different statuses
1071        let tx1 =
1072            create_test_transaction_with_status("tx-1", "relayer-1", TransactionStatus::Pending);
1073        let tx2 = create_test_transaction_with_status("tx-2", "relayer-1", TransactionStatus::Sent);
1074        let tx3 =
1075            create_test_transaction_with_status("tx-3", "relayer-1", TransactionStatus::Pending);
1076        let tx4 =
1077            create_test_transaction_with_status("tx-4", "relayer-2", TransactionStatus::Pending);
1078
1079        storage.create(tx1).await?;
1080        storage.create(tx2).await?;
1081        storage.create(tx3).await?;
1082        storage.create(tx4).await?;
1083
1084        // Find by status
1085        let statuses = vec![TransactionStatus::Pending];
1086        let result = storage.find_by_status("relayer-1", &statuses).await?;
1087
1088        assert_eq!(result.len(), 2);
1089
1090        // Verify all transactions have Pending status and belong to relayer-1
1091        for tx in result {
1092            assert_eq!(tx.status, TransactionStatus::Pending);
1093            assert_eq!(tx.relayer_id, "relayer-1");
1094        }
1095
1096        Ok(())
1097    }
1098
1099    #[tokio::test]
1100    async fn test_find_by_nonce_in_memory() -> Result<()> {
1101        let storage = TransactionRepositoryStorage::new_in_memory();
1102
1103        // Add transactions with different nonces
1104        let tx1 = create_test_transaction_with_nonce("tx-1", "relayer-1", 10);
1105        let tx2 = create_test_transaction_with_nonce("tx-2", "relayer-1", 20);
1106        let tx3 = create_test_transaction_with_nonce("tx-3", "relayer-2", 10);
1107
1108        storage.create(tx1).await?;
1109        storage.create(tx2).await?;
1110        storage.create(tx3).await?;
1111
1112        // Find by nonce
1113        let result = storage.find_by_nonce("relayer-1", 10).await?;
1114
1115        assert!(result.is_some());
1116        let found_tx = result.unwrap();
1117        assert_eq!(found_tx.id, "tx-1");
1118        assert_eq!(found_tx.relayer_id, "relayer-1");
1119
1120        // Check EVM nonce
1121        if let NetworkTransactionData::Evm(evm_data) = found_tx.network_data {
1122            assert_eq!(evm_data.nonce, Some(10));
1123        }
1124
1125        // Test not found
1126        let not_found = storage.find_by_nonce("relayer-1", 99).await?;
1127        assert!(not_found.is_none());
1128
1129        Ok(())
1130    }
1131
1132    #[tokio::test]
1133    async fn test_update_status_in_memory() -> Result<()> {
1134        let storage = TransactionRepositoryStorage::new_in_memory();
1135        let transaction = create_test_transaction("test-tx", "test-relayer");
1136
1137        // Create transaction first
1138        storage.create(transaction).await?;
1139
1140        // Update status
1141        let updated = storage
1142            .update_status("test-tx".to_string(), TransactionStatus::Sent)
1143            .await?;
1144
1145        assert_eq!(updated.id, "test-tx");
1146        assert_eq!(updated.status, TransactionStatus::Sent);
1147
1148        // Verify the update persisted
1149        let retrieved = storage.get_by_id("test-tx".to_string()).await?;
1150        assert_eq!(retrieved.status, TransactionStatus::Sent);
1151
1152        Ok(())
1153    }
1154
1155    #[tokio::test]
1156    async fn test_partial_update_in_memory() -> Result<()> {
1157        let storage = TransactionRepositoryStorage::new_in_memory();
1158        let transaction = create_test_transaction("test-tx", "test-relayer");
1159
1160        // Create transaction first
1161        storage.create(transaction).await?;
1162
1163        // Partial update
1164        let update_request = create_test_update_request();
1165        let updated = storage
1166            .partial_update("test-tx".to_string(), update_request)
1167            .await?;
1168
1169        assert_eq!(updated.id, "test-tx");
1170        assert_eq!(updated.status, TransactionStatus::Sent);
1171        assert_eq!(updated.status_reason, Some("Test reason".to_string()));
1172        assert!(updated.sent_at.is_some());
1173        assert_eq!(updated.hashes, vec!["test_hash".to_string()]);
1174
1175        Ok(())
1176    }
1177
1178    #[tokio::test]
1179    async fn test_update_network_data_in_memory() -> Result<()> {
1180        let storage = TransactionRepositoryStorage::new_in_memory();
1181        let transaction = create_test_transaction("test-tx", "test-relayer");
1182
1183        // Create transaction first
1184        storage.create(transaction).await?;
1185
1186        // Update network data
1187        let new_evm_data = EvmTransactionData {
1188            nonce: Some(42),
1189            gas_limit: Some(21000),
1190            ..Default::default()
1191        };
1192        let new_network_data = NetworkTransactionData::Evm(new_evm_data);
1193
1194        let updated = storage
1195            .update_network_data("test-tx".to_string(), new_network_data)
1196            .await?;
1197
1198        assert_eq!(updated.id, "test-tx");
1199        if let NetworkTransactionData::Evm(evm_data) = updated.network_data {
1200            assert_eq!(evm_data.nonce, Some(42));
1201            assert_eq!(evm_data.gas_limit, Some(21000));
1202        } else {
1203            panic!("Expected EVM network data");
1204        }
1205
1206        Ok(())
1207    }
1208
1209    #[tokio::test]
1210    async fn test_set_sent_at_in_memory() -> Result<()> {
1211        let storage = TransactionRepositoryStorage::new_in_memory();
1212        let transaction = create_test_transaction("test-tx", "test-relayer");
1213
1214        // Create transaction first
1215        storage.create(transaction).await?;
1216
1217        // Set sent_at
1218        let sent_at = Utc::now().to_string();
1219        let updated = storage
1220            .set_sent_at("test-tx".to_string(), sent_at.clone())
1221            .await?;
1222
1223        assert_eq!(updated.id, "test-tx");
1224        assert_eq!(updated.sent_at, Some(sent_at));
1225
1226        Ok(())
1227    }
1228
1229    #[tokio::test]
1230    async fn test_set_confirmed_at_in_memory() -> Result<()> {
1231        let storage = TransactionRepositoryStorage::new_in_memory();
1232        let transaction = create_test_transaction("test-tx", "test-relayer");
1233
1234        // Create transaction first
1235        storage.create(transaction).await?;
1236
1237        // Set confirmed_at
1238        let confirmed_at = Utc::now().to_string();
1239        let updated = storage
1240            .set_confirmed_at("test-tx".to_string(), confirmed_at.clone())
1241            .await?;
1242
1243        assert_eq!(updated.id, "test-tx");
1244        assert_eq!(updated.confirmed_at, Some(confirmed_at));
1245
1246        Ok(())
1247    }
1248
1249    #[tokio::test]
1250    async fn test_create_duplicate_id_in_memory() -> Result<()> {
1251        let storage = TransactionRepositoryStorage::new_in_memory();
1252        let transaction = create_test_transaction("duplicate-id", "test-relayer");
1253
1254        // Create first transaction
1255        storage.create(transaction.clone()).await?;
1256
1257        // Try to create another with same ID - should fail
1258        let result = storage.create(transaction.clone()).await;
1259        assert!(result.is_err());
1260
1261        Ok(())
1262    }
1263
1264    #[tokio::test]
1265    async fn test_workflow_in_memory() -> Result<()> {
1266        let storage = TransactionRepositoryStorage::new_in_memory();
1267
1268        // 1. Start with empty storage
1269        assert!(!storage.has_entries().await?);
1270        assert_eq!(storage.count().await?, 0);
1271
1272        // 2. Create transaction
1273        let transaction = create_test_transaction("workflow-test", "test-relayer");
1274        let created = storage.create(transaction.clone()).await?;
1275        assert_eq!(created.id, "workflow-test");
1276
1277        // 3. Verify it exists
1278        assert!(storage.has_entries().await?);
1279        assert_eq!(storage.count().await?, 1);
1280
1281        // 4. Retrieve it
1282        let retrieved = storage.get_by_id("workflow-test".to_string()).await?;
1283        assert_eq!(retrieved.id, "workflow-test");
1284
1285        // 5. Update status
1286        let updated = storage
1287            .update_status("workflow-test".to_string(), TransactionStatus::Sent)
1288            .await?;
1289        assert_eq!(updated.status, TransactionStatus::Sent);
1290
1291        // 6. Verify update
1292        let retrieved_updated = storage.get_by_id("workflow-test".to_string()).await?;
1293        assert_eq!(retrieved_updated.status, TransactionStatus::Sent);
1294
1295        // 7. Delete it
1296        storage.delete_by_id("workflow-test".to_string()).await?;
1297
1298        // 8. Verify it's gone
1299        assert!(!storage.has_entries().await?);
1300        assert_eq!(storage.count().await?, 0);
1301
1302        let result = storage.get_by_id("workflow-test".to_string()).await;
1303        assert!(result.is_err());
1304
1305        Ok(())
1306    }
1307
1308    #[tokio::test]
1309    async fn test_multiple_relayers_workflow() -> Result<()> {
1310        let storage = TransactionRepositoryStorage::new_in_memory();
1311
1312        // Add transactions for multiple relayers
1313        let tx1 =
1314            create_test_transaction_with_status("tx-1", "relayer-1", TransactionStatus::Pending);
1315        let tx2 = create_test_transaction_with_status("tx-2", "relayer-1", TransactionStatus::Sent);
1316        let tx3 =
1317            create_test_transaction_with_status("tx-3", "relayer-2", TransactionStatus::Pending);
1318
1319        storage.create(tx1).await?;
1320        storage.create(tx2).await?;
1321        storage.create(tx3).await?;
1322
1323        // Test find_by_relayer_id
1324        let query = PaginationQuery {
1325            page: 1,
1326            per_page: 10,
1327        };
1328        let relayer1_txs = storage.find_by_relayer_id("relayer-1", query).await?;
1329        assert_eq!(relayer1_txs.items.len(), 2);
1330
1331        // Test find_by_status
1332        let pending_txs = storage
1333            .find_by_status("relayer-1", &[TransactionStatus::Pending])
1334            .await?;
1335        assert_eq!(pending_txs.len(), 1);
1336        assert_eq!(pending_txs[0].id, "tx-1");
1337
1338        // Test count remains accurate
1339        assert_eq!(storage.count().await?, 3);
1340
1341        Ok(())
1342    }
1343
1344    #[tokio::test]
1345    async fn test_pagination_edge_cases_in_memory() -> Result<()> {
1346        let storage = TransactionRepositoryStorage::new_in_memory();
1347
1348        // Test pagination with empty storage
1349        let query = PaginationQuery {
1350            page: 1,
1351            per_page: 10,
1352        };
1353        let page = storage.list_paginated(query).await?;
1354        assert_eq!(page.items.len(), 0);
1355        assert_eq!(page.total, 0);
1356        assert_eq!(page.page, 1);
1357        assert_eq!(page.per_page, 10);
1358
1359        // Add one transaction
1360        let transaction = create_test_transaction("single-tx", "test-relayer");
1361        storage.create(transaction).await?;
1362
1363        // Test pagination with single item
1364        let query = PaginationQuery {
1365            page: 1,
1366            per_page: 10,
1367        };
1368        let page = storage.list_paginated(query).await?;
1369        assert_eq!(page.items.len(), 1);
1370        assert_eq!(page.total, 1);
1371        assert_eq!(page.page, 1);
1372        assert_eq!(page.per_page, 10);
1373
1374        // Test pagination with page beyond total
1375        let query = PaginationQuery {
1376            page: 3,
1377            per_page: 10,
1378        };
1379        let page = storage.list_paginated(query).await?;
1380        assert_eq!(page.items.len(), 0);
1381        assert_eq!(page.total, 1);
1382        assert_eq!(page.page, 3);
1383        assert_eq!(page.per_page, 10);
1384
1385        Ok(())
1386    }
1387
1388    #[tokio::test]
1389    async fn test_find_by_relayer_id_pagination() -> Result<()> {
1390        let storage = TransactionRepositoryStorage::new_in_memory();
1391
1392        // Add many transactions for one relayer
1393        for i in 1..=10 {
1394            let tx = create_test_transaction(&format!("tx-{i}"), "test-relayer");
1395            storage.create(tx).await?;
1396        }
1397
1398        // Test first page
1399        let query = PaginationQuery {
1400            page: 1,
1401            per_page: 3,
1402        };
1403        let page1 = storage.find_by_relayer_id("test-relayer", query).await?;
1404        assert_eq!(page1.items.len(), 3);
1405        assert_eq!(page1.total, 10);
1406        assert_eq!(page1.page, 1);
1407        assert_eq!(page1.per_page, 3);
1408
1409        // Test second page
1410        let query = PaginationQuery {
1411            page: 2,
1412            per_page: 3,
1413        };
1414        let page2 = storage.find_by_relayer_id("test-relayer", query).await?;
1415        assert_eq!(page2.items.len(), 3);
1416        assert_eq!(page2.total, 10);
1417        assert_eq!(page2.page, 2);
1418        assert_eq!(page2.per_page, 3);
1419
1420        Ok(())
1421    }
1422
1423    #[tokio::test]
1424    async fn test_find_by_multiple_statuses() -> Result<()> {
1425        let storage = TransactionRepositoryStorage::new_in_memory();
1426
1427        // Add transactions with different statuses
1428        let tx1 =
1429            create_test_transaction_with_status("tx-1", "test-relayer", TransactionStatus::Pending);
1430        let tx2 =
1431            create_test_transaction_with_status("tx-2", "test-relayer", TransactionStatus::Sent);
1432        let tx3 = create_test_transaction_with_status(
1433            "tx-3",
1434            "test-relayer",
1435            TransactionStatus::Confirmed,
1436        );
1437        let tx4 =
1438            create_test_transaction_with_status("tx-4", "test-relayer", TransactionStatus::Failed);
1439
1440        storage.create(tx1).await?;
1441        storage.create(tx2).await?;
1442        storage.create(tx3).await?;
1443        storage.create(tx4).await?;
1444
1445        // Find by multiple statuses
1446        let statuses = vec![TransactionStatus::Pending, TransactionStatus::Sent];
1447        let result = storage.find_by_status("test-relayer", &statuses).await?;
1448
1449        assert_eq!(result.len(), 2);
1450
1451        // Verify all transactions have the correct statuses
1452        let found_statuses: Vec<TransactionStatus> =
1453            result.iter().map(|tx| tx.status.clone()).collect();
1454        assert!(found_statuses.contains(&TransactionStatus::Pending));
1455        assert!(found_statuses.contains(&TransactionStatus::Sent));
1456
1457        Ok(())
1458    }
1459
1460    #[tokio::test]
1461    async fn test_record_stellar_try_again_later_retry_in_memory() -> Result<()> {
1462        let storage = TransactionRepositoryStorage::new_in_memory();
1463        let mut transaction = create_test_transaction("test-tx", "test-relayer");
1464        transaction.status = TransactionStatus::Sent;
1465        storage.create(transaction).await?;
1466
1467        let sent_at = "2025-03-18T10:00:00Z".to_string();
1468        let updated = storage
1469            .record_stellar_try_again_later_retry("test-tx".to_string(), sent_at.clone())
1470            .await?;
1471
1472        assert_eq!(updated.id, "test-tx");
1473        assert_eq!(updated.sent_at, Some(sent_at));
1474        let meta = updated.metadata.expect("metadata should be set");
1475        assert_eq!(meta.try_again_later_retries, 1);
1476        assert_eq!(meta.consecutive_failures, 0);
1477        assert_eq!(meta.total_failures, 0);
1478        assert_eq!(meta.insufficient_fee_retries, 0);
1479
1480        Ok(())
1481    }
1482
1483    #[tokio::test]
1484    async fn test_record_stellar_try_again_later_retry_accumulates_in_memory() -> Result<()> {
1485        let storage = TransactionRepositoryStorage::new_in_memory();
1486        let mut transaction = create_test_transaction("test-tx", "test-relayer");
1487        transaction.status = TransactionStatus::Sent;
1488        storage.create(transaction).await?;
1489
1490        storage
1491            .record_stellar_try_again_later_retry(
1492                "test-tx".to_string(),
1493                "2025-03-18T10:00:00Z".to_string(),
1494            )
1495            .await?;
1496
1497        let updated = storage
1498            .record_stellar_try_again_later_retry(
1499                "test-tx".to_string(),
1500                "2025-03-18T10:01:00Z".to_string(),
1501            )
1502            .await?;
1503
1504        assert_eq!(updated.sent_at.as_deref(), Some("2025-03-18T10:01:00Z"));
1505        let meta = updated.metadata.unwrap();
1506        assert_eq!(meta.try_again_later_retries, 2);
1507
1508        Ok(())
1509    }
1510
1511    #[tokio::test]
1512    async fn test_record_stellar_try_again_later_retry_noop_on_final_state_in_memory() -> Result<()>
1513    {
1514        let storage = TransactionRepositoryStorage::new_in_memory();
1515        let mut transaction = create_test_transaction("test-tx", "test-relayer");
1516        transaction.status = TransactionStatus::Confirmed;
1517        transaction.sent_at = Some("old-time".to_string());
1518        storage.create(transaction).await?;
1519
1520        let result = storage
1521            .record_stellar_try_again_later_retry("test-tx".to_string(), "new-time".to_string())
1522            .await?;
1523
1524        assert_eq!(result.sent_at.as_deref(), Some("old-time"));
1525        assert!(result.metadata.is_none());
1526
1527        Ok(())
1528    }
1529
1530    #[tokio::test]
1531    async fn test_record_stellar_try_again_later_retry_not_found_in_memory() -> Result<()> {
1532        let storage = TransactionRepositoryStorage::new_in_memory();
1533
1534        let result = storage
1535            .record_stellar_try_again_later_retry(
1536                "nonexistent".to_string(),
1537                "2025-03-18T10:00:00Z".to_string(),
1538            )
1539            .await;
1540
1541        assert!(matches!(result, Err(RepositoryError::NotFound(_))));
1542
1543        Ok(())
1544    }
1545}