openzeppelin_relayer/repositories/transaction/
transaction_redis.rs

1//! Redis-backed implementation of the TransactionRepository.
2
3use crate::config::ServerConfig;
4use crate::constants::FINAL_TRANSACTION_STATUSES;
5use crate::domain::transaction::common::is_final_state;
6use crate::metrics::{
7    TRANSACTIONS_BY_STATUS, TRANSACTIONS_CREATED, TRANSACTIONS_FAILED,
8    TRANSACTIONS_INSUFFICIENT_FEE_FAILED, TRANSACTIONS_INSUFFICIENT_FEE_SUCCESS,
9    TRANSACTIONS_SUBMITTED, TRANSACTIONS_SUCCESS, TRANSACTIONS_TRY_AGAIN_LATER_FAILED,
10    TRANSACTIONS_TRY_AGAIN_LATER_SUCCESS, TRANSACTION_PROCESSING_TIME,
11};
12use crate::models::{
13    NetworkTransactionData, PaginationQuery, RepositoryError, TransactionRepoModel,
14    TransactionStatus, TransactionUpdateRequest,
15};
16use crate::repositories::redis_base::RedisRepository;
17use crate::repositories::{
18    BatchDeleteResult, BatchRetrievalResult, PaginatedResult, Repository, TransactionDeleteRequest,
19    TransactionRepository,
20};
21use crate::utils::RedisConnections;
22use async_trait::async_trait;
23use chrono::Utc;
24use redis::{AsyncCommands, Script};
25use std::fmt;
26use std::sync::Arc;
27use tracing::{debug, error, warn};
28
29const RELAYER_PREFIX: &str = "relayer";
30const TX_PREFIX: &str = "tx";
31const STATUS_PREFIX: &str = "status";
32const STATUS_SORTED_PREFIX: &str = "status_sorted";
33const NONCE_PREFIX: &str = "nonce";
34const TX_TO_RELAYER_PREFIX: &str = "tx_to_relayer";
35const RELAYER_LIST_KEY: &str = "relayer_list";
36const TX_BY_CREATED_AT_PREFIX: &str = "tx_by_created_at";
37
38#[derive(Clone)]
39pub struct RedisTransactionRepository {
40    pub connections: Arc<RedisConnections>,
41    pub key_prefix: String,
42}
43
44impl RedisRepository for RedisTransactionRepository {}
45
46impl RedisTransactionRepository {
47    pub fn new(
48        connections: Arc<RedisConnections>,
49        key_prefix: String,
50    ) -> Result<Self, RepositoryError> {
51        if key_prefix.is_empty() {
52            return Err(RepositoryError::InvalidData(
53                "Redis key prefix cannot be empty".to_string(),
54            ));
55        }
56
57        Ok(Self {
58            connections,
59            key_prefix,
60        })
61    }
62
63    /// Generate key for transaction data: relayer:{relayer_id}:tx:{tx_id}
64    fn tx_key(&self, relayer_id: &str, tx_id: &str) -> String {
65        format!(
66            "{}:{}:{}:{}:{}",
67            self.key_prefix, RELAYER_PREFIX, relayer_id, TX_PREFIX, tx_id
68        )
69    }
70
71    /// Generate key for reverse lookup: tx_to_relayer:{tx_id}
72    fn tx_to_relayer_key(&self, tx_id: &str) -> String {
73        format!(
74            "{}:{}:{}:{}",
75            self.key_prefix, RELAYER_PREFIX, TX_TO_RELAYER_PREFIX, tx_id
76        )
77    }
78
79    /// Generate key for relayer status index (legacy SET): relayer:{relayer_id}:status:{status}
80    fn relayer_status_key(&self, relayer_id: &str, status: &TransactionStatus) -> String {
81        format!(
82            "{}:{}:{}:{}:{}",
83            self.key_prefix, RELAYER_PREFIX, relayer_id, STATUS_PREFIX, status
84        )
85    }
86
87    /// Generate key for relayer status sorted index (SORTED SET): relayer:{relayer_id}:status_sorted:{status}
88    /// Score is created_at timestamp in milliseconds for efficient ordering.
89    fn relayer_status_sorted_key(&self, relayer_id: &str, status: &TransactionStatus) -> String {
90        format!(
91            "{}:{}:{}:{}:{}",
92            self.key_prefix, RELAYER_PREFIX, relayer_id, STATUS_SORTED_PREFIX, status
93        )
94    }
95
96    /// Generate key for relayer nonce index: relayer:{relayer_id}:nonce:{nonce}
97    fn relayer_nonce_key(&self, relayer_id: &str, nonce: u64) -> String {
98        format!(
99            "{}:{}:{}:{}:{}",
100            self.key_prefix, RELAYER_PREFIX, relayer_id, NONCE_PREFIX, nonce
101        )
102    }
103
104    /// Generate key for relayer list: relayer_list (set of all relayer IDs)
105    fn relayer_list_key(&self) -> String {
106        format!("{}:{}", self.key_prefix, RELAYER_LIST_KEY)
107    }
108
109    /// Generate key for relayer's sorted set by created_at: relayer:{relayer_id}:tx_by_created_at
110    fn relayer_tx_by_created_at_key(&self, relayer_id: &str) -> String {
111        format!(
112            "{}:{}:{}:{}",
113            self.key_prefix, RELAYER_PREFIX, relayer_id, TX_BY_CREATED_AT_PREFIX
114        )
115    }
116
117    /// Returns the components needed for Lua scripts to resolve a tx key from
118    /// only the tx_id: (tx_to_relayer lookup key, key prefix, key suffix).
119    /// The Lua script does: `GET KEYS[1]` to get the relayer_id, then
120    /// constructs the tx key as `ARGV[1] .. relayer_id .. ARGV[2]`.
121    fn tx_key_parts(&self, tx_id: &str) -> (String, String, String) {
122        let lookup_key = self.tx_to_relayer_key(tx_id);
123        let key_prefix = format!("{}:{}:", self.key_prefix, RELAYER_PREFIX);
124        let key_suffix = format!(":{TX_PREFIX}:{tx_id}");
125        (lookup_key, key_prefix, key_suffix)
126    }
127
128    /// Executes an atomic Lua script with retry/backoff for transient Redis failures.
129    ///
130    /// Every script receives `KEYS[1]` = tx_to_relayer lookup key and
131    /// `ARGV[1..2]` = key prefix/suffix. `extra_args` are appended as `ARGV[3..]`.
132    /// The script must return the (possibly updated) JSON string or `false` for
133    /// not-found.
134    async fn run_atomic_script(
135        &self,
136        lua: &str,
137        tx_id: &str,
138        extra_args: &[&str],
139        op_name: &str,
140    ) -> Result<TransactionRepoModel, RepositoryError> {
141        const MAX_RETRIES: u32 = 3;
142        const BASE_BACKOFF_MS: u64 = 100;
143
144        let (lookup_key, key_prefix, key_suffix) = self.tx_key_parts(tx_id);
145        let script = Script::new(lua);
146        let mut last_error = None;
147
148        for attempt in 0..MAX_RETRIES {
149            let backoff = BASE_BACKOFF_MS * 2u64.pow(attempt);
150
151            let mut conn = match self
152                .get_connection(self.connections.primary(), op_name)
153                .await
154            {
155                Ok(conn) => conn,
156                Err(e) => {
157                    last_error = Some(e);
158                    if attempt < MAX_RETRIES - 1 {
159                        warn!(tx_id = %tx_id, attempt, op = %op_name, "connection failed, retrying");
160                        tokio::time::sleep(tokio::time::Duration::from_millis(backoff)).await;
161                        continue;
162                    }
163                    return Err(last_error.unwrap());
164                }
165            };
166
167            let mut invocation = script.prepare_invoke();
168            invocation
169                .key(&lookup_key)
170                .arg(&key_prefix)
171                .arg(&key_suffix);
172            for arg in extra_args {
173                invocation.arg(*arg);
174            }
175
176            match invocation.invoke_async::<Option<String>>(&mut conn).await {
177                Ok(result) => {
178                    let json = result.ok_or_else(|| {
179                        RepositoryError::NotFound(format!("Transaction with ID {tx_id} not found"))
180                    })?;
181                    return self.deserialize_entity::<TransactionRepoModel>(
182                        &json,
183                        tx_id,
184                        "transaction",
185                    );
186                }
187                Err(e) => {
188                    last_error = Some(self.map_redis_error(e, op_name));
189                    if attempt < MAX_RETRIES - 1 {
190                        warn!(
191                            tx_id = %tx_id, attempt, op = %op_name,
192                            "atomic script failed, retrying"
193                        );
194                        tokio::time::sleep(tokio::time::Duration::from_millis(backoff)).await;
195                        continue;
196                    }
197                    return Err(last_error.unwrap());
198                }
199            }
200        }
201        Err(last_error.unwrap_or_else(|| {
202            RepositoryError::UnexpectedError(format!("retry loop exhausted for {op_name}"))
203        }))
204    }
205
206    /// Executes a Lua script with retry/backoff, returning a Vec<String> result
207    /// (for scripts that return Lua tables / multi-bulk replies).
208    /// Returns `Ok(None)` when the script returns `false`.
209    async fn run_script_with_retry_vec(
210        &self,
211        script: &Script,
212        lookup_key: &str,
213        key_prefix: &str,
214        key_suffix: &str,
215        extra_args: &[&str],
216        op_name: &str,
217    ) -> Result<Option<Vec<String>>, RepositoryError> {
218        const MAX_RETRIES: u32 = 3;
219        const BASE_BACKOFF_MS: u64 = 100;
220
221        let mut last_error = None;
222
223        for attempt in 0..MAX_RETRIES {
224            let backoff = BASE_BACKOFF_MS * 2u64.pow(attempt);
225
226            let mut conn = match self
227                .get_connection(self.connections.primary(), op_name)
228                .await
229            {
230                Ok(conn) => conn,
231                Err(e) => {
232                    last_error = Some(e);
233                    if attempt < MAX_RETRIES - 1 {
234                        warn!(op = %op_name, attempt, "connection failed, retrying");
235                        tokio::time::sleep(tokio::time::Duration::from_millis(backoff)).await;
236                        continue;
237                    }
238                    return Err(last_error.unwrap());
239                }
240            };
241
242            let mut invocation = script.prepare_invoke();
243            invocation.key(lookup_key).arg(key_prefix).arg(key_suffix);
244            for arg in extra_args {
245                invocation.arg(*arg);
246            }
247
248            // Redis returns `false` from Lua as a Nil bulk reply, which
249            // redis-rs maps to `None` for `Option<Vec<String>>`.
250            match invocation
251                .invoke_async::<Option<Vec<String>>>(&mut conn)
252                .await
253            {
254                Ok(result) => return Ok(result),
255                Err(e) => {
256                    last_error = Some(self.map_redis_error(e, op_name));
257                    if attempt < MAX_RETRIES - 1 {
258                        warn!(op = %op_name, attempt, "script failed, retrying");
259                        tokio::time::sleep(tokio::time::Duration::from_millis(backoff)).await;
260                        continue;
261                    }
262                    return Err(last_error.unwrap());
263                }
264            }
265        }
266        Err(last_error.unwrap_or_else(|| {
267            RepositoryError::UnexpectedError(format!("retry loop exhausted for {op_name}"))
268        }))
269    }
270
271    /// Parse timestamp string to score for sorted set (milliseconds since epoch)
272    fn timestamp_to_score(&self, timestamp: &str) -> f64 {
273        chrono::DateTime::parse_from_rfc3339(timestamp)
274            .map(|dt| dt.timestamp_millis() as f64)
275            .unwrap_or_else(|_| {
276                warn!(timestamp = %timestamp, "failed to parse timestamp, using 0");
277                0.0
278            })
279    }
280
281    /// Compute the appropriate score for a transaction's status sorted set.
282    /// - For Confirmed status: use confirmed_at (on-chain confirmation order)
283    /// - For all other statuses: use created_at (queue/processing order)
284    fn status_sorted_score(&self, tx: &TransactionRepoModel) -> f64 {
285        if tx.status == TransactionStatus::Confirmed {
286            // For Confirmed, prefer confirmed_at for accurate on-chain ordering
287            if let Some(ref confirmed_at) = tx.confirmed_at {
288                return self.timestamp_to_score(confirmed_at);
289            }
290            // Fallback to created_at if confirmed_at not set (shouldn't happen)
291            warn!(tx_id = %tx.id, "Confirmed transaction missing confirmed_at, using created_at");
292        }
293        self.timestamp_to_score(&tx.created_at)
294    }
295
296    /// Batch fetch transactions by IDs using reverse lookup
297    async fn get_transactions_by_ids(
298        &self,
299        ids: &[String],
300    ) -> Result<BatchRetrievalResult<TransactionRepoModel>, RepositoryError> {
301        if ids.is_empty() {
302            debug!("no transaction IDs provided for batch fetch");
303            return Ok(BatchRetrievalResult {
304                results: vec![],
305                failed_ids: vec![],
306            });
307        }
308
309        let mut conn = self
310            .get_connection(self.connections.reader(), "batch_fetch_transactions")
311            .await?;
312
313        let reverse_keys: Vec<String> = ids.iter().map(|id| self.tx_to_relayer_key(id)).collect();
314
315        debug!(count = %ids.len(), "fetching relayer IDs for transactions");
316
317        let relayer_ids: Vec<Option<String>> = conn
318            .mget(&reverse_keys)
319            .await
320            .map_err(|e| self.map_redis_error(e, "batch_fetch_relayer_ids"))?;
321
322        let mut tx_keys = Vec::new();
323        let mut valid_ids = Vec::new();
324        let mut failed_ids = Vec::new();
325        for (i, relayer_id) in relayer_ids.into_iter().enumerate() {
326            match relayer_id {
327                Some(relayer_id) => {
328                    tx_keys.push(self.tx_key(&relayer_id, &ids[i]));
329                    valid_ids.push(ids[i].clone());
330                }
331                None => {
332                    warn!(tx_id = %ids[i], "no relayer found for transaction");
333                    failed_ids.push(ids[i].clone());
334                }
335            }
336        }
337
338        if tx_keys.is_empty() {
339            debug!("no valid transactions found for batch fetch");
340            return Ok(BatchRetrievalResult {
341                results: vec![],
342                failed_ids,
343            });
344        }
345
346        debug!(count = %tx_keys.len(), "batch fetching transaction data");
347
348        let values: Vec<Option<String>> = conn
349            .mget(&tx_keys)
350            .await
351            .map_err(|e| self.map_redis_error(e, "batch_fetch_transactions"))?;
352
353        let mut transactions = Vec::new();
354        let mut failed_count = 0;
355        for (i, value) in values.into_iter().enumerate() {
356            match value {
357                Some(json) => {
358                    match self.deserialize_entity::<TransactionRepoModel>(
359                        &json,
360                        &valid_ids[i],
361                        "transaction",
362                    ) {
363                        Ok(tx) => transactions.push(tx),
364                        Err(e) => {
365                            failed_count += 1;
366                            error!(tx_id = %valid_ids[i], error = %e, "failed to deserialize transaction");
367                            // Continue processing other transactions
368                        }
369                    }
370                }
371                None => {
372                    warn!(tx_id = %valid_ids[i], "transaction not found in batch fetch");
373                    failed_ids.push(valid_ids[i].clone());
374                }
375            }
376        }
377
378        if failed_count > 0 {
379            warn!(failed_count = %failed_count, total_count = %valid_ids.len(), "failed to deserialize transactions in batch");
380        }
381
382        debug!(count = %transactions.len(), "successfully fetched transactions");
383        Ok(BatchRetrievalResult {
384            results: transactions,
385            failed_ids,
386        })
387    }
388
389    /// Extract nonce from EVM transaction data
390    fn extract_nonce(&self, network_data: &NetworkTransactionData) -> Option<u64> {
391        match network_data.get_evm_transaction_data() {
392            Ok(tx_data) => tx_data.nonce,
393            Err(_) => {
394                debug!("no EVM transaction data available for nonce extraction");
395                None
396            }
397        }
398    }
399
400    /// Ensures the status sorted set exists, migrating from legacy SET if needed.
401    ///
402    /// This handles the transition from unordered SETs to sorted SETs for status indexing.
403    /// If the sorted set is empty but the legacy set has data, it migrates the data
404    /// by looking up each transaction's created_at timestamp to compute the score.
405    ///
406    /// # Concurrency
407    /// This function is safe for concurrent calls. If multiple calls race to migrate
408    /// the same status set:
409    /// - ZADD is idempotent (same member + score = no-op)
410    /// - DEL on non-existent key is safe (returns 0)
411    /// - After first successful migration, subsequent calls hit the fast path (ZCARD > 0)
412    ///
413    /// The only downside of concurrent migrations is wasted work, not data corruption.
414    ///
415    /// Returns the count of items in the sorted set after migration.
416    async fn ensure_status_sorted_set(
417        &self,
418        relayer_id: &str,
419        status: &TransactionStatus,
420    ) -> Result<u64, RepositoryError> {
421        let sorted_key = self.relayer_status_sorted_key(relayer_id, status);
422        let legacy_key = self.relayer_status_key(relayer_id, status);
423
424        // Phase 1: Check if migration is needed
425        let legacy_ids = {
426            let mut conn = self
427                .get_connection(self.connections.primary(), "ensure_status_sorted_set_check")
428                .await?;
429
430            // Always check if legacy set has data that needs migration
431            let legacy_count: u64 = conn
432                .scard(&legacy_key)
433                .await
434                .map_err(|e| self.map_redis_error(e, "ensure_status_sorted_set_scard"))?;
435
436            if legacy_count == 0 {
437                // No legacy data to migrate, return current ZSET count
438                let sorted_count: u64 = conn
439                    .zcard(&sorted_key)
440                    .await
441                    .map_err(|e| self.map_redis_error(e, "ensure_status_sorted_set_zcard"))?;
442                return Ok(sorted_count);
443            }
444
445            // Migration needed: get all IDs from legacy set
446            debug!(
447                relayer_id = %relayer_id,
448                status = %status,
449                legacy_count = %legacy_count,
450                "migrating status set to sorted set"
451            );
452
453            let ids: Vec<String> = conn
454                .smembers(&legacy_key)
455                .await
456                .map_err(|e| self.map_redis_error(e, "ensure_status_sorted_set_smembers"))?;
457
458            ids
459            // Connection dropped here before nested call to avoid connection doubling
460        };
461
462        if legacy_ids.is_empty() {
463            return Ok(0);
464        }
465
466        // Phase 2: Fetch transactions (uses its own connection internally)
467        let transactions = self.get_transactions_by_ids(&legacy_ids).await?;
468
469        // Phase 3: Perform migration with a new connection
470        let mut conn = self
471            .get_connection(
472                self.connections.primary(),
473                "ensure_status_sorted_set_migrate",
474            )
475            .await?;
476
477        if transactions.results.is_empty() {
478            // All transactions were stale/deleted, clean up legacy set
479            let _: () = conn
480                .del(&legacy_key)
481                .await
482                .map_err(|e| self.map_redis_error(e, "ensure_status_sorted_set_del_stale"))?;
483            return Ok(0);
484        }
485
486        // Build sorted set entries and migrate atomically
487        // Use status-aware scoring: confirmed_at for Confirmed, created_at for others
488        let mut pipe = redis::pipe();
489        pipe.atomic();
490
491        for tx in &transactions.results {
492            let score = self.status_sorted_score(tx);
493            pipe.zadd(&sorted_key, &tx.id, score);
494        }
495
496        // Delete legacy set after migration
497        pipe.del(&legacy_key);
498
499        pipe.query_async::<()>(&mut conn)
500            .await
501            .map_err(|e| self.map_redis_error(e, "ensure_status_sorted_set_migrate"))?;
502
503        let migrated_count = transactions.results.len() as u64;
504        debug!(
505            relayer_id = %relayer_id,
506            status = %status,
507            migrated_count = %migrated_count,
508            "completed migration of status set to sorted set"
509        );
510
511        Ok(migrated_count)
512    }
513
514    /// Update indexes atomically with comprehensive error handling
515    async fn update_indexes(
516        &self,
517        tx: &TransactionRepoModel,
518        old_tx: Option<&TransactionRepoModel>,
519    ) -> Result<(), RepositoryError> {
520        let mut conn = self
521            .get_connection(self.connections.primary(), "update_indexes")
522            .await?;
523        let mut pipe = redis::pipe();
524        pipe.atomic();
525
526        debug!(tx_id = %tx.id, "updating indexes for transaction");
527
528        // Add relayer to the global relayer list
529        let relayer_list_key = self.relayer_list_key();
530        pipe.sadd(&relayer_list_key, &tx.relayer_id);
531
532        // Compute scores for sorted sets
533        // Status sorted set: uses confirmed_at for Confirmed status, created_at for others
534        let status_score = self.status_sorted_score(tx);
535        // Global tx_by_created_at: always uses created_at for consistent ordering
536        let created_at_score = self.timestamp_to_score(&tx.created_at);
537
538        // Handle status index updates - write to SORTED SET (new format)
539        let new_status_sorted_key = self.relayer_status_sorted_key(&tx.relayer_id, &tx.status);
540        pipe.zadd(&new_status_sorted_key, &tx.id, status_score);
541        debug!(tx_id = %tx.id, status = %tx.status, score = %status_score, "added transaction to status sorted set");
542
543        if let Some(nonce) = self.extract_nonce(&tx.network_data) {
544            let nonce_key = self.relayer_nonce_key(&tx.relayer_id, nonce);
545            pipe.set(&nonce_key, &tx.id);
546            debug!(tx_id = %tx.id, nonce = %nonce, "added nonce index for transaction");
547        }
548
549        // Add to per-relayer sorted set by created_at (for efficient sorted pagination)
550        let relayer_sorted_key = self.relayer_tx_by_created_at_key(&tx.relayer_id);
551        pipe.zadd(&relayer_sorted_key, &tx.id, created_at_score);
552        debug!(tx_id = %tx.id, score = %created_at_score, "added transaction to sorted set by created_at");
553
554        // Remove old indexes if updating
555        if let Some(old) = old_tx {
556            if old.status != tx.status {
557                // Remove from old status sorted set (new format)
558                let old_status_sorted_key =
559                    self.relayer_status_sorted_key(&old.relayer_id, &old.status);
560                pipe.zrem(&old_status_sorted_key, &tx.id);
561
562                // Also clean up legacy SET if it exists (for migration cleanup)
563                let old_status_legacy_key = self.relayer_status_key(&old.relayer_id, &old.status);
564                pipe.srem(&old_status_legacy_key, &tx.id);
565
566                debug!(tx_id = %tx.id, old_status = %old.status, new_status = %tx.status, "removing old status indexes for transaction");
567            }
568
569            // Handle nonce index cleanup
570            if let Some(old_nonce) = self.extract_nonce(&old.network_data) {
571                let new_nonce = self.extract_nonce(&tx.network_data);
572                if Some(old_nonce) != new_nonce {
573                    let old_nonce_key = self.relayer_nonce_key(&old.relayer_id, old_nonce);
574                    pipe.del(&old_nonce_key);
575                    debug!(tx_id = %tx.id, old_nonce = %old_nonce, new_nonce = ?new_nonce, "removing old nonce index for transaction");
576                }
577            }
578        }
579
580        // Execute all operations in a single pipeline
581        pipe.exec_async(&mut conn).await.map_err(|e| {
582            error!(tx_id = %tx.id, error = %e, "index update pipeline failed for transaction");
583            self.map_redis_error(e, &format!("update_indexes_for_tx_{}", tx.id))
584        })?;
585
586        debug!(tx_id = %tx.id, "successfully updated indexes for transaction");
587        Ok(())
588    }
589
590    /// Remove all indexes with error recovery
591    async fn remove_all_indexes(&self, tx: &TransactionRepoModel) -> Result<(), RepositoryError> {
592        let mut conn = self
593            .get_connection(self.connections.primary(), "remove_all_indexes")
594            .await?;
595        let mut pipe = redis::pipe();
596        pipe.atomic();
597
598        debug!(tx_id = %tx.id, "removing all indexes for transaction");
599
600        // Remove from ALL possible status indexes to ensure complete cleanup
601        // This handles cases where a transaction might be in multiple status sets
602        // due to race conditions, partial failures, or bugs
603        for status in &[
604            TransactionStatus::Canceled,
605            TransactionStatus::Pending,
606            TransactionStatus::Sent,
607            TransactionStatus::Submitted,
608            TransactionStatus::Mined,
609            TransactionStatus::Confirmed,
610            TransactionStatus::Failed,
611            TransactionStatus::Expired,
612        ] {
613            // Remove from sorted status set (new format)
614            let status_sorted_key = self.relayer_status_sorted_key(&tx.relayer_id, status);
615            pipe.zrem(&status_sorted_key, &tx.id);
616
617            // Remove from legacy status set (for migration cleanup)
618            let status_legacy_key = self.relayer_status_key(&tx.relayer_id, status);
619            pipe.srem(&status_legacy_key, &tx.id);
620        }
621
622        // Remove nonce index if exists
623        if let Some(nonce) = self.extract_nonce(&tx.network_data) {
624            let nonce_key = self.relayer_nonce_key(&tx.relayer_id, nonce);
625            pipe.del(&nonce_key);
626            debug!(tx_id = %tx.id, nonce = %nonce, "removing nonce index for transaction");
627        }
628
629        // Remove from per-relayer sorted set by created_at
630        let relayer_sorted_key = self.relayer_tx_by_created_at_key(&tx.relayer_id);
631        pipe.zrem(&relayer_sorted_key, &tx.id);
632        debug!(tx_id = %tx.id, "removing transaction from sorted set by created_at");
633
634        // Remove reverse lookup
635        let reverse_key = self.tx_to_relayer_key(&tx.id);
636        pipe.del(&reverse_key);
637
638        pipe.exec_async(&mut conn).await.map_err(|e| {
639            error!(tx_id = %tx.id, error = %e, "index removal failed for transaction");
640            self.map_redis_error(e, &format!("remove_indexes_for_tx_{}", tx.id))
641        })?;
642
643        debug!(tx_id = %tx.id, "successfully removed all indexes for transaction");
644        Ok(())
645    }
646
647    /// Track Prometheus metrics when a transaction status changes.
648    fn track_status_change_metrics(
649        &self,
650        _original_tx: &TransactionRepoModel,
651        updated_tx: &TransactionRepoModel,
652        old_status: &TransactionStatus,
653        new_status: &TransactionStatus,
654    ) {
655        let network_type = format!("{:?}", updated_tx.network_type).to_lowercase();
656        let relayer_id = updated_tx.relayer_id.as_str();
657
658        // Track submission (when status changes to Submitted)
659        if *old_status != TransactionStatus::Submitted
660            && *new_status == TransactionStatus::Submitted
661        {
662            TRANSACTIONS_SUBMITTED
663                .with_label_values(&[relayer_id, &network_type])
664                .inc();
665
666            if let Ok(created_time) = chrono::DateTime::parse_from_rfc3339(&updated_tx.created_at) {
667                let processing_seconds =
668                    (Utc::now() - created_time.with_timezone(&Utc)).num_seconds() as f64;
669                TRANSACTION_PROCESSING_TIME
670                    .with_label_values(&[relayer_id, &network_type, "creation_to_submission"])
671                    .observe(processing_seconds);
672            }
673        }
674
675        // Track status distribution (update gauge when status changes)
676        if old_status != new_status {
677            let old_status_str = format!("{old_status:?}").to_lowercase();
678            let old_status_gauge = TRANSACTIONS_BY_STATUS.with_label_values(&[
679                relayer_id,
680                &network_type,
681                &old_status_str,
682            ]);
683            let clamped_value = (old_status_gauge.get() - 1.0).max(0.0);
684            old_status_gauge.set(clamped_value);
685
686            let new_status_str = format!("{new_status:?}").to_lowercase();
687            TRANSACTIONS_BY_STATUS
688                .with_label_values(&[relayer_id, &network_type, &new_status_str])
689                .inc();
690        }
691
692        // Track metrics for final transaction states
693        let was_final = is_final_state(old_status);
694        let is_final = is_final_state(new_status);
695
696        if !was_final && is_final {
697            let previous_status = format!("{old_status:?}").to_lowercase();
698            let meta = updated_tx.metadata.as_ref();
699            let had_insufficient_fee = meta.is_some_and(|m| m.insufficient_fee_retries > 0);
700            let had_try_again_later = meta.is_some_and(|m| m.try_again_later_retries > 0);
701
702            match new_status {
703                TransactionStatus::Confirmed => {
704                    TRANSACTIONS_SUCCESS
705                        .with_label_values(&[relayer_id, &network_type])
706                        .inc();
707                    if had_insufficient_fee {
708                        TRANSACTIONS_INSUFFICIENT_FEE_SUCCESS
709                            .with_label_values(&[relayer_id, &network_type])
710                            .inc();
711                    }
712                    if had_try_again_later {
713                        TRANSACTIONS_TRY_AGAIN_LATER_SUCCESS
714                            .with_label_values(&[relayer_id, &network_type])
715                            .inc();
716                    }
717
718                    if let (Some(sent_at_str), Some(confirmed_at_str)) =
719                        (&updated_tx.sent_at, &updated_tx.confirmed_at)
720                    {
721                        if let (Ok(sent_time), Ok(confirmed_time)) = (
722                            chrono::DateTime::parse_from_rfc3339(sent_at_str),
723                            chrono::DateTime::parse_from_rfc3339(confirmed_at_str),
724                        ) {
725                            let processing_seconds = (confirmed_time.with_timezone(&Utc)
726                                - sent_time.with_timezone(&Utc))
727                            .num_seconds()
728                                as f64;
729                            TRANSACTION_PROCESSING_TIME
730                                .with_label_values(&[
731                                    relayer_id,
732                                    &network_type,
733                                    "submission_to_confirmation",
734                                ])
735                                .observe(processing_seconds);
736                        }
737                    }
738
739                    if let Ok(created_time) =
740                        chrono::DateTime::parse_from_rfc3339(&updated_tx.created_at)
741                    {
742                        if let Some(confirmed_at_str) = &updated_tx.confirmed_at {
743                            if let Ok(confirmed_time) =
744                                chrono::DateTime::parse_from_rfc3339(confirmed_at_str)
745                            {
746                                let processing_seconds = (confirmed_time.with_timezone(&Utc)
747                                    - created_time.with_timezone(&Utc))
748                                .num_seconds()
749                                    as f64;
750                                TRANSACTION_PROCESSING_TIME
751                                    .with_label_values(&[
752                                        relayer_id,
753                                        &network_type,
754                                        "creation_to_confirmation",
755                                    ])
756                                    .observe(processing_seconds);
757                            }
758                        }
759                    }
760                }
761                TransactionStatus::Failed => {
762                    let failure_reason = updated_tx
763                        .status_reason
764                        .as_deref()
765                        .map(|reason| {
766                            if reason.starts_with("Submission failed:") {
767                                "submission_failed"
768                            } else if reason.starts_with("Preparation failed:") {
769                                "preparation_failed"
770                            } else {
771                                "failed"
772                            }
773                        })
774                        .unwrap_or("failed");
775                    TRANSACTIONS_FAILED
776                        .with_label_values(&[
777                            relayer_id,
778                            &network_type,
779                            failure_reason,
780                            &previous_status,
781                        ])
782                        .inc();
783                }
784                TransactionStatus::Expired => {
785                    TRANSACTIONS_FAILED
786                        .with_label_values(&[
787                            relayer_id,
788                            &network_type,
789                            "expired",
790                            &previous_status,
791                        ])
792                        .inc();
793                }
794                TransactionStatus::Canceled => {
795                    TRANSACTIONS_FAILED
796                        .with_label_values(&[
797                            relayer_id,
798                            &network_type,
799                            "canceled",
800                            &previous_status,
801                        ])
802                        .inc();
803                }
804                _ => {}
805            }
806
807            // Track retry-related failure metrics for all non-success final states
808            if *new_status != TransactionStatus::Confirmed {
809                if had_insufficient_fee {
810                    TRANSACTIONS_INSUFFICIENT_FEE_FAILED
811                        .with_label_values(&[relayer_id, &network_type])
812                        .inc();
813                }
814                if had_try_again_later {
815                    TRANSACTIONS_TRY_AGAIN_LATER_FAILED
816                        .with_label_values(&[relayer_id, &network_type])
817                        .inc();
818                }
819            }
820        }
821    }
822}
823
824impl fmt::Debug for RedisTransactionRepository {
825    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
826        f.debug_struct("RedisTransactionRepository")
827            .field("connections", &"<RedisConnections>")
828            .field("key_prefix", &self.key_prefix)
829            .finish()
830    }
831}
832
833#[async_trait]
834impl Repository<TransactionRepoModel, String> for RedisTransactionRepository {
835    async fn create(
836        &self,
837        entity: TransactionRepoModel,
838    ) -> Result<TransactionRepoModel, RepositoryError> {
839        if entity.id.is_empty() {
840            return Err(RepositoryError::InvalidData(
841                "Transaction ID cannot be empty".to_string(),
842            ));
843        }
844
845        let key = self.tx_key(&entity.relayer_id, &entity.id);
846        let reverse_key = self.tx_to_relayer_key(&entity.id);
847        let mut conn = self
848            .get_connection(self.connections.primary(), "create")
849            .await?;
850
851        debug!(tx_id = %entity.id, "creating transaction");
852
853        let value = self.serialize_entity(&entity, |t| &t.id, "transaction")?;
854
855        // Check if transaction already exists by checking reverse lookup
856        let existing: Option<String> = conn
857            .get(&reverse_key)
858            .await
859            .map_err(|e| self.map_redis_error(e, "create_transaction_check"))?;
860
861        if existing.is_some() {
862            return Err(RepositoryError::ConstraintViolation(format!(
863                "Transaction with ID {} already exists",
864                entity.id
865            )));
866        }
867
868        // Use atomic pipeline for consistency
869        let mut pipe = redis::pipe();
870        pipe.atomic();
871        pipe.set(&key, &value);
872        pipe.set(&reverse_key, &entity.relayer_id);
873
874        pipe.exec_async(&mut conn)
875            .await
876            .map_err(|e| self.map_redis_error(e, "create_transaction"))?;
877
878        // Update indexes separately to handle partial failures gracefully
879        if let Err(e) = self.update_indexes(&entity, None).await {
880            error!(tx_id = %entity.id, error = %e, "failed to update indexes for new transaction");
881            return Err(e);
882        }
883
884        // Track transaction creation metric
885        let network_type = format!("{:?}", entity.network_type).to_lowercase();
886        let relayer_id = entity.relayer_id.as_str();
887        TRANSACTIONS_CREATED
888            .with_label_values(&[relayer_id, &network_type])
889            .inc();
890
891        // Track initial status distribution (Pending)
892        let status = &entity.status;
893        let status_str = format!("{status:?}").to_lowercase();
894        TRANSACTIONS_BY_STATUS
895            .with_label_values(&[relayer_id, &network_type, &status_str])
896            .inc();
897
898        debug!(tx_id = %entity.id, "successfully created transaction");
899        Ok(entity)
900    }
901
902    async fn get_by_id(&self, id: String) -> Result<TransactionRepoModel, RepositoryError> {
903        if id.is_empty() {
904            return Err(RepositoryError::InvalidData(
905                "Transaction ID cannot be empty".to_string(),
906            ));
907        }
908
909        let mut conn = self
910            .get_connection(self.connections.reader(), "get_by_id")
911            .await?;
912
913        debug!(tx_id = %id, "fetching transaction");
914
915        let reverse_key = self.tx_to_relayer_key(&id);
916        let relayer_id: Option<String> = conn
917            .get(&reverse_key)
918            .await
919            .map_err(|e| self.map_redis_error(e, "get_transaction_reverse_lookup"))?;
920
921        let relayer_id = match relayer_id {
922            Some(relayer_id) => relayer_id,
923            None => {
924                debug!(tx_id = %id, "transaction not found (no reverse lookup)");
925                return Err(RepositoryError::NotFound(format!(
926                    "Transaction with ID {id} not found"
927                )));
928            }
929        };
930
931        let key = self.tx_key(&relayer_id, &id);
932        let value: Option<String> = conn
933            .get(&key)
934            .await
935            .map_err(|e| self.map_redis_error(e, "get_transaction_by_id"))?;
936
937        match value {
938            Some(json) => {
939                let tx =
940                    self.deserialize_entity::<TransactionRepoModel>(&json, &id, "transaction")?;
941                debug!(tx_id = %id, "successfully fetched transaction");
942                Ok(tx)
943            }
944            None => {
945                debug!(tx_id = %id, "transaction not found");
946                Err(RepositoryError::NotFound(format!(
947                    "Transaction with ID {id} not found"
948                )))
949            }
950        }
951    }
952
953    // Unoptimized implementation of list_paginated. Rarely used. find_by_relayer_id is preferred.
954    async fn list_all(&self) -> Result<Vec<TransactionRepoModel>, RepositoryError> {
955        let mut conn = self
956            .get_connection(self.connections.reader(), "list_all")
957            .await?;
958
959        debug!("fetching all transactions sorted by created_at (newest first)");
960
961        // Get all relayer IDs
962        let relayer_list_key = self.relayer_list_key();
963        let relayer_ids: Vec<String> = conn
964            .smembers(&relayer_list_key)
965            .await
966            .map_err(|e| self.map_redis_error(e, "list_all_relayer_ids"))?;
967
968        debug!(count = %relayer_ids.len(), "found relayers");
969
970        // Collect all transaction IDs from all relayers using their sorted sets
971        let mut all_tx_ids = Vec::new();
972        for relayer_id in relayer_ids {
973            let relayer_sorted_key = self.relayer_tx_by_created_at_key(&relayer_id);
974            let tx_ids: Vec<String> = redis::cmd("ZRANGE")
975                .arg(&relayer_sorted_key)
976                .arg(0)
977                .arg(-1)
978                .arg("REV")
979                .query_async(&mut conn)
980                .await
981                .map_err(|e| self.map_redis_error(e, "list_all_relayer_sorted"))?;
982
983            all_tx_ids.extend(tx_ids);
984        }
985
986        // Release connection before nested call to avoid connection doubling
987        drop(conn);
988
989        // Batch fetch all transactions at once
990        let batch_result = self.get_transactions_by_ids(&all_tx_ids).await?;
991        let mut all_transactions = batch_result.results;
992
993        // Sort all transactions by created_at (newest first)
994        all_transactions.sort_by(|a, b| b.created_at.cmp(&a.created_at));
995
996        debug!(count = %all_transactions.len(), "found transactions");
997        Ok(all_transactions)
998    }
999
1000    // Unoptimized implementation of list_paginated. Rarely used. find_by_relayer_id is preferred.
1001    async fn list_paginated(
1002        &self,
1003        query: PaginationQuery,
1004    ) -> Result<PaginatedResult<TransactionRepoModel>, RepositoryError> {
1005        if query.per_page == 0 {
1006            return Err(RepositoryError::InvalidData(
1007                "per_page must be greater than 0".to_string(),
1008            ));
1009        }
1010
1011        let mut conn = self
1012            .get_connection(self.connections.reader(), "list_paginated")
1013            .await?;
1014
1015        debug!(page = %query.page, per_page = %query.per_page, "fetching paginated transactions sorted by created_at (newest first)");
1016
1017        // Get all relayer IDs
1018        let relayer_list_key = self.relayer_list_key();
1019        let relayer_ids: Vec<String> = conn
1020            .smembers(&relayer_list_key)
1021            .await
1022            .map_err(|e| self.map_redis_error(e, "list_paginated_relayer_ids"))?;
1023
1024        // Collect all transaction IDs from all relayers using their sorted sets
1025        let mut all_tx_ids = Vec::new();
1026        for relayer_id in relayer_ids {
1027            let relayer_sorted_key = self.relayer_tx_by_created_at_key(&relayer_id);
1028            let tx_ids: Vec<String> = redis::cmd("ZRANGE")
1029                .arg(&relayer_sorted_key)
1030                .arg(0)
1031                .arg(-1)
1032                .arg("REV")
1033                .query_async(&mut conn)
1034                .await
1035                .map_err(|e| self.map_redis_error(e, "list_paginated_relayer_sorted"))?;
1036
1037            all_tx_ids.extend(tx_ids);
1038        }
1039
1040        // Release connection before nested call to avoid connection doubling
1041        drop(conn);
1042
1043        // Batch fetch all transactions at once
1044        let batch_result = self.get_transactions_by_ids(&all_tx_ids).await?;
1045        let mut all_transactions = batch_result.results;
1046
1047        // Sort all transactions by created_at (newest first)
1048        all_transactions.sort_by(|a, b| b.created_at.cmp(&a.created_at));
1049
1050        let total = all_transactions.len() as u64;
1051        let start = ((query.page - 1) * query.per_page) as usize;
1052        let end = (start + query.per_page as usize).min(all_transactions.len());
1053
1054        if start >= all_transactions.len() {
1055            debug!(page = %query.page, total = %total, "page is beyond available data");
1056            return Ok(PaginatedResult {
1057                items: vec![],
1058                total,
1059                page: query.page,
1060                per_page: query.per_page,
1061            });
1062        }
1063
1064        let items = all_transactions[start..end].to_vec();
1065
1066        debug!(count = %items.len(), page = %query.page, "successfully fetched transactions for page");
1067
1068        Ok(PaginatedResult {
1069            items,
1070            total,
1071            page: query.page,
1072            per_page: query.per_page,
1073        })
1074    }
1075
1076    async fn update(
1077        &self,
1078        id: String,
1079        entity: TransactionRepoModel,
1080    ) -> Result<TransactionRepoModel, RepositoryError> {
1081        if id.is_empty() {
1082            return Err(RepositoryError::InvalidData(
1083                "Transaction ID cannot be empty".to_string(),
1084            ));
1085        }
1086
1087        debug!(tx_id = %id, "updating transaction");
1088
1089        // Get the old transaction for index cleanup
1090        let old_tx = self.get_by_id(id.clone()).await?;
1091
1092        let key = self.tx_key(&entity.relayer_id, &id);
1093        let mut conn = self
1094            .get_connection(self.connections.primary(), "update")
1095            .await?;
1096
1097        let value = self.serialize_entity(&entity, |t| &t.id, "transaction")?;
1098
1099        // Update transaction
1100        let _: () = conn
1101            .set(&key, value)
1102            .await
1103            .map_err(|e| self.map_redis_error(e, "update_transaction"))?;
1104
1105        // Update indexes
1106        self.update_indexes(&entity, Some(&old_tx)).await?;
1107
1108        debug!(tx_id = %id, "successfully updated transaction");
1109        Ok(entity)
1110    }
1111
1112    async fn delete_by_id(&self, id: String) -> Result<(), RepositoryError> {
1113        if id.is_empty() {
1114            return Err(RepositoryError::InvalidData(
1115                "Transaction ID cannot be empty".to_string(),
1116            ));
1117        }
1118
1119        debug!(tx_id = %id, "deleting transaction");
1120
1121        // Get transaction first for index cleanup
1122        let tx = self.get_by_id(id.clone()).await?;
1123
1124        let key = self.tx_key(&tx.relayer_id, &id);
1125        let reverse_key = self.tx_to_relayer_key(&id);
1126        let mut conn = self
1127            .get_connection(self.connections.primary(), "delete_by_id")
1128            .await?;
1129
1130        let mut pipe = redis::pipe();
1131        pipe.atomic();
1132        pipe.del(&key);
1133        pipe.del(&reverse_key);
1134
1135        pipe.exec_async(&mut conn)
1136            .await
1137            .map_err(|e| self.map_redis_error(e, "delete_transaction"))?;
1138
1139        // Remove indexes (log errors but don't fail the delete)
1140        if let Err(e) = self.remove_all_indexes(&tx).await {
1141            error!(tx_id = %id, error = %e, "failed to remove indexes for deleted transaction");
1142        }
1143
1144        debug!(tx_id = %id, "successfully deleted transaction");
1145        Ok(())
1146    }
1147
1148    // Unoptimized implementation of count. Rarely used. find_by_relayer_id is preferred.
1149    async fn count(&self) -> Result<usize, RepositoryError> {
1150        let mut conn = self
1151            .get_connection(self.connections.reader(), "count")
1152            .await?;
1153
1154        debug!("counting transactions");
1155
1156        // Get all relayer IDs and sum their sorted set counts
1157        let relayer_list_key = self.relayer_list_key();
1158        let relayer_ids: Vec<String> = conn
1159            .smembers(&relayer_list_key)
1160            .await
1161            .map_err(|e| self.map_redis_error(e, "count_relayer_ids"))?;
1162
1163        let mut total_count = 0usize;
1164        for relayer_id in relayer_ids {
1165            let relayer_sorted_key = self.relayer_tx_by_created_at_key(&relayer_id);
1166            let count: usize = conn
1167                .zcard(&relayer_sorted_key)
1168                .await
1169                .map_err(|e| self.map_redis_error(e, "count_relayer_transactions"))?;
1170            total_count += count;
1171        }
1172
1173        debug!(count = %total_count, "transaction count");
1174        Ok(total_count)
1175    }
1176
1177    async fn has_entries(&self) -> Result<bool, RepositoryError> {
1178        let mut conn = self
1179            .get_connection(self.connections.reader(), "has_entries")
1180            .await?;
1181        let relayer_list_key = self.relayer_list_key();
1182
1183        debug!("checking if transaction entries exist");
1184
1185        let exists: bool = conn
1186            .exists(&relayer_list_key)
1187            .await
1188            .map_err(|e| self.map_redis_error(e, "has_entries_check"))?;
1189
1190        debug!(exists = %exists, "transaction entries exist");
1191        Ok(exists)
1192    }
1193
1194    async fn drop_all_entries(&self) -> Result<(), RepositoryError> {
1195        let mut conn = self
1196            .get_connection(self.connections.primary(), "drop_all_entries")
1197            .await?;
1198        let relayer_list_key = self.relayer_list_key();
1199
1200        debug!("dropping all transaction entries");
1201
1202        // Get all relayer IDs first
1203        let relayer_ids: Vec<String> = conn
1204            .smembers(&relayer_list_key)
1205            .await
1206            .map_err(|e| self.map_redis_error(e, "drop_all_entries_get_relayer_ids"))?;
1207
1208        if relayer_ids.is_empty() {
1209            debug!("no transaction entries to drop");
1210            return Ok(());
1211        }
1212
1213        // Use pipeline for atomic operations
1214        let mut pipe = redis::pipe();
1215        pipe.atomic();
1216
1217        // Delete all transactions and their indexes for each relayer
1218        for relayer_id in &relayer_ids {
1219            // Get all transaction IDs for this relayer
1220            let pattern = format!(
1221                "{}:{}:{}:{}:*",
1222                self.key_prefix, RELAYER_PREFIX, relayer_id, TX_PREFIX
1223            );
1224            let mut cursor = 0;
1225            let mut tx_ids = Vec::new();
1226
1227            loop {
1228                let (next_cursor, keys): (u64, Vec<String>) = redis::cmd("SCAN")
1229                    .cursor_arg(cursor)
1230                    .arg("MATCH")
1231                    .arg(&pattern)
1232                    .query_async(&mut conn)
1233                    .await
1234                    .map_err(|e| self.map_redis_error(e, "drop_all_entries_scan"))?;
1235
1236                // Extract transaction IDs from keys and delete keys
1237                for key in keys {
1238                    pipe.del(&key);
1239                    if let Some(tx_id) = key.split(':').next_back() {
1240                        tx_ids.push(tx_id.to_string());
1241                    }
1242                }
1243
1244                cursor = next_cursor;
1245                if cursor == 0 {
1246                    break;
1247                }
1248            }
1249
1250            // Delete reverse lookup keys and indexes
1251            for tx_id in tx_ids {
1252                let reverse_key = self.tx_to_relayer_key(&tx_id);
1253                pipe.del(&reverse_key);
1254
1255                // Delete status indexes (we can't know the specific status, so we'll clean up all possible ones)
1256                // This ensures complete cleanup even if there are orphaned entries
1257                for status in &[
1258                    TransactionStatus::Canceled,
1259                    TransactionStatus::Pending,
1260                    TransactionStatus::Sent,
1261                    TransactionStatus::Submitted,
1262                    TransactionStatus::Mined,
1263                    TransactionStatus::Confirmed,
1264                    TransactionStatus::Failed,
1265                    TransactionStatus::Expired,
1266                ] {
1267                    // Remove from sorted status set (new format)
1268                    let status_sorted_key = self.relayer_status_sorted_key(relayer_id, status);
1269                    pipe.zrem(&status_sorted_key, &tx_id);
1270
1271                    // Remove from legacy status set (for migration cleanup)
1272                    let status_key = self.relayer_status_key(relayer_id, status);
1273                    pipe.srem(&status_key, &tx_id);
1274                }
1275            }
1276
1277            // Delete the relayer's sorted set by created_at
1278            let relayer_sorted_key = self.relayer_tx_by_created_at_key(relayer_id);
1279            pipe.del(&relayer_sorted_key);
1280        }
1281
1282        // Delete the relayer list key
1283        pipe.del(&relayer_list_key);
1284
1285        pipe.exec_async(&mut conn)
1286            .await
1287            .map_err(|e| self.map_redis_error(e, "drop_all_entries_pipeline"))?;
1288
1289        debug!(count = %relayer_ids.len(), "dropped all transaction entries for relayers");
1290        Ok(())
1291    }
1292}
1293
1294#[async_trait]
1295impl TransactionRepository for RedisTransactionRepository {
1296    async fn find_by_relayer_id(
1297        &self,
1298        relayer_id: &str,
1299        query: PaginationQuery,
1300    ) -> Result<PaginatedResult<TransactionRepoModel>, RepositoryError> {
1301        let mut conn = self
1302            .get_connection(self.connections.reader(), "find_by_relayer_id")
1303            .await?;
1304
1305        debug!(relayer_id = %relayer_id, page = %query.page, per_page = %query.per_page, "fetching transactions for relayer sorted by created_at (newest first)");
1306
1307        let relayer_sorted_key = self.relayer_tx_by_created_at_key(relayer_id);
1308
1309        // Get total count from relayer's sorted set
1310        let sorted_set_count: u64 = conn
1311            .zcard(&relayer_sorted_key)
1312            .await
1313            .map_err(|e| self.map_redis_error(e, "find_by_relayer_id_count"))?;
1314
1315        // If sorted set is empty, return empty result immediately
1316        // All new transactions are automatically added to the sorted set
1317        if sorted_set_count == 0 {
1318            debug!(relayer_id = %relayer_id, "no transactions found for relayer (sorted set is empty)");
1319            return Ok(PaginatedResult {
1320                items: vec![],
1321                total: 0,
1322                page: query.page,
1323                per_page: query.per_page,
1324            });
1325        }
1326
1327        let total = sorted_set_count;
1328
1329        // Calculate pagination range (0-indexed for Redis ZRANGE with REV)
1330        let start = ((query.page - 1) * query.per_page) as isize;
1331        let end = start + query.per_page as isize - 1;
1332
1333        if start as u64 >= total {
1334            debug!(relayer_id = %relayer_id, page = %query.page, total = %total, "page is beyond available data");
1335            return Ok(PaginatedResult {
1336                items: vec![],
1337                total,
1338                page: query.page,
1339                per_page: query.per_page,
1340            });
1341        }
1342
1343        // Get page of transaction IDs from sorted set (newest first using ZRANGE with REV)
1344        let page_ids: Vec<String> = redis::cmd("ZRANGE")
1345            .arg(&relayer_sorted_key)
1346            .arg(start)
1347            .arg(end)
1348            .arg("REV")
1349            .query_async(&mut conn)
1350            .await
1351            .map_err(|e| self.map_redis_error(e, "find_by_relayer_id_sorted"))?;
1352
1353        // Release connection before nested call to avoid connection doubling
1354        drop(conn);
1355
1356        let items = self.get_transactions_by_ids(&page_ids).await?;
1357
1358        debug!(relayer_id = %relayer_id, count = %items.results.len(), page = %query.page, "successfully fetched transactions for relayer");
1359
1360        Ok(PaginatedResult {
1361            items: items.results,
1362            total,
1363            page: query.page,
1364            per_page: query.per_page,
1365        })
1366    }
1367
1368    // Unoptimized implementation of find_by_status. Rarely used. find_by_status_paginated is preferred.
1369    async fn find_by_status(
1370        &self,
1371        relayer_id: &str,
1372        statuses: &[TransactionStatus],
1373    ) -> Result<Vec<TransactionRepoModel>, RepositoryError> {
1374        // Ensure all status sorted sets are migrated first (releases connection after each)
1375        for status in statuses {
1376            self.ensure_status_sorted_set(relayer_id, status).await?;
1377        }
1378
1379        // Now get a connection and collect all IDs
1380        let mut conn = self
1381            .get_connection(self.connections.reader(), "find_by_status")
1382            .await?;
1383
1384        let mut all_ids: Vec<String> = Vec::new();
1385        for status in statuses {
1386            // Get IDs from sorted set (already ordered by created_at)
1387            let sorted_key = self.relayer_status_sorted_key(relayer_id, status);
1388            let ids: Vec<String> = redis::cmd("ZRANGE")
1389                .arg(&sorted_key)
1390                .arg(0)
1391                .arg(-1)
1392                .arg("REV") // Newest first
1393                .query_async(&mut conn)
1394                .await
1395                .map_err(|e| self.map_redis_error(e, "find_by_status"))?;
1396
1397            all_ids.extend(ids);
1398        }
1399
1400        // Release connection before nested call to avoid connection doubling
1401        drop(conn);
1402
1403        if all_ids.is_empty() {
1404            return Ok(vec![]);
1405        }
1406
1407        // Remove duplicates (can happen if a transaction is in multiple status sets due to partial failures)
1408        all_ids.sort();
1409        all_ids.dedup();
1410
1411        // Fetch all transactions and sort by created_at (newest first)
1412        let mut transactions = self.get_transactions_by_ids(&all_ids).await?;
1413
1414        // Sort by created_at descending (newest first)
1415        transactions
1416            .results
1417            .sort_by(|a, b| b.created_at.cmp(&a.created_at));
1418
1419        Ok(transactions.results)
1420    }
1421
1422    async fn find_by_status_paginated(
1423        &self,
1424        relayer_id: &str,
1425        statuses: &[TransactionStatus],
1426        query: PaginationQuery,
1427        oldest_first: bool,
1428    ) -> Result<PaginatedResult<TransactionRepoModel>, RepositoryError> {
1429        // Ensure all status sorted sets are migrated first (releases connection after each)
1430        for status in statuses {
1431            self.ensure_status_sorted_set(relayer_id, status).await?;
1432        }
1433
1434        let mut conn = self
1435            .get_connection(self.connections.reader(), "find_by_status_paginated")
1436            .await?;
1437
1438        // For single status, we can paginate directly from the sorted set
1439        if statuses.len() == 1 {
1440            let sorted_key = self.relayer_status_sorted_key(relayer_id, &statuses[0]);
1441
1442            // Get total count
1443            let total: u64 = conn
1444                .zcard(&sorted_key)
1445                .await
1446                .map_err(|e| self.map_redis_error(e, "find_by_status_paginated_count"))?;
1447
1448            if total == 0 {
1449                return Ok(PaginatedResult {
1450                    items: vec![],
1451                    total: 0,
1452                    page: query.page,
1453                    per_page: query.per_page,
1454                });
1455            }
1456
1457            // Calculate pagination bounds
1458            let start = ((query.page.saturating_sub(1)) * query.per_page) as isize;
1459            let end = start + query.per_page as isize - 1;
1460
1461            // Get page of IDs directly from sorted set
1462            // REV = newest first (descending), no REV = oldest first (ascending)
1463            let mut cmd = redis::cmd("ZRANGE");
1464            cmd.arg(&sorted_key).arg(start).arg(end);
1465            if !oldest_first {
1466                cmd.arg("REV");
1467            }
1468            let page_ids: Vec<String> = cmd
1469                .query_async(&mut conn)
1470                .await
1471                .map_err(|e| self.map_redis_error(e, "find_by_status_paginated"))?;
1472
1473            // Release connection before nested call to avoid connection doubling
1474            drop(conn);
1475
1476            let transactions = self.get_transactions_by_ids(&page_ids).await?;
1477
1478            debug!(
1479                relayer_id = %relayer_id,
1480                status = %statuses[0],
1481                total = %total,
1482                page = %query.page,
1483                page_size = %transactions.results.len(),
1484                "fetched paginated transactions by single status"
1485            );
1486
1487            return Ok(PaginatedResult {
1488                items: transactions.results,
1489                total,
1490                page: query.page,
1491                per_page: query.per_page,
1492            });
1493        }
1494
1495        // For multiple statuses, collect all IDs and merge
1496        let mut all_ids: Vec<(String, f64)> = Vec::new();
1497        for status in statuses {
1498            let sorted_key = self.relayer_status_sorted_key(relayer_id, status);
1499
1500            // Get IDs with scores for proper sorting
1501            let ids_with_scores: Vec<(String, f64)> = redis::cmd("ZRANGE")
1502                .arg(&sorted_key)
1503                .arg(0)
1504                .arg(-1)
1505                .arg("WITHSCORES")
1506                .query_async(&mut conn)
1507                .await
1508                .map_err(|e| self.map_redis_error(e, "find_by_status_paginated_multi"))?;
1509
1510            all_ids.extend(ids_with_scores);
1511        }
1512
1513        // Release connection before nested call to avoid connection doubling
1514        drop(conn);
1515
1516        // Remove duplicates (keep highest/lowest score based on sort order)
1517        let mut id_map: std::collections::HashMap<String, f64> = std::collections::HashMap::new();
1518        for (id, score) in all_ids {
1519            id_map
1520                .entry(id)
1521                .and_modify(|s| {
1522                    // For oldest_first, keep the lowest score; otherwise keep highest
1523                    if oldest_first {
1524                        if score < *s {
1525                            *s = score
1526                        }
1527                    } else if score > *s {
1528                        *s = score
1529                    }
1530                })
1531                .or_insert(score);
1532        }
1533
1534        // Sort by score: descending for newest first, ascending for oldest first
1535        let mut sorted_ids: Vec<(String, f64)> = id_map.into_iter().collect();
1536        if oldest_first {
1537            sorted_ids.sort_by(|a, b| a.1.partial_cmp(&b.1).unwrap_or(std::cmp::Ordering::Equal));
1538        } else {
1539            sorted_ids.sort_by(|a, b| b.1.partial_cmp(&a.1).unwrap_or(std::cmp::Ordering::Equal));
1540        }
1541
1542        let total = sorted_ids.len() as u64;
1543
1544        if total == 0 {
1545            return Ok(PaginatedResult {
1546                items: vec![],
1547                total: 0,
1548                page: query.page,
1549                per_page: query.per_page,
1550            });
1551        }
1552
1553        // Apply pagination
1554        let start = ((query.page.saturating_sub(1)) * query.per_page) as usize;
1555        let page_ids: Vec<String> = sorted_ids
1556            .into_iter()
1557            .skip(start)
1558            .take(query.per_page as usize)
1559            .map(|(id, _)| id)
1560            .collect();
1561
1562        // Fetch only the transactions for this page
1563        let transactions = self.get_transactions_by_ids(&page_ids).await?;
1564
1565        debug!(
1566            relayer_id = %relayer_id,
1567            total = %total,
1568            page = %query.page,
1569            page_size = %transactions.results.len(),
1570            "fetched paginated transactions by status"
1571        );
1572
1573        Ok(PaginatedResult {
1574            items: transactions.results,
1575            total,
1576            page: query.page,
1577            per_page: query.per_page,
1578        })
1579    }
1580
1581    async fn find_by_status_paginated_filtered(
1582        &self,
1583        relayer_id: &str,
1584        statuses: &[TransactionStatus],
1585        query: PaginationQuery,
1586        oldest_first: bool,
1587        exclude_canceled: bool,
1588    ) -> Result<PaginatedResult<TransactionRepoModel>, RepositoryError> {
1589        if !exclude_canceled {
1590            return self
1591                .find_by_status_paginated(relayer_id, statuses, query, oldest_first)
1592                .await;
1593        }
1594
1595        // Excluding cancelled-in-progress transactions means removing them BEFORE
1596        // counting and paginating, which the per-status sorted-set fast path (zcard +
1597        // ZRANGE) cannot do. So for this listing-only path we load all matching
1598        // transactions, drop the cancelled ones, then sort/paginate in memory. The hot
1599        // internal callers use the unfiltered fast method above.
1600        for status in statuses {
1601            self.ensure_status_sorted_set(relayer_id, status).await?;
1602        }
1603
1604        let mut conn = self
1605            .get_connection(
1606                self.connections.reader(),
1607                "find_by_status_paginated_filtered",
1608            )
1609            .await?;
1610
1611        // Collect all ids (with scores) across the requested statuses.
1612        let mut all_ids: Vec<(String, f64)> = Vec::new();
1613        for status in statuses {
1614            let sorted_key = self.relayer_status_sorted_key(relayer_id, status);
1615            let ids_with_scores: Vec<(String, f64)> = redis::cmd("ZRANGE")
1616                .arg(&sorted_key)
1617                .arg(0)
1618                .arg(-1)
1619                .arg("WITHSCORES")
1620                .query_async(&mut conn)
1621                .await
1622                .map_err(|e| self.map_redis_error(e, "find_by_status_paginated_filtered"))?;
1623            all_ids.extend(ids_with_scores);
1624        }
1625
1626        drop(conn);
1627
1628        // Dedup across statuses, keeping the extreme score for the sort direction.
1629        let mut id_map: std::collections::HashMap<String, f64> = std::collections::HashMap::new();
1630        for (id, score) in all_ids {
1631            id_map
1632                .entry(id)
1633                .and_modify(|s| {
1634                    if oldest_first {
1635                        if score < *s {
1636                            *s = score
1637                        }
1638                    } else if score > *s {
1639                        *s = score
1640                    }
1641                })
1642                .or_insert(score);
1643        }
1644
1645        let mut sorted_ids: Vec<(String, f64)> = id_map.into_iter().collect();
1646        if oldest_first {
1647            sorted_ids.sort_by(|a, b| a.1.partial_cmp(&b.1).unwrap_or(std::cmp::Ordering::Equal));
1648        } else {
1649            sorted_ids.sort_by(|a, b| b.1.partial_cmp(&a.1).unwrap_or(std::cmp::Ordering::Equal));
1650        }
1651        let ordered_ids: Vec<String> = sorted_ids.into_iter().map(|(id, _)| id).collect();
1652
1653        // Load all matching transactions, then filter + paginate, preserving sort order.
1654        let fetched = self.get_transactions_by_ids(&ordered_ids).await?;
1655        let mut by_id: std::collections::HashMap<String, TransactionRepoModel> = fetched
1656            .results
1657            .into_iter()
1658            .map(|t| (t.id.clone(), t))
1659            .collect();
1660
1661        let filtered: Vec<TransactionRepoModel> = ordered_ids
1662            .iter()
1663            .filter_map(|id| by_id.remove(id))
1664            .filter(|t| t.is_canceled != Some(true))
1665            .collect();
1666
1667        let total = filtered.len() as u64;
1668        let start = ((query.page.saturating_sub(1)) * query.per_page) as usize;
1669        let items: Vec<TransactionRepoModel> = filtered
1670            .into_iter()
1671            .skip(start)
1672            .take(query.per_page as usize)
1673            .collect();
1674
1675        Ok(PaginatedResult {
1676            items,
1677            total,
1678            page: query.page,
1679            per_page: query.per_page,
1680        })
1681    }
1682
1683    async fn find_by_nonce(
1684        &self,
1685        relayer_id: &str,
1686        nonce: u64,
1687    ) -> Result<Option<TransactionRepoModel>, RepositoryError> {
1688        let mut conn = self
1689            .get_connection(self.connections.reader(), "find_by_nonce")
1690            .await?;
1691        let nonce_key = self.relayer_nonce_key(relayer_id, nonce);
1692
1693        // Get transaction ID with this nonce for this relayer (should be single value)
1694        let tx_id: Option<String> = conn
1695            .get(nonce_key)
1696            .await
1697            .map_err(|e| self.map_redis_error(e, "find_by_nonce"))?;
1698
1699        match tx_id {
1700            Some(tx_id) => {
1701                match self.get_by_id(tx_id.clone()).await {
1702                    Ok(tx) => Ok(Some(tx)),
1703                    Err(RepositoryError::NotFound(_)) => {
1704                        // Transaction was deleted but index wasn't cleaned up
1705                        warn!(relayer_id = %relayer_id, nonce = %nonce, "stale nonce index found for relayer");
1706                        Ok(None)
1707                    }
1708                    Err(e) => Err(e),
1709                }
1710            }
1711            None => Ok(None),
1712        }
1713    }
1714
1715    async fn get_nonce_occupancy(
1716        &self,
1717        relayer_id: &str,
1718        from_nonce: u64,
1719        to_nonce: u64,
1720    ) -> Result<Vec<(u64, Option<TransactionStatus>)>, RepositoryError> {
1721        if from_nonce >= to_nonce {
1722            return Ok(vec![]);
1723        }
1724
1725        let nonces: Vec<u64> = (from_nonce..to_nonce).collect();
1726        let nonce_keys: Vec<String> = nonces
1727            .iter()
1728            .map(|n| self.relayer_nonce_key(relayer_id, *n))
1729            .collect();
1730
1731        // Phase 1: MGET nonce keys → tx_ids (single round trip)
1732        // Uses primary to avoid replica lag fabricating false gaps.
1733        let mut conn = self
1734            .get_connection(self.connections.primary(), "get_nonce_occupancy")
1735            .await?;
1736        let tx_ids: Vec<Option<String>> = redis::cmd("MGET")
1737            .arg(&nonce_keys)
1738            .query_async(&mut conn)
1739            .await
1740            .map_err(|e| self.map_redis_error(e, "get_nonce_occupancy:mget_nonces"))?;
1741
1742        // Build tx data keys for non-None slots. We know the relayer_id so we
1743        // skip the reverse lookup and go straight to the data key.
1744        let mut tx_key_entries: Vec<(usize, String)> = Vec::new();
1745        for (i, tx_id) in tx_ids.iter().enumerate() {
1746            if let Some(id) = tx_id {
1747                tx_key_entries.push((i, self.tx_key(relayer_id, id)));
1748            }
1749        }
1750
1751        // Phase 2: MGET tx data keys → JSON blobs (single round trip)
1752        let tx_statuses: Vec<Option<TransactionStatus>> = if tx_key_entries.is_empty() {
1753            vec![]
1754        } else {
1755            let data_keys: Vec<&str> = tx_key_entries.iter().map(|(_, k)| k.as_str()).collect();
1756            let raw_values: Vec<Option<String>> = redis::cmd("MGET")
1757                .arg(&data_keys)
1758                .query_async(&mut conn)
1759                .await
1760                .map_err(|e| self.map_redis_error(e, "get_nonce_occupancy:mget_txs"))?;
1761
1762            raw_values
1763                .into_iter()
1764                .enumerate()
1765                .map(|(i, v)| {
1766                    v.and_then(|json| {
1767                        match serde_json::from_str::<TransactionRepoModel>(&json) {
1768                            Ok(tx) => Some(tx.status),
1769                            Err(e) => {
1770                                let nonce = tx_key_entries.get(i).map(|(idx, _)| nonces[*idx]);
1771                                warn!(
1772                                    relayer_id = %relayer_id,
1773                                    nonce = ?nonce,
1774                                    error = %e,
1775                                    "get_nonce_occupancy: failed to deserialize transaction, treating as empty"
1776                                );
1777                                None
1778                            }
1779                        }
1780                    })
1781                })
1782                .collect()
1783        };
1784
1785        // Assemble results
1786        let mut results: Vec<(u64, Option<TransactionStatus>)> =
1787            nonces.iter().map(|n| (*n, None)).collect();
1788
1789        for (idx, (original_idx, _)) in tx_key_entries.iter().enumerate() {
1790            if let Some(status) = tx_statuses.get(idx).and_then(|s| s.clone()) {
1791                results[*original_idx].1 = Some(status);
1792            }
1793        }
1794
1795        Ok(results)
1796    }
1797
1798    async fn update_status(
1799        &self,
1800        tx_id: String,
1801        status: TransactionStatus,
1802    ) -> Result<TransactionRepoModel, RepositoryError> {
1803        let update = TransactionUpdateRequest {
1804            status: Some(status),
1805            ..Default::default()
1806        };
1807        self.partial_update(tx_id, update).await
1808    }
1809
1810    async fn partial_update(
1811        &self,
1812        tx_id: String,
1813        update: TransactionUpdateRequest,
1814    ) -> Result<TransactionRepoModel, RepositoryError> {
1815        // Serialize only the non-None fields as a JSON patch.
1816        let patch_json = serde_json::to_string(&update).map_err(|e| {
1817            RepositoryError::InvalidData(format!("Failed to serialize update patch: {e}"))
1818        })?;
1819
1820        // If the update sets a final status, compute delete_at in Rust (depends on server config)
1821        // and include it in the patch so the Lua script applies it atomically.
1822        let delete_at_value = if let Some(ref status) = update.status {
1823            if FINAL_TRANSACTION_STATUSES.contains(status) {
1824                let expiration_hours = ServerConfig::get_transaction_expiration_hours();
1825                let seconds = (expiration_hours * 3600.0) as i64;
1826                let delete_time = Utc::now() + chrono::Duration::seconds(seconds);
1827                Some(delete_time.to_rfc3339())
1828            } else {
1829                None
1830            }
1831        } else {
1832            None
1833        };
1834        let delete_at_arg = delete_at_value.as_deref().unwrap_or("");
1835
1836        let (lookup_key, key_prefix, key_suffix) = self.tx_key_parts(&tx_id);
1837
1838        // Lua script: atomically applies a JSON patch to the stored transaction.
1839        // Guards: rejects status changes on already-finalized transactions.
1840        // Returns a two-element array {old_json, new_json} so Rust has the full
1841        // pre-update state for index cleanup and metrics.
1842        // Returns false if tx not found.
1843        let patch_script = Script::new(
1844            r#"
1845            local relayer_id = redis.call('GET', KEYS[1])
1846            if not relayer_id then return false end
1847
1848            local tx_key = ARGV[1] .. relayer_id .. ARGV[2]
1849            local current = redis.call('GET', tx_key)
1850            if not current then return false end
1851
1852            local tx = cjson.decode(current)
1853            local patch = cjson.decode(ARGV[3])
1854
1855            -- Guard: reject status changes on finalized transactions.
1856            -- A stale worker must not resurrect a tx that another worker
1857            -- already moved to a terminal state.
1858            local final_states = {confirmed=true, failed=true, expired=true, canceled=true}
1859            if final_states[tx["status"]] and patch["status"] then
1860                return {current, current}
1861            end
1862
1863            local old_snapshot = current
1864
1865            -- lua-cjson cannot distinguish empty Lua tables from empty
1866            -- arrays, so a decode/encode round-trip turns [] into {}.
1867            -- Record which keys held [] in the stored doc and the patch
1868            -- so we can restore them after cjson.encode.
1869            -- NOTE: this relies on each array-typed field having a unique key
1870            -- name across the entire JSON document (including nested objects).
1871            -- If the model ever introduces duplicate key names at different
1872            -- nesting levels (e.g. metadata.hashes), the gsub below could
1873            -- restore the wrong occurrence.
1874            local empty_arrs = {}
1875            for k in string.gmatch(current, '"([^"]+)"%s*:%s*%[%s*%]') do
1876                empty_arrs[k] = true
1877            end
1878            for k in string.gmatch(ARGV[3], '"([^"]+)"%s*:%s*%[%s*%]') do
1879                empty_arrs[k] = true
1880            end
1881
1882            for k, v in pairs(patch) do
1883                tx[k] = v
1884            end
1885
1886            -- Apply delete_at if transitioning to a final state and not already set
1887            if ARGV[4] ~= '' and (not tx["delete_at"] or tx["delete_at"] == cjson.null) then
1888                tx["delete_at"] = ARGV[4]
1889            end
1890
1891            local updated = cjson.encode(tx)
1892
1893            -- Restore empty arrays that cjson.encode converted to {}
1894            for k, _ in pairs(empty_arrs) do
1895                updated = string.gsub(
1896                    updated, '"'..k..'"%s*:%s*{}', '"'..k..'":[]', 1
1897                )
1898            end
1899
1900            redis.call('SET', tx_key, updated)
1901            return {old_snapshot, updated}
1902            "#,
1903        );
1904
1905        let result: Option<Vec<String>> = self
1906            .run_script_with_retry_vec(
1907                &patch_script,
1908                &lookup_key,
1909                &key_prefix,
1910                &key_suffix,
1911                &[&patch_json, delete_at_arg],
1912                "partial_update",
1913            )
1914            .await?;
1915
1916        let parts = result.ok_or_else(|| {
1917            RepositoryError::NotFound(format!("Transaction with ID {tx_id} not found"))
1918        })?;
1919
1920        if parts.len() != 2 {
1921            return Err(RepositoryError::UnexpectedError(format!(
1922                "partial_update script returned {} elements, expected 2",
1923                parts.len()
1924            )));
1925        }
1926
1927        let old_json = &parts[0];
1928        let new_json = &parts[1];
1929
1930        let original_tx =
1931            self.deserialize_entity::<TransactionRepoModel>(old_json, &tx_id, "transaction")?;
1932        let updated_tx =
1933            self.deserialize_entity::<TransactionRepoModel>(new_json, &tx_id, "transaction")?;
1934
1935        // Update indexes using the full pre-update state (status, network_data, nonce, etc.)
1936        self.update_indexes(&updated_tx, Some(&original_tx)).await?;
1937
1938        debug!(tx_id = %tx_id, "successfully updated transaction via patch");
1939
1940        // Track metrics only when the persisted status actually changed.
1941        // The Lua script may silently reject a status patch on already-final
1942        // transactions, so we compare the deserialized before/after states.
1943        if original_tx.status != updated_tx.status {
1944            self.track_status_change_metrics(
1945                &original_tx,
1946                &updated_tx,
1947                &original_tx.status,
1948                &updated_tx.status,
1949            );
1950        }
1951
1952        Ok(updated_tx)
1953    }
1954
1955    async fn update_network_data(
1956        &self,
1957        tx_id: String,
1958        network_data: NetworkTransactionData,
1959    ) -> Result<TransactionRepoModel, RepositoryError> {
1960        let update = TransactionUpdateRequest {
1961            network_data: Some(network_data),
1962            ..Default::default()
1963        };
1964        self.partial_update(tx_id, update).await
1965    }
1966
1967    async fn set_sent_at(
1968        &self,
1969        tx_id: String,
1970        sent_at: String,
1971    ) -> Result<TransactionRepoModel, RepositoryError> {
1972        let update = TransactionUpdateRequest {
1973            sent_at: Some(sent_at),
1974            ..Default::default()
1975        };
1976        self.partial_update(tx_id, update).await
1977    }
1978
1979    async fn increment_status_check_failures(
1980        &self,
1981        tx_id: String,
1982    ) -> Result<TransactionRepoModel, RepositoryError> {
1983        self.run_atomic_script(
1984            r#"
1985            local function set_obj(json, key, tbl)
1986                local enc = cjson.encode(tbl)
1987                local r, n = string.gsub(json, '"'..key..'"%s*:%s*%b{}', '"'..key..'":'..enc, 1)
1988                if n > 0 then return r end
1989                r, n = string.gsub(json, '"'..key..'"%s*:%s*null', '"'..key..'":'..enc, 1)
1990                if n > 0 then return r end
1991                return string.gsub(json, '}%s*$', ',"'..key..'":'..enc..'}', 1)
1992            end
1993
1994            local relayer_id = redis.call('GET', KEYS[1])
1995            if not relayer_id then return false end
1996
1997            local tx_key = ARGV[1] .. relayer_id .. ARGV[2]
1998            local current = redis.call('GET', tx_key)
1999            if not current then return false end
2000
2001            local tx = cjson.decode(current)
2002            local final_states = {confirmed=true, failed=true, expired=true, canceled=true}
2003            if final_states[tx["status"]] then return current end
2004
2005            local metadata = tx["metadata"]
2006            if type(metadata) ~= 'table' then metadata = {} end
2007            metadata["consecutive_failures"] = (metadata["consecutive_failures"] or 0) + 1
2008            metadata["total_failures"] = (metadata["total_failures"] or 0) + 1
2009
2010            local updated = set_obj(current, "metadata", metadata)
2011            redis.call('SET', tx_key, updated)
2012            return updated
2013            "#,
2014            &tx_id,
2015            &[],
2016            "increment_status_check_failures",
2017        )
2018        .await
2019    }
2020
2021    async fn reset_status_check_consecutive_failures(
2022        &self,
2023        tx_id: String,
2024    ) -> Result<TransactionRepoModel, RepositoryError> {
2025        self.run_atomic_script(
2026            r#"
2027            local function set_obj(json, key, tbl)
2028                local enc = cjson.encode(tbl)
2029                local r, n = string.gsub(json, '"'..key..'"%s*:%s*%b{}', '"'..key..'":'..enc, 1)
2030                if n > 0 then return r end
2031                r, n = string.gsub(json, '"'..key..'"%s*:%s*null', '"'..key..'":'..enc, 1)
2032                if n > 0 then return r end
2033                return string.gsub(json, '}%s*$', ',"'..key..'":'..enc..'}', 1)
2034            end
2035
2036            local relayer_id = redis.call('GET', KEYS[1])
2037            if not relayer_id then return false end
2038
2039            local tx_key = ARGV[1] .. relayer_id .. ARGV[2]
2040            local current = redis.call('GET', tx_key)
2041            if not current then return false end
2042
2043            local tx = cjson.decode(current)
2044            local final_states = {confirmed=true, failed=true, expired=true, canceled=true}
2045            if final_states[tx["status"]] then return current end
2046
2047            local metadata = tx["metadata"]
2048            if type(metadata) ~= 'table' then metadata = {} end
2049            metadata["consecutive_failures"] = 0
2050
2051            local updated = set_obj(current, "metadata", metadata)
2052            redis.call('SET', tx_key, updated)
2053            return updated
2054            "#,
2055            &tx_id,
2056            &[],
2057            "reset_status_check_consecutive_failures",
2058        )
2059        .await
2060    }
2061
2062    async fn record_stellar_insufficient_fee_retry(
2063        &self,
2064        tx_id: String,
2065        sent_at: String,
2066    ) -> Result<TransactionRepoModel, RepositoryError> {
2067        self.run_atomic_script(
2068            r#"
2069            local function set_str(json, key, val)
2070                local enc = cjson.encode(val)
2071                local r, n = string.gsub(json, '"'..key..'"%s*:%s*"[^"]*"', '"'..key..'":'..enc, 1)
2072                if n > 0 then return r end
2073                r, n = string.gsub(json, '"'..key..'"%s*:%s*null', '"'..key..'":'..enc, 1)
2074                if n > 0 then return r end
2075                return string.gsub(json, '}%s*$', ',"'..key..'":'..enc..'}', 1)
2076            end
2077            local function set_obj(json, key, tbl)
2078                local enc = cjson.encode(tbl)
2079                local r, n = string.gsub(json, '"'..key..'"%s*:%s*%b{}', '"'..key..'":'..enc, 1)
2080                if n > 0 then return r end
2081                r, n = string.gsub(json, '"'..key..'"%s*:%s*null', '"'..key..'":'..enc, 1)
2082                if n > 0 then return r end
2083                return string.gsub(json, '}%s*$', ',"'..key..'":'..enc..'}', 1)
2084            end
2085
2086            local relayer_id = redis.call('GET', KEYS[1])
2087            if not relayer_id then return false end
2088
2089            local tx_key = ARGV[1] .. relayer_id .. ARGV[2]
2090            local current = redis.call('GET', tx_key)
2091            if not current then return false end
2092
2093            local tx = cjson.decode(current)
2094            local final_states = {confirmed=true, failed=true, expired=true, canceled=true}
2095            if final_states[tx["status"]] then return current end
2096
2097            local metadata = tx["metadata"]
2098            if type(metadata) ~= 'table' then metadata = {} end
2099            metadata["insufficient_fee_retries"] = (metadata["insufficient_fee_retries"] or 0) + 1
2100
2101            local updated = set_str(current, "sent_at", ARGV[3])
2102            updated = set_obj(updated, "metadata", metadata)
2103            redis.call('SET', tx_key, updated)
2104            return updated
2105            "#,
2106            &tx_id,
2107            &[&sent_at],
2108            "record_stellar_insufficient_fee_retry",
2109        )
2110        .await
2111    }
2112
2113    async fn record_stellar_try_again_later_retry(
2114        &self,
2115        tx_id: String,
2116        sent_at: String,
2117    ) -> Result<TransactionRepoModel, RepositoryError> {
2118        self.run_atomic_script(
2119            r#"
2120            local function set_str(json, key, val)
2121                local enc = cjson.encode(val)
2122                local r, n = string.gsub(json, '"'..key..'"%s*:%s*"[^"]*"', '"'..key..'":'..enc, 1)
2123                if n > 0 then return r end
2124                r, n = string.gsub(json, '"'..key..'"%s*:%s*null', '"'..key..'":'..enc, 1)
2125                if n > 0 then return r end
2126                return string.gsub(json, '}%s*$', ',"'..key..'":'..enc..'}', 1)
2127            end
2128            local function set_obj(json, key, tbl)
2129                local enc = cjson.encode(tbl)
2130                local r, n = string.gsub(json, '"'..key..'"%s*:%s*%b{}', '"'..key..'":'..enc, 1)
2131                if n > 0 then return r end
2132                r, n = string.gsub(json, '"'..key..'"%s*:%s*null', '"'..key..'":'..enc, 1)
2133                if n > 0 then return r end
2134                return string.gsub(json, '}%s*$', ',"'..key..'":'..enc..'}', 1)
2135            end
2136
2137            local relayer_id = redis.call('GET', KEYS[1])
2138            if not relayer_id then return false end
2139
2140            local tx_key = ARGV[1] .. relayer_id .. ARGV[2]
2141            local current = redis.call('GET', tx_key)
2142            if not current then return false end
2143
2144            local tx = cjson.decode(current)
2145            local final_states = {confirmed=true, failed=true, expired=true, canceled=true}
2146            if final_states[tx["status"]] then return current end
2147
2148            local metadata = tx["metadata"]
2149            if type(metadata) ~= 'table' then metadata = {} end
2150            metadata["try_again_later_retries"] = (metadata["try_again_later_retries"] or 0) + 1
2151
2152            local updated = set_str(current, "sent_at", ARGV[3])
2153            updated = set_obj(updated, "metadata", metadata)
2154            redis.call('SET', tx_key, updated)
2155            return updated
2156            "#,
2157            &tx_id,
2158            &[&sent_at],
2159            "record_stellar_try_again_later_retry",
2160        )
2161        .await
2162    }
2163
2164    async fn set_confirmed_at(
2165        &self,
2166        tx_id: String,
2167        confirmed_at: String,
2168    ) -> Result<TransactionRepoModel, RepositoryError> {
2169        let update = TransactionUpdateRequest {
2170            confirmed_at: Some(confirmed_at),
2171            ..Default::default()
2172        };
2173        self.partial_update(tx_id, update).await
2174    }
2175
2176    /// Count transactions by status using Redis ZCARD (O(1) per sorted set).
2177    /// Much more efficient than find_by_status when you only need the count.
2178    /// Triggers migration from legacy SETs if needed.
2179    async fn count_by_status(
2180        &self,
2181        relayer_id: &str,
2182        statuses: &[TransactionStatus],
2183    ) -> Result<u64, RepositoryError> {
2184        let mut conn = self
2185            .get_connection(self.connections.reader(), "count_by_status")
2186            .await?;
2187        let mut total_count: u64 = 0;
2188
2189        for status in statuses {
2190            // Ensure sorted set is migrated
2191            self.ensure_status_sorted_set(relayer_id, status).await?;
2192
2193            let sorted_key = self.relayer_status_sorted_key(relayer_id, status);
2194            let count: u64 = conn
2195                .zcard(&sorted_key)
2196                .await
2197                .map_err(|e| self.map_redis_error(e, "count_by_status"))?;
2198            total_count += count;
2199        }
2200
2201        debug!(relayer_id = %relayer_id, count = %total_count, "counted transactions by status");
2202        Ok(total_count)
2203    }
2204
2205    async fn delete_by_ids(&self, ids: Vec<String>) -> Result<BatchDeleteResult, RepositoryError> {
2206        if ids.is_empty() {
2207            debug!("no transaction IDs provided for batch delete");
2208            return Ok(BatchDeleteResult::default());
2209        }
2210
2211        debug!(count = %ids.len(), "batch deleting transactions by IDs (with fetch)");
2212
2213        // Fetch transactions to get their data for index cleanup
2214        let batch_result = self.get_transactions_by_ids(&ids).await?;
2215
2216        // Convert to delete requests
2217        let requests: Vec<TransactionDeleteRequest> = batch_result
2218            .results
2219            .iter()
2220            .map(|tx| TransactionDeleteRequest {
2221                id: tx.id.clone(),
2222                relayer_id: tx.relayer_id.clone(),
2223                nonce: self.extract_nonce(&tx.network_data),
2224            })
2225            .collect();
2226
2227        // Track IDs that weren't found
2228        let mut result = self.delete_by_requests(requests).await?;
2229
2230        // Add the IDs that weren't found during fetch
2231        for id in batch_result.failed_ids {
2232            result
2233                .failed
2234                .push((id.clone(), format!("Transaction with ID {id} not found")));
2235        }
2236
2237        Ok(result)
2238    }
2239
2240    async fn delete_by_requests(
2241        &self,
2242        requests: Vec<TransactionDeleteRequest>,
2243    ) -> Result<BatchDeleteResult, RepositoryError> {
2244        if requests.is_empty() {
2245            debug!("no delete requests provided for batch delete");
2246            return Ok(BatchDeleteResult::default());
2247        }
2248
2249        debug!(count = %requests.len(), "batch deleting transactions by requests (no fetch)");
2250        let mut conn = self
2251            .get_connection(self.connections.primary(), "batch_delete_no_fetch")
2252            .await?;
2253        let mut pipe = redis::pipe();
2254        pipe.atomic();
2255
2256        // All possible statuses for index cleanup
2257        let all_statuses = [
2258            TransactionStatus::Canceled,
2259            TransactionStatus::Pending,
2260            TransactionStatus::Sent,
2261            TransactionStatus::Submitted,
2262            TransactionStatus::Mined,
2263            TransactionStatus::Confirmed,
2264            TransactionStatus::Failed,
2265            TransactionStatus::Expired,
2266        ];
2267
2268        // Build pipeline for all deletions and index removals
2269        for req in &requests {
2270            // Delete transaction data
2271            let tx_key = self.tx_key(&req.relayer_id, &req.id);
2272            pipe.del(&tx_key);
2273
2274            // Delete reverse lookup
2275            let reverse_key = self.tx_to_relayer_key(&req.id);
2276            pipe.del(&reverse_key);
2277
2278            // Remove from all possible status indexes
2279            for status in &all_statuses {
2280                let status_sorted_key = self.relayer_status_sorted_key(&req.relayer_id, status);
2281                pipe.zrem(&status_sorted_key, &req.id);
2282
2283                let status_legacy_key = self.relayer_status_key(&req.relayer_id, status);
2284                pipe.srem(&status_legacy_key, &req.id);
2285            }
2286
2287            // Remove nonce index if exists
2288            if let Some(nonce) = req.nonce {
2289                let nonce_key = self.relayer_nonce_key(&req.relayer_id, nonce);
2290                pipe.del(&nonce_key);
2291            }
2292
2293            // Remove from per-relayer sorted set by created_at
2294            let relayer_sorted_key = self.relayer_tx_by_created_at_key(&req.relayer_id);
2295            pipe.zrem(&relayer_sorted_key, &req.id);
2296        }
2297
2298        // Execute the entire pipeline in one round-trip
2299        match pipe.exec_async(&mut conn).await {
2300            Ok(_) => {
2301                let deleted_count = requests.len();
2302                debug!(
2303                    deleted_count = %deleted_count,
2304                    "batch delete completed"
2305                );
2306                Ok(BatchDeleteResult {
2307                    deleted_count,
2308                    failed: vec![],
2309                })
2310            }
2311            Err(e) => {
2312                error!(error = %e, "batch delete pipeline failed");
2313                // Mark all requests as failed
2314                let failed: Vec<(String, String)> = requests
2315                    .iter()
2316                    .map(|req| (req.id.clone(), format!("Redis pipeline error: {e}")))
2317                    .collect();
2318                Ok(BatchDeleteResult {
2319                    deleted_count: 0,
2320                    failed,
2321                })
2322            }
2323        }
2324    }
2325}
2326
2327#[cfg(test)]
2328mod tests {
2329    use super::*;
2330    use crate::models::{
2331        evm::Speed, EvmTransactionData, MemoSpec, NetworkType, StellarTransactionData,
2332        TransactionInput,
2333    };
2334    use alloy::primitives::U256;
2335    use deadpool_redis::{Config, Runtime};
2336    use lazy_static::lazy_static;
2337    use std::str::FromStr;
2338    use tokio;
2339    use uuid::Uuid;
2340
2341    use tokio::sync::Mutex;
2342
2343    // Use a mutex to ensure tests don't run in parallel when modifying env vars
2344    lazy_static! {
2345        static ref ENV_MUTEX: Mutex<()> = Mutex::new(());
2346    }
2347
2348    // Helper function to create test transactions
2349    fn create_test_transaction(id: &str) -> TransactionRepoModel {
2350        TransactionRepoModel {
2351            id: id.to_string(),
2352            relayer_id: "relayer-1".to_string(),
2353            status: TransactionStatus::Pending,
2354            status_reason: None,
2355            created_at: "2025-01-27T15:31:10.777083+00:00".to_string(),
2356            sent_at: Some("2025-01-27T15:31:10.777083+00:00".to_string()),
2357            confirmed_at: Some("2025-01-27T15:31:10.777083+00:00".to_string()),
2358            valid_until: None,
2359            delete_at: None,
2360            network_type: NetworkType::Evm,
2361            priced_at: None,
2362            hashes: vec![],
2363            network_data: NetworkTransactionData::Evm(EvmTransactionData {
2364                gas_price: Some(1000000000),
2365                gas_limit: Some(21000),
2366                nonce: Some(1),
2367                value: U256::from_str("1000000000000000000").unwrap(),
2368                data: Some("0x".to_string()),
2369                from: "0xSender".to_string(),
2370                to: Some("0xRecipient".to_string()),
2371                chain_id: 1,
2372                signature: None,
2373                hash: Some(format!("0x{id}")),
2374                speed: Some(Speed::Fast),
2375                max_fee_per_gas: None,
2376                max_priority_fee_per_gas: None,
2377                raw: None,
2378            }),
2379            noop_count: None,
2380            is_canceled: Some(false),
2381            metadata: None,
2382        }
2383    }
2384
2385    fn create_test_transaction_with_relayer(id: &str, relayer_id: &str) -> TransactionRepoModel {
2386        let mut tx = create_test_transaction(id);
2387        tx.relayer_id = relayer_id.to_string();
2388        tx
2389    }
2390
2391    fn create_stellar_test_transaction(
2392        id: &str,
2393        relayer_id: &str,
2394        sequence_number: i64,
2395        max_fee: i64,
2396        memo_id: u64,
2397    ) -> TransactionRepoModel {
2398        let mut tx = create_test_transaction_with_relayer(id, relayer_id);
2399        tx.network_type = NetworkType::Stellar;
2400        tx.network_data = NetworkTransactionData::Stellar(StellarTransactionData {
2401            source_account: "GAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAWHF".to_string(),
2402            fee: Some(100),
2403            sequence_number: Some(sequence_number),
2404            memo: Some(MemoSpec::Id { value: memo_id }),
2405            valid_until: None,
2406            network_passphrase: "Test SDF Network ; September 2015".to_string(),
2407            signatures: vec![],
2408            hash: None,
2409            simulation_transaction_data: None,
2410            transaction_input: TransactionInput::SignedXdr {
2411                xdr: "AAAA".to_string(),
2412                max_fee,
2413            },
2414            signed_envelope_xdr: None,
2415            transaction_result_xdr: None,
2416        });
2417        tx
2418    }
2419
2420    fn create_test_transaction_with_status(
2421        id: &str,
2422        relayer_id: &str,
2423        status: TransactionStatus,
2424    ) -> TransactionRepoModel {
2425        let mut tx = create_test_transaction_with_relayer(id, relayer_id);
2426        tx.status = status;
2427        tx
2428    }
2429
2430    fn create_test_transaction_with_nonce(
2431        id: &str,
2432        nonce: u64,
2433        relayer_id: &str,
2434    ) -> TransactionRepoModel {
2435        let mut tx = create_test_transaction_with_relayer(id, relayer_id);
2436        if let NetworkTransactionData::Evm(ref mut evm_data) = tx.network_data {
2437            evm_data.nonce = Some(nonce);
2438        }
2439        tx
2440    }
2441
2442    async fn setup_test_repo() -> RedisTransactionRepository {
2443        // Use a mock Redis URL - in real integration tests, this would connect to a test Redis instance
2444        let redis_url = std::env::var("REDIS_TEST_URL")
2445            .unwrap_or_else(|_| "redis://127.0.0.1:6379".to_string());
2446
2447        let cfg = Config::from_url(&redis_url);
2448        let pool = Arc::new(
2449            cfg.builder()
2450                .expect("Failed to create pool builder")
2451                .max_size(16)
2452                .runtime(Runtime::Tokio1)
2453                .build()
2454                .expect("Failed to build Redis pool"),
2455        );
2456
2457        // Create RedisConnections with same pool for both primary and reader (for testing)
2458        let connections = Arc::new(RedisConnections::new_single_pool(pool));
2459
2460        let random_id = Uuid::new_v4().to_string();
2461        let key_prefix = format!("test_prefix:{random_id}");
2462
2463        RedisTransactionRepository::new(connections, key_prefix)
2464            .expect("Failed to create RedisTransactionRepository")
2465    }
2466
2467    #[tokio::test]
2468    #[ignore = "Requires active Redis instance"]
2469    async fn test_new_repository_creation() {
2470        let repo = setup_test_repo().await;
2471        assert!(repo.key_prefix.contains("test_prefix"));
2472    }
2473
2474    #[tokio::test]
2475    #[ignore = "Requires active Redis instance"]
2476    async fn test_new_repository_empty_prefix_fails() {
2477        let redis_url = std::env::var("REDIS_TEST_URL")
2478            .unwrap_or_else(|_| "redis://127.0.0.1:6379".to_string());
2479        let cfg = Config::from_url(&redis_url);
2480        let pool = Arc::new(
2481            cfg.builder()
2482                .expect("Failed to create pool builder")
2483                .max_size(16)
2484                .runtime(Runtime::Tokio1)
2485                .build()
2486                .expect("Failed to build Redis pool"),
2487        );
2488        let connections = Arc::new(RedisConnections::new_single_pool(pool));
2489
2490        let result = RedisTransactionRepository::new(connections, "".to_string());
2491        assert!(matches!(result, Err(RepositoryError::InvalidData(_))));
2492    }
2493
2494    #[tokio::test]
2495    #[ignore = "Requires active Redis instance"]
2496    async fn test_key_generation() {
2497        let repo = setup_test_repo().await;
2498
2499        assert!(repo
2500            .tx_key("relayer-1", "test-id")
2501            .contains(":relayer:relayer-1:tx:test-id"));
2502        assert!(repo
2503            .tx_to_relayer_key("test-id")
2504            .contains(":relayer:tx_to_relayer:test-id"));
2505        assert!(repo.relayer_list_key().contains(":relayer_list"));
2506        assert!(repo
2507            .relayer_status_key("relayer-1", &TransactionStatus::Pending)
2508            .contains(":relayer:relayer-1:status:Pending"));
2509        assert!(repo
2510            .relayer_nonce_key("relayer-1", 42)
2511            .contains(":relayer:relayer-1:nonce:42"));
2512    }
2513
2514    #[tokio::test]
2515    #[ignore = "Requires active Redis instance"]
2516    async fn test_serialize_deserialize_transaction() {
2517        let repo = setup_test_repo().await;
2518        let tx = create_test_transaction("test-1");
2519
2520        let serialized = repo
2521            .serialize_entity(&tx, |t| &t.id, "transaction")
2522            .expect("Serialization should succeed");
2523        let deserialized: TransactionRepoModel = repo
2524            .deserialize_entity(&serialized, "test-1", "transaction")
2525            .expect("Deserialization should succeed");
2526
2527        assert_eq!(tx.id, deserialized.id);
2528        assert_eq!(tx.relayer_id, deserialized.relayer_id);
2529        assert_eq!(tx.status, deserialized.status);
2530    }
2531
2532    #[tokio::test]
2533    #[ignore = "Requires active Redis instance"]
2534    async fn test_extract_nonce() {
2535        let repo = setup_test_repo().await;
2536        let random_id = Uuid::new_v4().to_string();
2537        let relayer_id = Uuid::new_v4().to_string();
2538        let tx_with_nonce = create_test_transaction_with_nonce(&random_id, 42, &relayer_id);
2539
2540        let nonce = repo.extract_nonce(&tx_with_nonce.network_data);
2541        assert_eq!(nonce, Some(42));
2542    }
2543
2544    #[tokio::test]
2545    #[ignore = "Requires active Redis instance"]
2546    async fn test_create_transaction() {
2547        let repo = setup_test_repo().await;
2548        let random_id = Uuid::new_v4().to_string();
2549        let tx = create_test_transaction(&random_id);
2550
2551        let result = repo.create(tx.clone()).await.unwrap();
2552        assert_eq!(result.id, tx.id);
2553    }
2554
2555    #[tokio::test]
2556    #[ignore = "Requires active Redis instance"]
2557    async fn test_get_transaction() {
2558        let repo = setup_test_repo().await;
2559        let random_id = Uuid::new_v4().to_string();
2560        let tx = create_test_transaction(&random_id);
2561
2562        repo.create(tx.clone()).await.unwrap();
2563        let stored = repo.get_by_id(random_id.to_string()).await.unwrap();
2564        assert_eq!(stored.id, tx.id);
2565        assert_eq!(stored.relayer_id, tx.relayer_id);
2566    }
2567
2568    #[tokio::test]
2569    #[ignore = "Requires active Redis instance"]
2570    async fn test_update_transaction() {
2571        let repo = setup_test_repo().await;
2572        let random_id = Uuid::new_v4().to_string();
2573        let mut tx = create_test_transaction(&random_id);
2574
2575        repo.create(tx.clone()).await.unwrap();
2576        tx.status = TransactionStatus::Confirmed;
2577
2578        let updated = repo.update(random_id.to_string(), tx).await.unwrap();
2579        assert!(matches!(updated.status, TransactionStatus::Confirmed));
2580    }
2581
2582    #[tokio::test]
2583    #[ignore = "Requires active Redis instance"]
2584    async fn test_delete_transaction() {
2585        let repo = setup_test_repo().await;
2586        let random_id = Uuid::new_v4().to_string();
2587        let tx = create_test_transaction(&random_id);
2588
2589        repo.create(tx).await.unwrap();
2590        repo.delete_by_id(random_id.to_string()).await.unwrap();
2591
2592        let result = repo.get_by_id(random_id.to_string()).await;
2593        assert!(matches!(result, Err(RepositoryError::NotFound(_))));
2594    }
2595
2596    #[tokio::test]
2597    #[ignore = "Requires active Redis instance"]
2598    async fn test_list_all_transactions() {
2599        let repo = setup_test_repo().await;
2600        let random_id = Uuid::new_v4().to_string();
2601        let random_id2 = Uuid::new_v4().to_string();
2602
2603        let tx1 = create_test_transaction(&random_id);
2604        let tx2 = create_test_transaction(&random_id2);
2605
2606        repo.create(tx1).await.unwrap();
2607        repo.create(tx2).await.unwrap();
2608
2609        let transactions = repo.list_all().await.unwrap();
2610        assert!(transactions.len() >= 2);
2611    }
2612
2613    #[tokio::test]
2614    #[ignore = "Requires active Redis instance"]
2615    async fn test_count_transactions() {
2616        let repo = setup_test_repo().await;
2617        let random_id = Uuid::new_v4().to_string();
2618        let tx = create_test_transaction(&random_id);
2619
2620        let count = repo.count().await.unwrap();
2621        repo.create(tx).await.unwrap();
2622        assert!(repo.count().await.unwrap() > count);
2623    }
2624
2625    #[tokio::test]
2626    #[ignore = "Requires active Redis instance"]
2627    async fn test_get_nonexistent_transaction() {
2628        let repo = setup_test_repo().await;
2629        let result = repo.get_by_id("nonexistent".to_string()).await;
2630        assert!(matches!(result, Err(RepositoryError::NotFound(_))));
2631    }
2632
2633    #[tokio::test]
2634    #[ignore = "Requires active Redis instance"]
2635    async fn test_duplicate_transaction_creation() {
2636        let repo = setup_test_repo().await;
2637        let random_id = Uuid::new_v4().to_string();
2638
2639        let tx = create_test_transaction(&random_id);
2640
2641        repo.create(tx.clone()).await.unwrap();
2642        let result = repo.create(tx).await;
2643
2644        assert!(matches!(
2645            result,
2646            Err(RepositoryError::ConstraintViolation(_))
2647        ));
2648    }
2649
2650    #[tokio::test]
2651    #[ignore = "Requires active Redis instance"]
2652    async fn test_update_nonexistent_transaction() {
2653        let repo = setup_test_repo().await;
2654        let tx = create_test_transaction("test-1");
2655
2656        let result = repo.update("nonexistent".to_string(), tx).await;
2657        assert!(matches!(result, Err(RepositoryError::NotFound(_))));
2658    }
2659
2660    #[tokio::test]
2661    #[ignore = "Requires active Redis instance"]
2662    async fn test_list_paginated() {
2663        let repo = setup_test_repo().await;
2664
2665        // Create multiple transactions
2666        for _ in 1..=10 {
2667            let random_id = Uuid::new_v4().to_string();
2668            let tx = create_test_transaction(&random_id);
2669            repo.create(tx).await.unwrap();
2670        }
2671
2672        // Test first page with 3 items per page
2673        let query = PaginationQuery {
2674            page: 1,
2675            per_page: 3,
2676        };
2677        let result = repo.list_paginated(query).await.unwrap();
2678        assert_eq!(result.items.len(), 3);
2679        assert!(result.total >= 10);
2680        assert_eq!(result.page, 1);
2681        assert_eq!(result.per_page, 3);
2682
2683        // Test empty page (beyond total items)
2684        let query = PaginationQuery {
2685            page: 1000,
2686            per_page: 3,
2687        };
2688        let result = repo.list_paginated(query).await.unwrap();
2689        assert_eq!(result.items.len(), 0);
2690    }
2691
2692    #[tokio::test]
2693    #[ignore = "Requires active Redis instance"]
2694    async fn test_find_by_relayer_id() {
2695        let repo = setup_test_repo().await;
2696        let random_id = Uuid::new_v4().to_string();
2697        let random_id2 = Uuid::new_v4().to_string();
2698        let random_id3 = Uuid::new_v4().to_string();
2699
2700        let tx1 = create_test_transaction_with_relayer(&random_id, "relayer-1");
2701        let tx2 = create_test_transaction_with_relayer(&random_id2, "relayer-1");
2702        let tx3 = create_test_transaction_with_relayer(&random_id3, "relayer-2");
2703
2704        repo.create(tx1).await.unwrap();
2705        repo.create(tx2).await.unwrap();
2706        repo.create(tx3).await.unwrap();
2707
2708        // Test finding transactions for relayer-1
2709        let query = PaginationQuery {
2710            page: 1,
2711            per_page: 10,
2712        };
2713        let result = repo
2714            .find_by_relayer_id("relayer-1", query.clone())
2715            .await
2716            .unwrap();
2717        assert!(result.total >= 2);
2718        assert!(result.items.len() >= 2);
2719        assert!(result.items.iter().all(|tx| tx.relayer_id == "relayer-1"));
2720
2721        // Test finding transactions for relayer-2
2722        let result = repo
2723            .find_by_relayer_id("relayer-2", query.clone())
2724            .await
2725            .unwrap();
2726        assert!(result.total >= 1);
2727        assert!(!result.items.is_empty());
2728        assert!(result.items.iter().all(|tx| tx.relayer_id == "relayer-2"));
2729
2730        // Test finding transactions for non-existent relayer
2731        let result = repo
2732            .find_by_relayer_id("non-existent", query.clone())
2733            .await
2734            .unwrap();
2735        assert_eq!(result.total, 0);
2736        assert_eq!(result.items.len(), 0);
2737    }
2738
2739    #[tokio::test]
2740    #[ignore = "Requires active Redis instance"]
2741    async fn test_find_by_relayer_id_sorted_by_created_at_newest_first() {
2742        let repo = setup_test_repo().await;
2743        let relayer_id = Uuid::new_v4().to_string();
2744
2745        // Create transactions with different created_at timestamps
2746        let mut tx1 = create_test_transaction_with_relayer("test-1", &relayer_id);
2747        tx1.created_at = "2025-01-27T10:00:00.000000+00:00".to_string(); // Oldest
2748
2749        let mut tx2 = create_test_transaction_with_relayer("test-2", &relayer_id);
2750        tx2.created_at = "2025-01-27T12:00:00.000000+00:00".to_string(); // Middle
2751
2752        let mut tx3 = create_test_transaction_with_relayer("test-3", &relayer_id);
2753        tx3.created_at = "2025-01-27T14:00:00.000000+00:00".to_string(); // Newest
2754
2755        // Create transactions in non-chronological order to ensure sorting works
2756        repo.create(tx2.clone()).await.unwrap(); // Middle first
2757        repo.create(tx1.clone()).await.unwrap(); // Oldest second
2758        repo.create(tx3.clone()).await.unwrap(); // Newest last
2759
2760        let query = PaginationQuery {
2761            page: 1,
2762            per_page: 10,
2763        };
2764        let result = repo.find_by_relayer_id(&relayer_id, query).await.unwrap();
2765
2766        assert_eq!(result.total, 3);
2767        assert_eq!(result.items.len(), 3);
2768
2769        // Verify transactions are sorted by created_at descending (newest first)
2770        assert_eq!(
2771            result.items[0].id, "test-3",
2772            "First item should be newest (test-3)"
2773        );
2774        assert_eq!(
2775            result.items[0].created_at,
2776            "2025-01-27T14:00:00.000000+00:00"
2777        );
2778
2779        assert_eq!(
2780            result.items[1].id, "test-2",
2781            "Second item should be middle (test-2)"
2782        );
2783        assert_eq!(
2784            result.items[1].created_at,
2785            "2025-01-27T12:00:00.000000+00:00"
2786        );
2787
2788        assert_eq!(
2789            result.items[2].id, "test-1",
2790            "Third item should be oldest (test-1)"
2791        );
2792        assert_eq!(
2793            result.items[2].created_at,
2794            "2025-01-27T10:00:00.000000+00:00"
2795        );
2796    }
2797
2798    #[tokio::test]
2799    #[ignore = "Requires active Redis instance"]
2800    async fn test_find_by_status() {
2801        let repo = setup_test_repo().await;
2802        let random_id = Uuid::new_v4().to_string();
2803        let random_id2 = Uuid::new_v4().to_string();
2804        let random_id3 = Uuid::new_v4().to_string();
2805        let relayer_id = Uuid::new_v4().to_string();
2806        let tx1 = create_test_transaction_with_status(
2807            &random_id,
2808            &relayer_id,
2809            TransactionStatus::Pending,
2810        );
2811        let tx2 =
2812            create_test_transaction_with_status(&random_id2, &relayer_id, TransactionStatus::Sent);
2813        let tx3 = create_test_transaction_with_status(
2814            &random_id3,
2815            &relayer_id,
2816            TransactionStatus::Confirmed,
2817        );
2818
2819        repo.create(tx1).await.unwrap();
2820        repo.create(tx2).await.unwrap();
2821        repo.create(tx3).await.unwrap();
2822
2823        // Test finding pending transactions
2824        let result = repo
2825            .find_by_status(&relayer_id, &[TransactionStatus::Pending])
2826            .await
2827            .unwrap();
2828        assert_eq!(result.len(), 1);
2829        assert_eq!(result[0].status, TransactionStatus::Pending);
2830
2831        // Test finding multiple statuses
2832        let result = repo
2833            .find_by_status(
2834                &relayer_id,
2835                &[TransactionStatus::Pending, TransactionStatus::Sent],
2836            )
2837            .await
2838            .unwrap();
2839        assert_eq!(result.len(), 2);
2840
2841        // Test finding non-existent status
2842        let result = repo
2843            .find_by_status(&relayer_id, &[TransactionStatus::Failed])
2844            .await
2845            .unwrap();
2846        assert_eq!(result.len(), 0);
2847    }
2848
2849    #[tokio::test]
2850    #[ignore = "Requires active Redis instance"]
2851    async fn test_find_by_status_paginated() {
2852        let repo = setup_test_repo().await;
2853        let relayer_id = Uuid::new_v4().to_string();
2854
2855        // Create 5 pending transactions with different timestamps
2856        for i in 1..=5 {
2857            let tx_id = Uuid::new_v4().to_string();
2858            let mut tx = create_test_transaction_with_status(
2859                &tx_id,
2860                &relayer_id,
2861                TransactionStatus::Pending,
2862            );
2863            tx.created_at = format!("2025-01-27T{:02}:00:00.000000+00:00", 10 + i);
2864            repo.create(tx).await.unwrap();
2865        }
2866
2867        // Create 2 confirmed transactions
2868        for i in 6..=7 {
2869            let tx_id = Uuid::new_v4().to_string();
2870            let mut tx = create_test_transaction_with_status(
2871                &tx_id,
2872                &relayer_id,
2873                TransactionStatus::Confirmed,
2874            );
2875            tx.created_at = format!("2025-01-27T{:02}:00:00.000000+00:00", 10 + i);
2876            repo.create(tx).await.unwrap();
2877        }
2878
2879        // Test first page (2 items per page)
2880        let query = PaginationQuery {
2881            page: 1,
2882            per_page: 2,
2883        };
2884        let result = repo
2885            .find_by_status_paginated(&relayer_id, &[TransactionStatus::Pending], query, false)
2886            .await
2887            .unwrap();
2888
2889        assert_eq!(result.total, 5);
2890        assert_eq!(result.items.len(), 2);
2891        assert_eq!(result.page, 1);
2892        assert_eq!(result.per_page, 2);
2893
2894        // Test second page
2895        let query = PaginationQuery {
2896            page: 2,
2897            per_page: 2,
2898        };
2899        let result = repo
2900            .find_by_status_paginated(&relayer_id, &[TransactionStatus::Pending], query, false)
2901            .await
2902            .unwrap();
2903
2904        assert_eq!(result.total, 5);
2905        assert_eq!(result.items.len(), 2);
2906        assert_eq!(result.page, 2);
2907
2908        // Test last page (partial)
2909        let query = PaginationQuery {
2910            page: 3,
2911            per_page: 2,
2912        };
2913        let result = repo
2914            .find_by_status_paginated(&relayer_id, &[TransactionStatus::Pending], query, false)
2915            .await
2916            .unwrap();
2917
2918        assert_eq!(result.total, 5);
2919        assert_eq!(result.items.len(), 1);
2920
2921        // Test multiple statuses
2922        let query = PaginationQuery {
2923            page: 1,
2924            per_page: 10,
2925        };
2926        let result = repo
2927            .find_by_status_paginated(
2928                &relayer_id,
2929                &[TransactionStatus::Pending, TransactionStatus::Confirmed],
2930                query,
2931                false,
2932            )
2933            .await
2934            .unwrap();
2935
2936        assert_eq!(result.total, 7);
2937        assert_eq!(result.items.len(), 7);
2938
2939        // Test empty result
2940        let query = PaginationQuery {
2941            page: 1,
2942            per_page: 10,
2943        };
2944        let result = repo
2945            .find_by_status_paginated(&relayer_id, &[TransactionStatus::Failed], query, false)
2946            .await
2947            .unwrap();
2948
2949        assert_eq!(result.total, 0);
2950        assert_eq!(result.items.len(), 0);
2951    }
2952
2953    #[tokio::test]
2954    #[ignore = "Requires active Redis instance"]
2955    async fn test_find_by_status_paginated_oldest_first() {
2956        let repo = setup_test_repo().await;
2957        let relayer_id = Uuid::new_v4().to_string();
2958
2959        // Create 5 pending transactions with ascending timestamps
2960        for i in 1..=5 {
2961            let tx_id = format!("tx{}-{}", i, Uuid::new_v4());
2962            let mut tx = create_test_transaction(&tx_id);
2963            tx.relayer_id = relayer_id.clone();
2964            tx.status = TransactionStatus::Pending;
2965            tx.created_at = format!("2025-01-27T{:02}:00:00.000000+00:00", 10 + i);
2966            repo.create(tx).await.unwrap();
2967        }
2968
2969        // Test oldest_first: true - should return oldest transactions first
2970        let query = PaginationQuery {
2971            page: 1,
2972            per_page: 3,
2973        };
2974        let result = repo
2975            .find_by_status_paginated(
2976                &relayer_id,
2977                &[TransactionStatus::Pending],
2978                query.clone(),
2979                true,
2980            )
2981            .await
2982            .unwrap();
2983
2984        assert_eq!(result.total, 5);
2985        assert_eq!(result.items.len(), 3);
2986        // Verify ordering: oldest first (11:00, 12:00, 13:00)
2987        assert!(
2988            result.items[0].created_at < result.items[1].created_at,
2989            "First item should be older than second"
2990        );
2991        assert!(
2992            result.items[1].created_at < result.items[2].created_at,
2993            "Second item should be older than third"
2994        );
2995
2996        // Contrast with oldest_first: false - should return newest first
2997        let result_newest = repo
2998            .find_by_status_paginated(&relayer_id, &[TransactionStatus::Pending], query, false)
2999            .await
3000            .unwrap();
3001
3002        assert_eq!(result_newest.items.len(), 3);
3003        // Verify ordering: newest first (15:00, 14:00, 13:00)
3004        assert!(
3005            result_newest.items[0].created_at > result_newest.items[1].created_at,
3006            "First item should be newer than second"
3007        );
3008        assert!(
3009            result_newest.items[1].created_at > result_newest.items[2].created_at,
3010            "Second item should be newer than third"
3011        );
3012    }
3013
3014    #[tokio::test]
3015    #[ignore = "Requires active Redis instance"]
3016    async fn test_find_by_status_paginated_oldest_first_single_item() {
3017        let repo = setup_test_repo().await;
3018        let relayer_id = Uuid::new_v4().to_string();
3019
3020        // Create transactions with specific timestamps
3021        let timestamps = [
3022            "2025-01-27T08:00:00.000000+00:00", // oldest
3023            "2025-01-27T10:00:00.000000+00:00", // middle
3024            "2025-01-27T12:00:00.000000+00:00", // newest
3025        ];
3026
3027        let mut oldest_id = String::new();
3028        let mut newest_id = String::new();
3029
3030        for (i, timestamp) in timestamps.iter().enumerate() {
3031            let tx_id = format!("tx-{}-{}", i, Uuid::new_v4());
3032            if i == 0 {
3033                oldest_id = tx_id.clone();
3034            }
3035            if i == 2 {
3036                newest_id = tx_id.clone();
3037            }
3038            let mut tx = create_test_transaction(&tx_id);
3039            tx.relayer_id = relayer_id.clone();
3040            tx.status = TransactionStatus::Pending;
3041            tx.created_at = timestamp.to_string();
3042            repo.create(tx).await.unwrap();
3043        }
3044
3045        // Request just 1 item with oldest_first: true
3046        let query = PaginationQuery {
3047            page: 1,
3048            per_page: 1,
3049        };
3050        let result = repo
3051            .find_by_status_paginated(
3052                &relayer_id,
3053                &[TransactionStatus::Pending],
3054                query.clone(),
3055                true,
3056            )
3057            .await
3058            .unwrap();
3059
3060        assert_eq!(result.total, 3);
3061        assert_eq!(result.items.len(), 1);
3062        assert_eq!(
3063            result.items[0].id, oldest_id,
3064            "With oldest_first=true and per_page=1, should return the oldest transaction"
3065        );
3066
3067        // Contrast with oldest_first: false
3068        let result = repo
3069            .find_by_status_paginated(&relayer_id, &[TransactionStatus::Pending], query, false)
3070            .await
3071            .unwrap();
3072
3073        assert_eq!(result.items.len(), 1);
3074        assert_eq!(
3075            result.items[0].id, newest_id,
3076            "With oldest_first=false and per_page=1, should return the newest transaction"
3077        );
3078    }
3079
3080    #[tokio::test]
3081    #[ignore = "Requires active Redis instance"]
3082    async fn test_find_by_nonce() {
3083        let repo = setup_test_repo().await;
3084        let random_id = Uuid::new_v4().to_string();
3085        let random_id2 = Uuid::new_v4().to_string();
3086        let relayer_id = Uuid::new_v4().to_string();
3087
3088        let tx1 = create_test_transaction_with_nonce(&random_id, 42, &relayer_id);
3089        let tx2 = create_test_transaction_with_nonce(&random_id2, 43, &relayer_id);
3090
3091        repo.create(tx1.clone()).await.unwrap();
3092        repo.create(tx2).await.unwrap();
3093
3094        // Test finding existing nonce
3095        let result = repo.find_by_nonce(&relayer_id, 42).await.unwrap();
3096        assert!(result.is_some());
3097        assert_eq!(result.unwrap().id, random_id);
3098
3099        // Test finding non-existent nonce
3100        let result = repo.find_by_nonce(&relayer_id, 99).await.unwrap();
3101        assert!(result.is_none());
3102
3103        // Test finding nonce for non-existent relayer
3104        let result = repo.find_by_nonce("non-existent", 42).await.unwrap();
3105        assert!(result.is_none());
3106    }
3107
3108    #[tokio::test]
3109    #[ignore = "Requires active Redis instance"]
3110    async fn test_get_nonce_occupancy_mixed_slots() {
3111        let repo = setup_test_repo().await;
3112        let relayer_id = Uuid::new_v4().to_string();
3113
3114        // nonce 10 → Pending (default), nonce 11 → Failed, nonce 12 → empty
3115        let tx1 = create_test_transaction_with_nonce(&Uuid::new_v4().to_string(), 10, &relayer_id);
3116        repo.create(tx1).await.unwrap();
3117
3118        let mut tx2 =
3119            create_test_transaction_with_nonce(&Uuid::new_v4().to_string(), 11, &relayer_id);
3120        tx2.status = TransactionStatus::Failed;
3121        repo.create(tx2).await.unwrap();
3122
3123        let result = repo.get_nonce_occupancy(&relayer_id, 10, 13).await.unwrap();
3124
3125        assert_eq!(result.len(), 3);
3126        assert_eq!(result[0], (10, Some(TransactionStatus::Pending)));
3127        assert_eq!(result[1], (11, Some(TransactionStatus::Failed)));
3128        assert_eq!(result[2], (12, None));
3129    }
3130
3131    #[tokio::test]
3132    #[ignore = "Requires active Redis instance"]
3133    async fn test_get_nonce_occupancy_empty_range() {
3134        let repo = setup_test_repo().await;
3135
3136        let result = repo.get_nonce_occupancy("any-relayer", 5, 5).await.unwrap();
3137        assert!(result.is_empty());
3138
3139        let result = repo
3140            .get_nonce_occupancy("any-relayer", 10, 5)
3141            .await
3142            .unwrap();
3143        assert!(result.is_empty());
3144    }
3145
3146    #[tokio::test]
3147    #[ignore = "Requires active Redis instance"]
3148    async fn test_get_nonce_occupancy_all_empty() {
3149        let repo = setup_test_repo().await;
3150        let relayer_id = Uuid::new_v4().to_string();
3151
3152        // No transactions created — all slots should be None
3153        let result = repo
3154            .get_nonce_occupancy(&relayer_id, 100, 103)
3155            .await
3156            .unwrap();
3157
3158        assert_eq!(result.len(), 3);
3159        assert!(result.iter().all(|(_, status)| status.is_none()));
3160    }
3161
3162    #[tokio::test]
3163    #[ignore = "Requires active Redis instance"]
3164    async fn test_update_status() {
3165        let repo = setup_test_repo().await;
3166        let random_id = Uuid::new_v4().to_string();
3167        let tx = create_test_transaction(&random_id);
3168
3169        repo.create(tx).await.unwrap();
3170        let updated = repo
3171            .update_status(random_id.to_string(), TransactionStatus::Confirmed)
3172            .await
3173            .unwrap();
3174        assert_eq!(updated.status, TransactionStatus::Confirmed);
3175    }
3176
3177    #[tokio::test]
3178    #[ignore = "Requires active Redis instance"]
3179    async fn test_partial_update() {
3180        let repo = setup_test_repo().await;
3181        let random_id = Uuid::new_v4().to_string();
3182        let tx = create_test_transaction(&random_id);
3183
3184        repo.create(tx).await.unwrap();
3185
3186        let update = TransactionUpdateRequest {
3187            status: Some(TransactionStatus::Sent),
3188            status_reason: Some("Transaction sent".to_string()),
3189            sent_at: Some("2025-01-27T16:00:00.000000+00:00".to_string()),
3190            confirmed_at: None,
3191            network_data: None,
3192            hashes: None,
3193            is_canceled: None,
3194            priced_at: None,
3195            noop_count: None,
3196            delete_at: None,
3197            metadata: None,
3198        };
3199
3200        let updated = repo
3201            .partial_update(random_id.to_string(), update)
3202            .await
3203            .unwrap();
3204        assert_eq!(updated.status, TransactionStatus::Sent);
3205        assert_eq!(updated.status_reason, Some("Transaction sent".to_string()));
3206        assert_eq!(
3207            updated.sent_at,
3208            Some("2025-01-27T16:00:00.000000+00:00".to_string())
3209        );
3210    }
3211
3212    #[tokio::test]
3213    #[ignore = "Requires active Redis instance"]
3214    async fn test_partial_update_preserves_large_stellar_i64_fields() {
3215        // End-to-end regression against a real Redis: large i64/u64 fields
3216        // (sequence_number, SignedXdr.max_fee, and the memo Id) must survive the
3217        // partial_update Lua `cjson` decode->encode round-trip. Before the
3218        // string-serde fix, cjson re-emitted these as floats (e.g.
3219        // 643918676885760.0) and the read-back failed with "invalid type:
3220        // floating point, expected i64".
3221        let repo = setup_test_repo().await;
3222        let relayer_id = Uuid::new_v4().to_string();
3223        let tx_id = Uuid::new_v4().to_string();
3224        // All values exceed 2^32 so they exercise the large-integer path that
3225        // the pre-fix code corrupted.
3226        let seq = 643918676885760_i64;
3227        let max_fee = 549755813888_i64;
3228        let memo_id = 1099511627776_u64;
3229
3230        let tx = create_stellar_test_transaction(&tx_id, &relayer_id, seq, max_fee, memo_id);
3231        repo.create(tx).await.unwrap();
3232
3233        let update = TransactionUpdateRequest {
3234            status: Some(TransactionStatus::Sent),
3235            sent_at: Some("2025-01-27T16:00:00.000000+00:00".to_string()),
3236            ..Default::default()
3237        };
3238
3239        // The Lua script returns the re-encoded JSON; this deserialization is
3240        // exactly where the float corruption used to blow up.
3241        let updated = repo.partial_update(tx_id.clone(), update).await.unwrap();
3242        assert_large_stellar_fields(&updated, seq, max_fee, memo_id);
3243
3244        // Re-read from Redis to confirm what was actually persisted, not just
3245        // the script's return value.
3246        let reloaded = repo.get_by_id(tx_id).await.unwrap();
3247        assert_large_stellar_fields(&reloaded, seq, max_fee, memo_id);
3248    }
3249
3250    fn assert_large_stellar_fields(
3251        tx: &TransactionRepoModel,
3252        seq: i64,
3253        max_fee: i64,
3254        memo_id: u64,
3255    ) {
3256        let stellar = tx.network_data.get_stellar_transaction_data().unwrap();
3257        assert_eq!(stellar.sequence_number, Some(seq));
3258        assert_eq!(stellar.memo, Some(MemoSpec::Id { value: memo_id }));
3259        match stellar.transaction_input {
3260            TransactionInput::SignedXdr { max_fee: mf, .. } => assert_eq!(mf, max_fee),
3261            other => panic!("expected SignedXdr, got {other:?}"),
3262        }
3263    }
3264
3265    #[tokio::test]
3266    #[ignore = "Requires active Redis instance"]
3267    async fn test_set_sent_at() {
3268        let repo = setup_test_repo().await;
3269        let random_id = Uuid::new_v4().to_string();
3270        let tx = create_test_transaction(&random_id);
3271
3272        repo.create(tx).await.unwrap();
3273        let updated = repo
3274            .set_sent_at(
3275                random_id.to_string(),
3276                "2025-01-27T16:00:00.000000+00:00".to_string(),
3277            )
3278            .await
3279            .unwrap();
3280        assert_eq!(
3281            updated.sent_at,
3282            Some("2025-01-27T16:00:00.000000+00:00".to_string())
3283        );
3284    }
3285
3286    #[tokio::test]
3287    #[ignore = "Requires active Redis instance"]
3288    async fn test_set_confirmed_at() {
3289        let repo = setup_test_repo().await;
3290        let random_id = Uuid::new_v4().to_string();
3291        let tx = create_test_transaction(&random_id);
3292
3293        repo.create(tx).await.unwrap();
3294        let updated = repo
3295            .set_confirmed_at(
3296                random_id.to_string(),
3297                "2025-01-27T16:00:00.000000+00:00".to_string(),
3298            )
3299            .await
3300            .unwrap();
3301        assert_eq!(
3302            updated.confirmed_at,
3303            Some("2025-01-27T16:00:00.000000+00:00".to_string())
3304        );
3305    }
3306
3307    #[tokio::test]
3308    #[ignore = "Requires active Redis instance"]
3309    async fn test_update_network_data() {
3310        let repo = setup_test_repo().await;
3311        let random_id = Uuid::new_v4().to_string();
3312        let tx = create_test_transaction(&random_id);
3313
3314        repo.create(tx).await.unwrap();
3315
3316        let new_network_data = NetworkTransactionData::Evm(EvmTransactionData {
3317            gas_price: Some(2000000000),
3318            gas_limit: Some(42000),
3319            nonce: Some(2),
3320            value: U256::from_str("2000000000000000000").unwrap(),
3321            data: Some("0x1234".to_string()),
3322            from: "0xNewSender".to_string(),
3323            to: Some("0xNewRecipient".to_string()),
3324            chain_id: 1,
3325            signature: None,
3326            hash: Some("0xnewhash".to_string()),
3327            speed: Some(Speed::SafeLow),
3328            max_fee_per_gas: None,
3329            max_priority_fee_per_gas: None,
3330            raw: None,
3331        });
3332
3333        let updated = repo
3334            .update_network_data(random_id.to_string(), new_network_data.clone())
3335            .await
3336            .unwrap();
3337        assert_eq!(
3338            updated
3339                .network_data
3340                .get_evm_transaction_data()
3341                .unwrap()
3342                .hash,
3343            new_network_data.get_evm_transaction_data().unwrap().hash
3344        );
3345    }
3346
3347    #[tokio::test]
3348    #[ignore = "Requires active Redis instance"]
3349    async fn test_debug_implementation() {
3350        let repo = setup_test_repo().await;
3351        let debug_str = format!("{repo:?}");
3352        assert!(debug_str.contains("RedisTransactionRepository"));
3353        assert!(debug_str.contains("test_prefix"));
3354    }
3355
3356    #[tokio::test]
3357    #[ignore = "Requires active Redis instance"]
3358    async fn test_error_handling_empty_id() {
3359        let repo = setup_test_repo().await;
3360
3361        let result = repo.get_by_id("".to_string()).await;
3362        assert!(matches!(result, Err(RepositoryError::InvalidData(_))));
3363
3364        let result = repo
3365            .update("".to_string(), create_test_transaction("test"))
3366            .await;
3367        assert!(matches!(result, Err(RepositoryError::InvalidData(_))));
3368
3369        let result = repo.delete_by_id("".to_string()).await;
3370        assert!(matches!(result, Err(RepositoryError::InvalidData(_))));
3371    }
3372
3373    #[tokio::test]
3374    #[ignore = "Requires active Redis instance"]
3375    async fn test_pagination_validation() {
3376        let repo = setup_test_repo().await;
3377
3378        let query = PaginationQuery {
3379            page: 1,
3380            per_page: 0,
3381        };
3382        let result = repo.list_paginated(query).await;
3383        assert!(matches!(result, Err(RepositoryError::InvalidData(_))));
3384    }
3385
3386    #[tokio::test]
3387    #[ignore = "Requires active Redis instance"]
3388    async fn test_index_consistency() {
3389        let repo = setup_test_repo().await;
3390        let random_id = Uuid::new_v4().to_string();
3391        let relayer_id = Uuid::new_v4().to_string();
3392        let tx = create_test_transaction_with_nonce(&random_id, 42, &relayer_id);
3393
3394        // Create transaction
3395        repo.create(tx.clone()).await.unwrap();
3396
3397        // Verify it can be found by nonce
3398        let found = repo.find_by_nonce(&relayer_id, 42).await.unwrap();
3399        assert!(found.is_some());
3400
3401        // Update the transaction with a new nonce
3402        let mut updated_tx = tx.clone();
3403        if let NetworkTransactionData::Evm(ref mut evm_data) = updated_tx.network_data {
3404            evm_data.nonce = Some(43);
3405        }
3406
3407        repo.update(random_id.to_string(), updated_tx)
3408            .await
3409            .unwrap();
3410
3411        // Verify old nonce index is cleaned up
3412        let old_nonce_result = repo.find_by_nonce(&relayer_id, 42).await.unwrap();
3413        assert!(old_nonce_result.is_none());
3414
3415        // Verify new nonce index works
3416        let new_nonce_result = repo.find_by_nonce(&relayer_id, 43).await.unwrap();
3417        assert!(new_nonce_result.is_some());
3418    }
3419
3420    #[tokio::test]
3421    #[ignore = "Requires active Redis instance"]
3422    async fn test_has_entries() {
3423        let repo = setup_test_repo().await;
3424        assert!(!repo.has_entries().await.unwrap());
3425
3426        let tx_id = uuid::Uuid::new_v4().to_string();
3427        let tx = create_test_transaction(&tx_id);
3428        repo.create(tx.clone()).await.unwrap();
3429
3430        assert!(repo.has_entries().await.unwrap());
3431    }
3432
3433    #[tokio::test]
3434    #[ignore = "Requires active Redis instance"]
3435    async fn test_drop_all_entries() {
3436        let repo = setup_test_repo().await;
3437        let tx_id = uuid::Uuid::new_v4().to_string();
3438        let tx = create_test_transaction(&tx_id);
3439        repo.create(tx.clone()).await.unwrap();
3440        assert!(repo.has_entries().await.unwrap());
3441
3442        repo.drop_all_entries().await.unwrap();
3443        assert!(!repo.has_entries().await.unwrap());
3444    }
3445
3446    // Tests for delete_at field setting on final status updates
3447    #[tokio::test]
3448    #[ignore = "Requires active Redis instance"]
3449    async fn test_update_status_sets_delete_at_for_final_statuses() {
3450        let _lock = ENV_MUTEX.lock().await;
3451
3452        use chrono::{DateTime, Duration, Utc};
3453        use std::env;
3454
3455        // Use a unique test environment variable to avoid conflicts
3456        env::set_var("TRANSACTION_EXPIRATION_HOURS", "6");
3457
3458        let repo = setup_test_repo().await;
3459
3460        let final_statuses = [
3461            TransactionStatus::Canceled,
3462            TransactionStatus::Confirmed,
3463            TransactionStatus::Failed,
3464            TransactionStatus::Expired,
3465        ];
3466
3467        for (i, status) in final_statuses.iter().enumerate() {
3468            let tx_id = format!("test-final-{}-{}", i, Uuid::new_v4());
3469            let mut tx = create_test_transaction(&tx_id);
3470
3471            // Ensure transaction has no delete_at initially and is in pending state
3472            tx.delete_at = None;
3473            tx.status = TransactionStatus::Pending;
3474
3475            repo.create(tx).await.unwrap();
3476
3477            let before_update = Utc::now();
3478
3479            // Update to final status
3480            let updated = repo
3481                .update_status(tx_id.clone(), status.clone())
3482                .await
3483                .unwrap();
3484
3485            // Should have delete_at set
3486            assert!(
3487                updated.delete_at.is_some(),
3488                "delete_at should be set for status: {status:?}"
3489            );
3490
3491            // Verify the timestamp is reasonable (approximately 6 hours from now)
3492            let delete_at_str = updated.delete_at.unwrap();
3493            let delete_at = DateTime::parse_from_rfc3339(&delete_at_str)
3494                .expect("delete_at should be valid RFC3339")
3495                .with_timezone(&Utc);
3496
3497            let duration_from_before = delete_at.signed_duration_since(before_update);
3498            let expected_duration = Duration::hours(6);
3499            let tolerance = Duration::minutes(5);
3500
3501            assert!(
3502                duration_from_before >= expected_duration - tolerance
3503                    && duration_from_before <= expected_duration + tolerance,
3504                "delete_at should be approximately 6 hours from now for status: {status:?}. Duration: {duration_from_before:?}"
3505            );
3506        }
3507
3508        // Cleanup
3509        env::remove_var("TRANSACTION_EXPIRATION_HOURS");
3510    }
3511
3512    #[tokio::test]
3513    #[ignore = "Requires active Redis instance"]
3514    async fn test_update_status_does_not_set_delete_at_for_non_final_statuses() {
3515        let _lock = ENV_MUTEX.lock().await;
3516
3517        use std::env;
3518
3519        env::set_var("TRANSACTION_EXPIRATION_HOURS", "4");
3520
3521        let repo = setup_test_repo().await;
3522
3523        let non_final_statuses = [
3524            TransactionStatus::Pending,
3525            TransactionStatus::Sent,
3526            TransactionStatus::Submitted,
3527            TransactionStatus::Mined,
3528        ];
3529
3530        for (i, status) in non_final_statuses.iter().enumerate() {
3531            let tx_id = format!("test-non-final-{}-{}", i, Uuid::new_v4());
3532            let mut tx = create_test_transaction(&tx_id);
3533            tx.delete_at = None;
3534            tx.status = TransactionStatus::Pending;
3535
3536            repo.create(tx).await.unwrap();
3537
3538            // Update to non-final status
3539            let updated = repo
3540                .update_status(tx_id.clone(), status.clone())
3541                .await
3542                .unwrap();
3543
3544            // Should NOT have delete_at set
3545            assert!(
3546                updated.delete_at.is_none(),
3547                "delete_at should NOT be set for status: {status:?}"
3548            );
3549        }
3550
3551        // Cleanup
3552        env::remove_var("TRANSACTION_EXPIRATION_HOURS");
3553    }
3554
3555    #[tokio::test]
3556    #[ignore = "Requires active Redis instance"]
3557    async fn test_partial_update_sets_delete_at_for_final_statuses() {
3558        let _lock = ENV_MUTEX.lock().await;
3559
3560        use chrono::{DateTime, Duration, Utc};
3561        use std::env;
3562
3563        env::set_var("TRANSACTION_EXPIRATION_HOURS", "8");
3564
3565        let repo = setup_test_repo().await;
3566        let tx_id = format!("test-partial-final-{}", Uuid::new_v4());
3567        let mut tx = create_test_transaction(&tx_id);
3568        tx.delete_at = None;
3569        tx.status = TransactionStatus::Pending;
3570
3571        repo.create(tx).await.unwrap();
3572
3573        let before_update = Utc::now();
3574
3575        // Use partial_update to set status to Confirmed (final status)
3576        let update = TransactionUpdateRequest {
3577            status: Some(TransactionStatus::Confirmed),
3578            status_reason: Some("Transaction completed".to_string()),
3579            confirmed_at: Some("2023-01-01T12:05:00Z".to_string()),
3580            ..Default::default()
3581        };
3582
3583        let updated = repo.partial_update(tx_id.clone(), update).await.unwrap();
3584
3585        // Should have delete_at set
3586        assert!(
3587            updated.delete_at.is_some(),
3588            "delete_at should be set when updating to Confirmed status"
3589        );
3590
3591        // Verify the timestamp is reasonable (approximately 8 hours from now)
3592        let delete_at_str = updated.delete_at.unwrap();
3593        let delete_at = DateTime::parse_from_rfc3339(&delete_at_str)
3594            .expect("delete_at should be valid RFC3339")
3595            .with_timezone(&Utc);
3596
3597        let duration_from_before = delete_at.signed_duration_since(before_update);
3598        let expected_duration = Duration::hours(8);
3599        let tolerance = Duration::minutes(5);
3600
3601        assert!(
3602            duration_from_before >= expected_duration - tolerance
3603                && duration_from_before <= expected_duration + tolerance,
3604            "delete_at should be approximately 8 hours from now. Duration: {duration_from_before:?}"
3605        );
3606
3607        // Also verify other fields were updated
3608        assert_eq!(updated.status, TransactionStatus::Confirmed);
3609        assert_eq!(
3610            updated.status_reason,
3611            Some("Transaction completed".to_string())
3612        );
3613        assert_eq!(
3614            updated.confirmed_at,
3615            Some("2023-01-01T12:05:00Z".to_string())
3616        );
3617
3618        // Cleanup
3619        env::remove_var("TRANSACTION_EXPIRATION_HOURS");
3620    }
3621
3622    #[tokio::test]
3623    #[ignore = "Requires active Redis instance"]
3624    async fn test_update_status_preserves_existing_delete_at() {
3625        let _lock = ENV_MUTEX.lock().await;
3626
3627        use std::env;
3628
3629        env::set_var("TRANSACTION_EXPIRATION_HOURS", "2");
3630
3631        let repo = setup_test_repo().await;
3632        let tx_id = format!("test-preserve-delete-at-{}", Uuid::new_v4());
3633        let mut tx = create_test_transaction(&tx_id);
3634
3635        // Set an existing delete_at value
3636        let existing_delete_at = "2025-01-01T12:00:00Z".to_string();
3637        tx.delete_at = Some(existing_delete_at.clone());
3638        tx.status = TransactionStatus::Pending;
3639
3640        repo.create(tx).await.unwrap();
3641
3642        // Update to final status
3643        let updated = repo
3644            .update_status(tx_id.clone(), TransactionStatus::Confirmed)
3645            .await
3646            .unwrap();
3647
3648        // Should preserve the existing delete_at value
3649        assert_eq!(
3650            updated.delete_at,
3651            Some(existing_delete_at),
3652            "Existing delete_at should be preserved when updating to final status"
3653        );
3654
3655        // Cleanup
3656        env::remove_var("TRANSACTION_EXPIRATION_HOURS");
3657    }
3658    #[tokio::test]
3659    #[ignore = "Requires active Redis instance"]
3660    async fn test_partial_update_without_status_change_preserves_delete_at() {
3661        let _lock = ENV_MUTEX.lock().await;
3662
3663        use std::env;
3664
3665        env::set_var("TRANSACTION_EXPIRATION_HOURS", "3");
3666
3667        let repo = setup_test_repo().await;
3668        let tx_id = format!("test-preserve-no-status-{}", Uuid::new_v4());
3669        let mut tx = create_test_transaction(&tx_id);
3670        tx.delete_at = None;
3671        tx.status = TransactionStatus::Pending;
3672
3673        repo.create(tx).await.unwrap();
3674
3675        // First, update to final status to set delete_at
3676        let updated1 = repo
3677            .update_status(tx_id.clone(), TransactionStatus::Confirmed)
3678            .await
3679            .unwrap();
3680
3681        assert!(updated1.delete_at.is_some());
3682        let original_delete_at = updated1.delete_at.clone();
3683
3684        // Now update other fields without changing status
3685        let update = TransactionUpdateRequest {
3686            status: None, // No status change
3687            status_reason: Some("Updated reason".to_string()),
3688            confirmed_at: Some("2023-01-01T12:10:00Z".to_string()),
3689            ..Default::default()
3690        };
3691
3692        let updated2 = repo.partial_update(tx_id.clone(), update).await.unwrap();
3693
3694        // delete_at should be preserved
3695        assert_eq!(
3696            updated2.delete_at, original_delete_at,
3697            "delete_at should be preserved when status is not updated"
3698        );
3699
3700        // Other fields should be updated
3701        assert_eq!(updated2.status, TransactionStatus::Confirmed); // Unchanged
3702        assert_eq!(updated2.status_reason, Some("Updated reason".to_string()));
3703        assert_eq!(
3704            updated2.confirmed_at,
3705            Some("2023-01-01T12:10:00Z".to_string())
3706        );
3707
3708        // Cleanup
3709        env::remove_var("TRANSACTION_EXPIRATION_HOURS");
3710    }
3711
3712    // Tests for delete_by_ids batch delete functionality
3713
3714    #[tokio::test]
3715    #[ignore = "Requires active Redis instance"]
3716    async fn test_delete_by_ids_empty_list() {
3717        let repo = setup_test_repo().await;
3718        let tx_id = format!("test-empty-{}", Uuid::new_v4());
3719
3720        // Create a transaction to ensure repo is not empty
3721        let tx = create_test_transaction(&tx_id);
3722        repo.create(tx).await.unwrap();
3723
3724        // Delete with empty list should succeed and not affect existing data
3725        let result = repo.delete_by_ids(vec![]).await.unwrap();
3726
3727        assert_eq!(result.deleted_count, 0);
3728        assert!(result.failed.is_empty());
3729
3730        // Original transaction should still exist
3731        assert!(repo.get_by_id(tx_id).await.is_ok());
3732    }
3733
3734    #[tokio::test]
3735    #[ignore = "Requires active Redis instance"]
3736    async fn test_delete_by_ids_single_transaction() {
3737        let repo = setup_test_repo().await;
3738        let tx_id = format!("test-single-{}", Uuid::new_v4());
3739
3740        let tx = create_test_transaction(&tx_id);
3741        repo.create(tx).await.unwrap();
3742
3743        let result = repo.delete_by_ids(vec![tx_id.clone()]).await.unwrap();
3744
3745        assert_eq!(result.deleted_count, 1);
3746        assert!(result.failed.is_empty());
3747
3748        // Verify transaction was deleted
3749        assert!(repo.get_by_id(tx_id).await.is_err());
3750    }
3751
3752    #[tokio::test]
3753    #[ignore = "Requires active Redis instance"]
3754    async fn test_delete_by_ids_multiple_transactions() {
3755        let repo = setup_test_repo().await;
3756        let base_id = Uuid::new_v4();
3757
3758        // Create multiple transactions
3759        let mut created_ids = Vec::new();
3760        for i in 1..=5 {
3761            let tx_id = format!("test-multi-{base_id}-{i}");
3762            let tx = create_test_transaction(&tx_id);
3763            repo.create(tx).await.unwrap();
3764            created_ids.push(tx_id);
3765        }
3766
3767        // Delete 3 of them
3768        let ids_to_delete = vec![
3769            created_ids[0].clone(),
3770            created_ids[2].clone(),
3771            created_ids[4].clone(),
3772        ];
3773        let result = repo.delete_by_ids(ids_to_delete).await.unwrap();
3774
3775        assert_eq!(result.deleted_count, 3);
3776        assert!(result.failed.is_empty());
3777
3778        // Verify correct transactions were deleted
3779        assert!(repo.get_by_id(created_ids[0].clone()).await.is_err());
3780        assert!(repo.get_by_id(created_ids[1].clone()).await.is_ok()); // Not deleted
3781        assert!(repo.get_by_id(created_ids[2].clone()).await.is_err());
3782        assert!(repo.get_by_id(created_ids[3].clone()).await.is_ok()); // Not deleted
3783        assert!(repo.get_by_id(created_ids[4].clone()).await.is_err());
3784    }
3785
3786    #[tokio::test]
3787    #[ignore = "Requires active Redis instance"]
3788    async fn test_delete_by_ids_nonexistent_transactions() {
3789        let repo = setup_test_repo().await;
3790        let base_id = Uuid::new_v4();
3791
3792        // Try to delete transactions that don't exist
3793        let ids_to_delete = vec![
3794            format!("nonexistent-{}-1", base_id),
3795            format!("nonexistent-{}-2", base_id),
3796        ];
3797        let result = repo.delete_by_ids(ids_to_delete.clone()).await.unwrap();
3798
3799        assert_eq!(result.deleted_count, 0);
3800        assert_eq!(result.failed.len(), 2);
3801
3802        // Verify error messages contain the IDs
3803        let failed_ids: Vec<&String> = result.failed.iter().map(|(id, _)| id).collect();
3804        assert!(failed_ids.contains(&&ids_to_delete[0]));
3805        assert!(failed_ids.contains(&&ids_to_delete[1]));
3806    }
3807
3808    #[tokio::test]
3809    #[ignore = "Requires active Redis instance"]
3810    async fn test_delete_by_ids_mixed_existing_and_nonexistent() {
3811        let repo = setup_test_repo().await;
3812        let base_id = Uuid::new_v4();
3813
3814        // Create some transactions
3815        let existing_ids: Vec<String> = (1..=3)
3816            .map(|i| format!("test-mixed-existing-{base_id}-{i}"))
3817            .collect();
3818
3819        for id in &existing_ids {
3820            let tx = create_test_transaction(id);
3821            repo.create(tx).await.unwrap();
3822        }
3823
3824        let nonexistent_ids: Vec<String> = (1..=2)
3825            .map(|i| format!("test-mixed-nonexistent-{base_id}-{i}"))
3826            .collect();
3827
3828        // Try to delete mix of existing and non-existing
3829        let ids_to_delete = vec![
3830            existing_ids[0].clone(),
3831            nonexistent_ids[0].clone(),
3832            existing_ids[1].clone(),
3833            nonexistent_ids[1].clone(),
3834        ];
3835        let result = repo.delete_by_ids(ids_to_delete).await.unwrap();
3836
3837        assert_eq!(result.deleted_count, 2);
3838        assert_eq!(result.failed.len(), 2);
3839
3840        // Verify existing transactions were deleted
3841        assert!(repo.get_by_id(existing_ids[0].clone()).await.is_err());
3842        assert!(repo.get_by_id(existing_ids[1].clone()).await.is_err());
3843
3844        // Verify remaining transaction still exists
3845        assert!(repo.get_by_id(existing_ids[2].clone()).await.is_ok());
3846    }
3847
3848    #[tokio::test]
3849    #[ignore = "Requires active Redis instance"]
3850    async fn test_delete_by_ids_removes_all_indexes() {
3851        let repo = setup_test_repo().await;
3852        let relayer_id = format!("relayer-{}", Uuid::new_v4());
3853        let tx_id = format!("test-indexes-{}", Uuid::new_v4());
3854
3855        // Create a transaction with specific status
3856        let mut tx = create_test_transaction(&tx_id);
3857        tx.relayer_id = relayer_id.clone();
3858        tx.status = TransactionStatus::Confirmed;
3859        repo.create(tx).await.unwrap();
3860
3861        // Verify transaction exists and is indexed
3862        let found = repo
3863            .find_by_status(&relayer_id, &[TransactionStatus::Confirmed])
3864            .await
3865            .unwrap();
3866        assert!(found.iter().any(|t| t.id == tx_id));
3867
3868        // Delete the transaction
3869        let result = repo.delete_by_ids(vec![tx_id.clone()]).await.unwrap();
3870        assert_eq!(result.deleted_count, 1);
3871
3872        // Verify transaction is no longer in status index
3873        let found_after = repo
3874            .find_by_status(&relayer_id, &[TransactionStatus::Confirmed])
3875            .await
3876            .unwrap();
3877        assert!(!found_after.iter().any(|t| t.id == tx_id));
3878
3879        // Verify transaction cannot be found
3880        assert!(repo.get_by_id(tx_id).await.is_err());
3881    }
3882
3883    #[tokio::test]
3884    #[ignore = "Requires active Redis instance"]
3885    async fn test_delete_by_ids_removes_nonce_index() {
3886        let repo = setup_test_repo().await;
3887        let relayer_id = format!("relayer-{}", Uuid::new_v4());
3888        let tx_id = format!("test-nonce-{}", Uuid::new_v4());
3889        let nonce = 12345u64;
3890
3891        // Create a transaction with a specific nonce
3892        let tx = create_test_transaction_with_nonce(&tx_id, nonce, &relayer_id);
3893        repo.create(tx).await.unwrap();
3894
3895        // Verify nonce index works
3896        let found = repo.find_by_nonce(&relayer_id, nonce).await.unwrap();
3897        assert!(found.is_some());
3898        assert_eq!(found.unwrap().id, tx_id);
3899
3900        // Delete the transaction
3901        let result = repo.delete_by_ids(vec![tx_id.clone()]).await.unwrap();
3902        assert_eq!(result.deleted_count, 1);
3903
3904        // Verify nonce index was cleaned up
3905        let found_after = repo.find_by_nonce(&relayer_id, nonce).await.unwrap();
3906        assert!(found_after.is_none());
3907    }
3908
3909    #[tokio::test]
3910    #[ignore = "Requires active Redis instance"]
3911    async fn test_delete_by_ids_large_batch() {
3912        let repo = setup_test_repo().await;
3913        let base_id = Uuid::new_v4();
3914
3915        // Create many transactions to test batch performance
3916        let count = 50;
3917        let mut created_ids = Vec::new();
3918
3919        for i in 0..count {
3920            let tx_id = format!("test-large-{base_id}-{i}");
3921            let tx = create_test_transaction(&tx_id);
3922            repo.create(tx).await.unwrap();
3923            created_ids.push(tx_id);
3924        }
3925
3926        // Delete all of them in one batch
3927        let result = repo.delete_by_ids(created_ids.clone()).await.unwrap();
3928
3929        assert_eq!(result.deleted_count, count);
3930        assert!(result.failed.is_empty());
3931
3932        // Verify all were deleted
3933        for id in created_ids {
3934            assert!(repo.get_by_id(id).await.is_err());
3935        }
3936    }
3937
3938    #[tokio::test]
3939    #[ignore = "Requires active Redis instance"]
3940    async fn test_delete_by_ids_preserves_other_relayer_transactions() {
3941        let repo = setup_test_repo().await;
3942        let relayer_1 = format!("relayer-1-{}", Uuid::new_v4());
3943        let relayer_2 = format!("relayer-2-{}", Uuid::new_v4());
3944        let tx_id_1 = format!("tx-relayer-1-{}", Uuid::new_v4());
3945        let tx_id_2 = format!("tx-relayer-2-{}", Uuid::new_v4());
3946
3947        // Create transactions for different relayers
3948        let tx1 = create_test_transaction_with_relayer(&tx_id_1, &relayer_1);
3949        let tx2 = create_test_transaction_with_relayer(&tx_id_2, &relayer_2);
3950
3951        repo.create(tx1).await.unwrap();
3952        repo.create(tx2).await.unwrap();
3953
3954        // Delete only relayer-1's transaction
3955        let result = repo.delete_by_ids(vec![tx_id_1.clone()]).await.unwrap();
3956
3957        assert_eq!(result.deleted_count, 1);
3958
3959        // relayer-1's transaction should be deleted
3960        assert!(repo.get_by_id(tx_id_1).await.is_err());
3961
3962        // relayer-2's transaction should still exist
3963        let remaining = repo.get_by_id(tx_id_2).await.unwrap();
3964        assert_eq!(remaining.relayer_id, relayer_2);
3965    }
3966
3967    // ── increment_status_check_failures ─────────────────────────────
3968
3969    #[tokio::test]
3970    #[ignore = "Requires active Redis instance"]
3971    async fn test_increment_status_check_failures_no_prior_metadata() {
3972        let _lock = ENV_MUTEX.lock().await;
3973        let repo = setup_test_repo().await;
3974        let relayer_id = Uuid::new_v4().to_string();
3975        let tx_id = Uuid::new_v4().to_string();
3976        let mut tx =
3977            create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Sent);
3978        tx.metadata = None;
3979        repo.create(tx).await.unwrap();
3980
3981        let updated = repo.increment_status_check_failures(tx_id).await.unwrap();
3982
3983        let meta = updated.metadata.expect("metadata should be set");
3984        assert_eq!(meta.consecutive_failures, 1);
3985        assert_eq!(meta.total_failures, 1);
3986        assert_eq!(meta.insufficient_fee_retries, 0);
3987    }
3988
3989    #[tokio::test]
3990    #[ignore = "Requires active Redis instance"]
3991    async fn test_increment_status_check_failures_accumulates() {
3992        let _lock = ENV_MUTEX.lock().await;
3993        let repo = setup_test_repo().await;
3994        let relayer_id = Uuid::new_v4().to_string();
3995        let tx_id = Uuid::new_v4().to_string();
3996        let tx = create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Sent);
3997        repo.create(tx).await.unwrap();
3998
3999        repo.increment_status_check_failures(tx_id.clone())
4000            .await
4001            .unwrap();
4002        repo.increment_status_check_failures(tx_id.clone())
4003            .await
4004            .unwrap();
4005        let updated = repo.increment_status_check_failures(tx_id).await.unwrap();
4006
4007        let meta = updated.metadata.unwrap();
4008        assert_eq!(meta.consecutive_failures, 3);
4009        assert_eq!(meta.total_failures, 3);
4010    }
4011
4012    #[tokio::test]
4013    #[ignore = "Requires active Redis instance"]
4014    async fn test_increment_status_check_failures_noop_on_final_state() {
4015        let _lock = ENV_MUTEX.lock().await;
4016        let repo = setup_test_repo().await;
4017        let relayer_id = Uuid::new_v4().to_string();
4018        let tx_id = Uuid::new_v4().to_string();
4019        let tx =
4020            create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Confirmed);
4021        repo.create(tx).await.unwrap();
4022
4023        let result = repo.increment_status_check_failures(tx_id).await.unwrap();
4024
4025        // Should return unchanged — no metadata mutation on final state
4026        assert!(result.metadata.is_none());
4027        assert_eq!(result.status, TransactionStatus::Confirmed);
4028    }
4029
4030    #[tokio::test]
4031    #[ignore = "Requires active Redis instance"]
4032    async fn test_increment_status_check_failures_not_found() {
4033        let _lock = ENV_MUTEX.lock().await;
4034        let repo = setup_test_repo().await;
4035
4036        let result = repo
4037            .increment_status_check_failures("nonexistent".to_string())
4038            .await;
4039
4040        assert!(matches!(result, Err(RepositoryError::NotFound(_))));
4041    }
4042
4043    // ── reset_status_check_consecutive_failures ─────────────────────
4044
4045    #[tokio::test]
4046    #[ignore = "Requires active Redis instance"]
4047    async fn test_reset_consecutive_failures() {
4048        let _lock = ENV_MUTEX.lock().await;
4049        let repo = setup_test_repo().await;
4050        let relayer_id = Uuid::new_v4().to_string();
4051        let tx_id = Uuid::new_v4().to_string();
4052        let tx = create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Sent);
4053        repo.create(tx).await.unwrap();
4054
4055        // Increment a few times first
4056        repo.increment_status_check_failures(tx_id.clone())
4057            .await
4058            .unwrap();
4059        repo.increment_status_check_failures(tx_id.clone())
4060            .await
4061            .unwrap();
4062
4063        let updated = repo
4064            .reset_status_check_consecutive_failures(tx_id)
4065            .await
4066            .unwrap();
4067
4068        let meta = updated.metadata.unwrap();
4069        assert_eq!(meta.consecutive_failures, 0);
4070        // total_failures should be preserved
4071        assert_eq!(meta.total_failures, 2);
4072    }
4073
4074    #[tokio::test]
4075    #[ignore = "Requires active Redis instance"]
4076    async fn test_reset_consecutive_failures_noop_on_final_state() {
4077        let _lock = ENV_MUTEX.lock().await;
4078        let repo = setup_test_repo().await;
4079        let relayer_id = Uuid::new_v4().to_string();
4080        let tx_id = Uuid::new_v4().to_string();
4081        let mut tx =
4082            create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Failed);
4083        tx.metadata = Some(crate::models::TransactionMetadata {
4084            consecutive_failures: 5,
4085            total_failures: 10,
4086            insufficient_fee_retries: 0,
4087            try_again_later_retries: 0,
4088            nonce_too_high_retries: 0,
4089        });
4090        repo.create(tx).await.unwrap();
4091
4092        let result = repo
4093            .reset_status_check_consecutive_failures(tx_id)
4094            .await
4095            .unwrap();
4096
4097        // Should return unchanged on final state
4098        let meta = result.metadata.unwrap();
4099        assert_eq!(meta.consecutive_failures, 5);
4100    }
4101
4102    #[tokio::test]
4103    #[ignore = "Requires active Redis instance"]
4104    async fn test_reset_consecutive_failures_not_found() {
4105        let _lock = ENV_MUTEX.lock().await;
4106        let repo = setup_test_repo().await;
4107
4108        let result = repo
4109            .reset_status_check_consecutive_failures("nonexistent".to_string())
4110            .await;
4111
4112        assert!(matches!(result, Err(RepositoryError::NotFound(_))));
4113    }
4114
4115    // ── record_stellar_insufficient_fee_retry ───────────────────────
4116
4117    #[tokio::test]
4118    #[ignore = "Requires active Redis instance"]
4119    async fn test_record_insufficient_fee_retry() {
4120        let _lock = ENV_MUTEX.lock().await;
4121        let repo = setup_test_repo().await;
4122        let relayer_id = Uuid::new_v4().to_string();
4123        let tx_id = Uuid::new_v4().to_string();
4124        let mut tx =
4125            create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Sent);
4126        tx.sent_at = None;
4127        repo.create(tx).await.unwrap();
4128
4129        let updated = repo
4130            .record_stellar_insufficient_fee_retry(tx_id, "2025-03-18T10:00:00Z".to_string())
4131            .await
4132            .unwrap();
4133
4134        assert_eq!(updated.sent_at.as_deref(), Some("2025-03-18T10:00:00Z"));
4135        let meta = updated.metadata.unwrap();
4136        assert_eq!(meta.insufficient_fee_retries, 1);
4137        assert_eq!(meta.consecutive_failures, 0);
4138        assert_eq!(meta.total_failures, 0);
4139    }
4140
4141    #[tokio::test]
4142    #[ignore = "Requires active Redis instance"]
4143    async fn test_record_insufficient_fee_retry_accumulates() {
4144        let _lock = ENV_MUTEX.lock().await;
4145        let repo = setup_test_repo().await;
4146        let relayer_id = Uuid::new_v4().to_string();
4147        let tx_id = Uuid::new_v4().to_string();
4148        let tx = create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Sent);
4149        repo.create(tx).await.unwrap();
4150
4151        repo.record_stellar_insufficient_fee_retry(
4152            tx_id.clone(),
4153            "2025-03-18T10:00:00Z".to_string(),
4154        )
4155        .await
4156        .unwrap();
4157
4158        let updated = repo
4159            .record_stellar_insufficient_fee_retry(tx_id, "2025-03-18T10:01:00Z".to_string())
4160            .await
4161            .unwrap();
4162
4163        assert_eq!(updated.sent_at.as_deref(), Some("2025-03-18T10:01:00Z"));
4164        let meta = updated.metadata.unwrap();
4165        assert_eq!(meta.insufficient_fee_retries, 2);
4166    }
4167
4168    #[tokio::test]
4169    #[ignore = "Requires active Redis instance"]
4170    async fn test_record_insufficient_fee_retry_noop_on_final_state() {
4171        let _lock = ENV_MUTEX.lock().await;
4172        let repo = setup_test_repo().await;
4173        let relayer_id = Uuid::new_v4().to_string();
4174        let tx_id = Uuid::new_v4().to_string();
4175        let mut tx =
4176            create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Confirmed);
4177        tx.sent_at = Some("old-time".to_string());
4178        repo.create(tx).await.unwrap();
4179
4180        let result = repo
4181            .record_stellar_insufficient_fee_retry(tx_id, "new-time".to_string())
4182            .await
4183            .unwrap();
4184
4185        // Should return unchanged on final state
4186        assert_eq!(result.sent_at.as_deref(), Some("old-time"));
4187        assert!(result.metadata.is_none());
4188    }
4189
4190    #[tokio::test]
4191    #[ignore = "Requires active Redis instance"]
4192    async fn test_record_insufficient_fee_retry_not_found() {
4193        let _lock = ENV_MUTEX.lock().await;
4194        let repo = setup_test_repo().await;
4195
4196        let result = repo
4197            .record_stellar_insufficient_fee_retry(
4198                "nonexistent".to_string(),
4199                "2025-03-18T10:00:00Z".to_string(),
4200            )
4201            .await;
4202
4203        assert!(matches!(result, Err(RepositoryError::NotFound(_))));
4204    }
4205
4206    // ── record_stellar_try_again_later_retry ────────────────────────
4207
4208    #[tokio::test]
4209    #[ignore = "Requires active Redis instance"]
4210    async fn test_record_try_again_later_retry() {
4211        let _lock = ENV_MUTEX.lock().await;
4212        let repo = setup_test_repo().await;
4213        let relayer_id = Uuid::new_v4().to_string();
4214        let tx_id = Uuid::new_v4().to_string();
4215        let mut tx =
4216            create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Sent);
4217        tx.sent_at = None;
4218        repo.create(tx).await.unwrap();
4219
4220        let updated = repo
4221            .record_stellar_try_again_later_retry(tx_id, "2025-03-18T10:00:00Z".to_string())
4222            .await
4223            .unwrap();
4224
4225        assert_eq!(updated.sent_at.as_deref(), Some("2025-03-18T10:00:00Z"));
4226        let meta = updated.metadata.unwrap();
4227        assert_eq!(meta.try_again_later_retries, 1);
4228        assert_eq!(meta.consecutive_failures, 0);
4229        assert_eq!(meta.total_failures, 0);
4230    }
4231
4232    #[tokio::test]
4233    #[ignore = "Requires active Redis instance"]
4234    async fn test_record_try_again_later_retry_accumulates() {
4235        let _lock = ENV_MUTEX.lock().await;
4236        let repo = setup_test_repo().await;
4237        let relayer_id = Uuid::new_v4().to_string();
4238        let tx_id = Uuid::new_v4().to_string();
4239        let tx = create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Sent);
4240        repo.create(tx).await.unwrap();
4241
4242        repo.record_stellar_try_again_later_retry(
4243            tx_id.clone(),
4244            "2025-03-18T10:00:00Z".to_string(),
4245        )
4246        .await
4247        .unwrap();
4248
4249        let updated = repo
4250            .record_stellar_try_again_later_retry(tx_id, "2025-03-18T10:01:00Z".to_string())
4251            .await
4252            .unwrap();
4253
4254        assert_eq!(updated.sent_at.as_deref(), Some("2025-03-18T10:01:00Z"));
4255        let meta = updated.metadata.unwrap();
4256        assert_eq!(meta.try_again_later_retries, 2);
4257    }
4258
4259    #[tokio::test]
4260    #[ignore = "Requires active Redis instance"]
4261    async fn test_record_try_again_later_retry_noop_on_final_state() {
4262        let _lock = ENV_MUTEX.lock().await;
4263        let repo = setup_test_repo().await;
4264        let relayer_id = Uuid::new_v4().to_string();
4265        let tx_id = Uuid::new_v4().to_string();
4266        let mut tx =
4267            create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Confirmed);
4268        tx.sent_at = Some("old-time".to_string());
4269        repo.create(tx).await.unwrap();
4270
4271        let result = repo
4272            .record_stellar_try_again_later_retry(tx_id, "new-time".to_string())
4273            .await
4274            .unwrap();
4275
4276        // Should return unchanged on final state
4277        assert_eq!(result.sent_at.as_deref(), Some("old-time"));
4278        assert!(result.metadata.is_none());
4279    }
4280
4281    #[tokio::test]
4282    #[ignore = "Requires active Redis instance"]
4283    async fn test_record_try_again_later_retry_not_found() {
4284        let _lock = ENV_MUTEX.lock().await;
4285        let repo = setup_test_repo().await;
4286
4287        let result = repo
4288            .record_stellar_try_again_later_retry(
4289                "nonexistent".to_string(),
4290                "2025-03-18T10:00:00Z".to_string(),
4291            )
4292            .await;
4293
4294        assert!(matches!(result, Err(RepositoryError::NotFound(_))));
4295    }
4296
4297    // ── metadata preservation across operations ─────────────────────
4298
4299    #[tokio::test]
4300    #[ignore = "Requires active Redis instance"]
4301    async fn test_increment_failures_preserves_try_again_later_retries() {
4302        let _lock = ENV_MUTEX.lock().await;
4303        let repo = setup_test_repo().await;
4304        let relayer_id = Uuid::new_v4().to_string();
4305        let tx_id = Uuid::new_v4().to_string();
4306        let tx = create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Sent);
4307        repo.create(tx).await.unwrap();
4308
4309        // Set try_again_later_retries = 1
4310        repo.record_stellar_try_again_later_retry(
4311            tx_id.clone(),
4312            "2025-03-18T10:00:00Z".to_string(),
4313        )
4314        .await
4315        .unwrap();
4316
4317        // Now increment failures — should NOT clobber try_again_later_retries
4318        let updated = repo.increment_status_check_failures(tx_id).await.unwrap();
4319
4320        let meta = updated.metadata.unwrap();
4321        assert_eq!(
4322            meta.try_again_later_retries, 1,
4323            "try_again_later_retries must survive increment_status_check_failures"
4324        );
4325        assert_eq!(meta.consecutive_failures, 1);
4326        assert_eq!(meta.total_failures, 1);
4327    }
4328
4329    #[tokio::test]
4330    #[ignore = "Requires active Redis instance"]
4331    async fn test_increment_failures_preserves_insufficient_fee_retries() {
4332        let _lock = ENV_MUTEX.lock().await;
4333        let repo = setup_test_repo().await;
4334        let relayer_id = Uuid::new_v4().to_string();
4335        let tx_id = Uuid::new_v4().to_string();
4336        let tx = create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Sent);
4337        repo.create(tx).await.unwrap();
4338
4339        // Set insufficient_fee_retries = 1
4340        repo.record_stellar_insufficient_fee_retry(
4341            tx_id.clone(),
4342            "2025-03-18T10:00:00Z".to_string(),
4343        )
4344        .await
4345        .unwrap();
4346
4347        // Now increment failures — should NOT clobber insufficient_fee_retries
4348        let updated = repo.increment_status_check_failures(tx_id).await.unwrap();
4349
4350        let meta = updated.metadata.unwrap();
4351        assert_eq!(
4352            meta.insufficient_fee_retries, 1,
4353            "insufficient_fee_retries must survive increment_status_check_failures"
4354        );
4355        assert_eq!(meta.consecutive_failures, 1);
4356    }
4357
4358    #[tokio::test]
4359    #[ignore = "Requires active Redis instance"]
4360    async fn test_reset_failures_preserves_retry_counters() {
4361        let _lock = ENV_MUTEX.lock().await;
4362        let repo = setup_test_repo().await;
4363        let relayer_id = Uuid::new_v4().to_string();
4364        let tx_id = Uuid::new_v4().to_string();
4365        let tx = create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Sent);
4366        repo.create(tx).await.unwrap();
4367
4368        // Set both retry counters
4369        repo.record_stellar_try_again_later_retry(
4370            tx_id.clone(),
4371            "2025-03-18T10:00:00Z".to_string(),
4372        )
4373        .await
4374        .unwrap();
4375        repo.record_stellar_insufficient_fee_retry(
4376            tx_id.clone(),
4377            "2025-03-18T10:01:00Z".to_string(),
4378        )
4379        .await
4380        .unwrap();
4381
4382        // Increment then reset consecutive failures
4383        repo.increment_status_check_failures(tx_id.clone())
4384            .await
4385            .unwrap();
4386        let updated = repo
4387            .reset_status_check_consecutive_failures(tx_id)
4388            .await
4389            .unwrap();
4390
4391        let meta = updated.metadata.unwrap();
4392        assert_eq!(meta.consecutive_failures, 0);
4393        assert_eq!(meta.total_failures, 1);
4394        assert_eq!(
4395            meta.try_again_later_retries, 1,
4396            "try_again_later_retries must survive reset"
4397        );
4398        assert_eq!(
4399            meta.insufficient_fee_retries, 1,
4400            "insufficient_fee_retries must survive reset"
4401        );
4402    }
4403
4404    #[tokio::test]
4405    #[ignore = "Requires active Redis instance"]
4406    async fn test_fee_and_try_again_later_retries_independent() {
4407        let _lock = ENV_MUTEX.lock().await;
4408        let repo = setup_test_repo().await;
4409        let relayer_id = Uuid::new_v4().to_string();
4410        let tx_id = Uuid::new_v4().to_string();
4411        let tx = create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Sent);
4412        repo.create(tx).await.unwrap();
4413
4414        // Set try_again_later_retries = 2
4415        repo.record_stellar_try_again_later_retry(
4416            tx_id.clone(),
4417            "2025-03-18T10:00:00Z".to_string(),
4418        )
4419        .await
4420        .unwrap();
4421        repo.record_stellar_try_again_later_retry(
4422            tx_id.clone(),
4423            "2025-03-18T10:01:00Z".to_string(),
4424        )
4425        .await
4426        .unwrap();
4427
4428        // Set insufficient_fee_retries = 1 — should NOT clobber try_again_later_retries
4429        let updated = repo
4430            .record_stellar_insufficient_fee_retry(tx_id, "2025-03-18T10:02:00Z".to_string())
4431            .await
4432            .unwrap();
4433
4434        let meta = updated.metadata.unwrap();
4435        assert_eq!(
4436            meta.try_again_later_retries, 2,
4437            "try_again_later_retries must survive insufficient_fee_retry"
4438        );
4439        assert_eq!(meta.insufficient_fee_retries, 1);
4440    }
4441}