1use crate::config::ServerConfig;
4use crate::constants::{ALL_TRANSACTION_STATUSES, 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::collections::{HashMap, HashSet};
26use std::fmt;
27use std::sync::Arc;
28use tracing::{debug, error, info, warn};
29
30const RELAYER_PREFIX: &str = "relayer";
31const TX_PREFIX: &str = "tx";
32const STATUS_PREFIX: &str = "status";
33const STATUS_SORTED_PREFIX: &str = "status_sorted";
34const NONCE_PREFIX: &str = "nonce";
35const TX_TO_RELAYER_PREFIX: &str = "tx_to_relayer";
36const RELAYER_LIST_KEY: &str = "relayer_list";
37const TX_BY_CREATED_AT_PREFIX: &str = "tx_by_created_at";
38const STALE_INDEX_MIN_AGE_SECS: i64 = 3600;
39const STALE_INDEX_PAGE_SIZE: usize = 100;
40const MAX_STALE_INDEX_RECONCILE_ITERATIONS_PER_STATUS: u32 = 1500;
41const NON_FINAL_STATUS_INDEXES: &[TransactionStatus] = &[
42 TransactionStatus::Pending,
43 TransactionStatus::Sent,
44 TransactionStatus::Submitted,
45 TransactionStatus::Mined,
46];
47
48#[derive(Clone, Copy, PartialEq)]
50enum ReconcileMode {
51 Promote,
54 PurgeOrphans,
58}
59
60#[derive(Clone)]
61pub struct RedisTransactionRepository {
62 pub connections: Arc<RedisConnections>,
63 pub key_prefix: String,
64}
65
66impl RedisRepository for RedisTransactionRepository {}
67
68impl RedisTransactionRepository {
69 pub fn new(
70 connections: Arc<RedisConnections>,
71 key_prefix: String,
72 ) -> Result<Self, RepositoryError> {
73 if key_prefix.is_empty() {
74 return Err(RepositoryError::InvalidData(
75 "Redis key prefix cannot be empty".to_string(),
76 ));
77 }
78
79 Ok(Self {
80 connections,
81 key_prefix,
82 })
83 }
84
85 fn tx_key(&self, relayer_id: &str, tx_id: &str) -> String {
87 format!(
88 "{}:{}:{}:{}:{}",
89 self.key_prefix, RELAYER_PREFIX, relayer_id, TX_PREFIX, tx_id
90 )
91 }
92
93 fn tx_to_relayer_key(&self, tx_id: &str) -> String {
95 format!(
96 "{}:{}:{}:{}",
97 self.key_prefix, RELAYER_PREFIX, TX_TO_RELAYER_PREFIX, tx_id
98 )
99 }
100
101 fn relayer_status_key(&self, relayer_id: &str, status: &TransactionStatus) -> String {
103 format!(
104 "{}:{}:{}{}",
105 self.key_prefix,
106 RELAYER_PREFIX,
107 relayer_id,
108 Self::relayer_status_key_suffix(status)
109 )
110 }
111
112 fn relayer_status_sorted_key(&self, relayer_id: &str, status: &TransactionStatus) -> String {
115 format!(
116 "{}:{}:{}{}",
117 self.key_prefix,
118 RELAYER_PREFIX,
119 relayer_id,
120 Self::relayer_status_sorted_key_suffix(status)
121 )
122 }
123
124 fn relayer_status_key_suffix(status: &TransactionStatus) -> String {
125 format!(":{STATUS_PREFIX}:{status}")
126 }
127
128 fn relayer_status_sorted_key_suffix(status: &TransactionStatus) -> String {
129 format!(":{STATUS_SORTED_PREFIX}:{status}")
130 }
131
132 fn relayer_nonce_key(&self, relayer_id: &str, nonce: u64) -> String {
134 format!(
135 "{}:{}:{}:{}:{}",
136 self.key_prefix, RELAYER_PREFIX, relayer_id, NONCE_PREFIX, nonce
137 )
138 }
139
140 fn relayer_list_key(&self) -> String {
142 format!("{}:{}", self.key_prefix, RELAYER_LIST_KEY)
143 }
144
145 fn relayer_tx_by_created_at_key(&self, relayer_id: &str) -> String {
147 format!(
148 "{}:{}:{}{}",
149 self.key_prefix,
150 RELAYER_PREFIX,
151 relayer_id,
152 Self::relayer_tx_by_created_at_key_suffix()
153 )
154 }
155
156 fn relayer_tx_by_created_at_key_suffix() -> String {
157 format!(":{TX_BY_CREATED_AT_PREFIX}")
158 }
159
160 fn tx_key_parts(&self, tx_id: &str) -> (String, String, String) {
165 let lookup_key = self.tx_to_relayer_key(tx_id);
166 let key_prefix = format!("{}:{}:", self.key_prefix, RELAYER_PREFIX);
167 let key_suffix = format!(":{TX_PREFIX}:{tx_id}");
168 (lookup_key, key_prefix, key_suffix)
169 }
170
171 async fn run_atomic_script(
178 &self,
179 lua: &str,
180 tx_id: &str,
181 extra_args: &[&str],
182 op_name: &str,
183 ) -> Result<TransactionRepoModel, RepositoryError> {
184 const MAX_RETRIES: u32 = 3;
185 const BASE_BACKOFF_MS: u64 = 100;
186
187 let (lookup_key, key_prefix, key_suffix) = self.tx_key_parts(tx_id);
188 let script = Script::new(lua);
189 let mut last_error = None;
190
191 for attempt in 0..MAX_RETRIES {
192 let backoff = BASE_BACKOFF_MS * 2u64.pow(attempt);
193
194 let mut conn = match self
195 .get_connection(self.connections.primary(), op_name)
196 .await
197 {
198 Ok(conn) => conn,
199 Err(e) => {
200 last_error = Some(e);
201 if attempt < MAX_RETRIES - 1 {
202 warn!(tx_id = %tx_id, attempt, op = %op_name, "connection failed, retrying");
203 tokio::time::sleep(tokio::time::Duration::from_millis(backoff)).await;
204 continue;
205 }
206 return Err(last_error.unwrap());
207 }
208 };
209
210 let mut invocation = script.prepare_invoke();
211 invocation
212 .key(&lookup_key)
213 .arg(&key_prefix)
214 .arg(&key_suffix);
215 for arg in extra_args {
216 invocation.arg(*arg);
217 }
218
219 match invocation.invoke_async::<Option<String>>(&mut conn).await {
220 Ok(result) => {
221 let json = result.ok_or_else(|| {
222 RepositoryError::NotFound(format!("Transaction with ID {tx_id} not found"))
223 })?;
224 return self.deserialize_entity::<TransactionRepoModel>(
225 &json,
226 tx_id,
227 "transaction",
228 );
229 }
230 Err(e) => {
231 last_error = Some(self.map_redis_error(e, op_name));
232 if attempt < MAX_RETRIES - 1 {
233 warn!(
234 tx_id = %tx_id, attempt, op = %op_name,
235 "atomic script failed, retrying"
236 );
237 tokio::time::sleep(tokio::time::Duration::from_millis(backoff)).await;
238 continue;
239 }
240 return Err(last_error.unwrap());
241 }
242 }
243 }
244 Err(last_error.unwrap_or_else(|| {
245 RepositoryError::UnexpectedError(format!("retry loop exhausted for {op_name}"))
246 }))
247 }
248
249 async fn run_script_with_retry_vec(
253 &self,
254 script: &Script,
255 lookup_key: &str,
256 key_prefix: &str,
257 key_suffix: &str,
258 extra_args: &[&str],
259 op_name: &str,
260 ) -> Result<Option<Vec<String>>, RepositoryError> {
261 const MAX_RETRIES: u32 = 3;
262 const BASE_BACKOFF_MS: u64 = 100;
263
264 let mut last_error = None;
265
266 for attempt in 0..MAX_RETRIES {
267 let backoff = BASE_BACKOFF_MS * 2u64.pow(attempt);
268
269 let mut conn = match self
270 .get_connection(self.connections.primary(), op_name)
271 .await
272 {
273 Ok(conn) => conn,
274 Err(e) => {
275 last_error = Some(e);
276 if attempt < MAX_RETRIES - 1 {
277 warn!(op = %op_name, attempt, "connection failed, retrying");
278 tokio::time::sleep(tokio::time::Duration::from_millis(backoff)).await;
279 continue;
280 }
281 return Err(last_error.unwrap());
282 }
283 };
284
285 let mut invocation = script.prepare_invoke();
286 invocation.key(lookup_key).arg(key_prefix).arg(key_suffix);
287 for arg in extra_args {
288 invocation.arg(*arg);
289 }
290
291 match invocation
294 .invoke_async::<Option<Vec<String>>>(&mut conn)
295 .await
296 {
297 Ok(result) => return Ok(result),
298 Err(e) => {
299 last_error = Some(self.map_redis_error(e, op_name));
300 if attempt < MAX_RETRIES - 1 {
301 warn!(op = %op_name, attempt, "script failed, retrying");
302 tokio::time::sleep(tokio::time::Duration::from_millis(backoff)).await;
303 continue;
304 }
305 return Err(last_error.unwrap());
306 }
307 }
308 }
309 Err(last_error.unwrap_or_else(|| {
310 RepositoryError::UnexpectedError(format!("retry loop exhausted for {op_name}"))
311 }))
312 }
313
314 fn timestamp_to_score(&self, timestamp: &str) -> f64 {
316 chrono::DateTime::parse_from_rfc3339(timestamp)
317 .map(|dt| dt.timestamp_millis() as f64)
318 .unwrap_or_else(|_| {
319 warn!(timestamp = %timestamp, "failed to parse timestamp, using 0");
320 0.0
321 })
322 }
323
324 fn status_sorted_score(&self, tx: &TransactionRepoModel) -> f64 {
328 if tx.status == TransactionStatus::Confirmed {
329 if let Some(ref confirmed_at) = tx.confirmed_at {
331 return self.timestamp_to_score(confirmed_at);
332 }
333 warn!(tx_id = %tx.id, "Confirmed transaction missing confirmed_at, using created_at");
335 }
336 self.timestamp_to_score(&tx.created_at)
337 }
338
339 fn status_json_key(status: &TransactionStatus) -> Result<String, RepositoryError> {
340 let value = serde_json::to_value(status).map_err(|e| {
341 RepositoryError::InvalidData(format!("Failed to serialize transaction status: {e}"))
342 })?;
343
344 value.as_str().map(str::to_string).ok_or_else(|| {
345 RepositoryError::InvalidData("Transaction status serialized as non-string".to_string())
346 })
347 }
348
349 fn partial_update_index_metadata(
350 &self,
351 tx_id: &str,
352 update: &TransactionUpdateRequest,
353 ) -> Result<String, RepositoryError> {
354 let mut sorted_suffixes = serde_json::Map::new();
355 let mut legacy_suffixes = serde_json::Map::new();
356
357 for status in ALL_TRANSACTION_STATUSES {
358 let status_key = Self::status_json_key(status)?;
359 sorted_suffixes.insert(
360 status_key.clone(),
361 serde_json::Value::String(Self::relayer_status_sorted_key_suffix(status)),
362 );
363 legacy_suffixes.insert(
364 status_key,
365 serde_json::Value::String(Self::relayer_status_key_suffix(status)),
366 );
367 }
368
369 let nonfinal_statuses = NON_FINAL_STATUS_INDEXES
370 .iter()
371 .map(Self::status_json_key)
372 .collect::<Result<Vec<_>, _>>()?;
373
374 let confirmed_score = if update.status.as_ref() == Some(&TransactionStatus::Confirmed) {
375 update
376 .confirmed_at
377 .as_deref()
378 .map(|confirmed_at| self.timestamp_to_score(confirmed_at).to_string())
379 .unwrap_or_default()
380 } else {
381 String::new()
382 };
383
384 serde_json::to_string(&serde_json::json!({
385 "sorted": sorted_suffixes,
386 "legacy": legacy_suffixes,
387 "nonfinal": nonfinal_statuses,
388 "created_key_suffix": Self::relayer_tx_by_created_at_key_suffix(),
389 "confirmed_score": confirmed_score,
390 "tx_id": tx_id,
391 }))
392 .map_err(|e| {
393 RepositoryError::InvalidData(format!(
394 "Failed to serialize partial update index metadata: {e}"
395 ))
396 })
397 }
398
399 async fn reconcile_status_index(
400 &self,
401 relayer_id: &str,
402 index_status: &TransactionStatus,
403 cutoff_ms: i64,
404 mode: ReconcileMode,
405 ) -> Result<usize, RepositoryError> {
406 let sorted_key = self.relayer_status_sorted_key(relayer_id, index_status);
407 let legacy_key = self.relayer_status_key(relayer_id, index_status);
408 let mut offset = 0usize;
409 let mut total_repaired = 0usize;
410 let mut iterations = 0u32;
411
412 loop {
413 if iterations >= MAX_STALE_INDEX_RECONCILE_ITERATIONS_PER_STATUS {
414 warn!(
415 relayer_id = %relayer_id,
416 status = %index_status,
417 iterations,
418 total_repaired,
419 "reached max stale status-index reconciliation iterations, stopping"
420 );
421 break;
422 }
423 iterations += 1;
424
425 let ids: Vec<String> = {
426 let mut conn = self
427 .get_connection(
428 self.connections.primary(),
429 "reconcile_stale_status_indexes_scan",
430 )
431 .await?;
432
433 redis::cmd("ZRANGEBYSCORE")
434 .arg(&sorted_key)
435 .arg("-inf")
436 .arg(cutoff_ms)
437 .arg("LIMIT")
438 .arg(offset)
439 .arg(STALE_INDEX_PAGE_SIZE)
440 .query_async(&mut conn)
441 .await
442 .map_err(|e| self.map_redis_error(e, "reconcile_stale_status_indexes_scan"))?
443 };
444
445 if ids.is_empty() {
446 break;
447 }
448
449 let batch = self.get_transactions_by_ids(&ids).await?;
450 let failed_ids: HashSet<&str> = batch.failed_ids.iter().map(String::as_str).collect();
451 let transactions_by_id: HashMap<&str, &TransactionRepoModel> = batch
452 .results
453 .iter()
454 .map(|tx| (tx.id.as_str(), tx))
455 .collect();
456
457 let mut pipe = redis::pipe();
458 pipe.atomic();
459 let mut repaired_on_page = 0usize;
460
461 for tx_id in &ids {
462 if failed_ids.contains(tx_id.as_str()) {
463 pipe.zrem(&sorted_key, tx_id);
464 pipe.srem(&legacy_key, tx_id);
465 repaired_on_page += 1;
466 continue;
467 }
468
469 if mode == ReconcileMode::PurgeOrphans {
474 continue;
475 }
476
477 let Some(tx) = transactions_by_id.get(tx_id.as_str()) else {
478 warn!(
479 relayer_id = %relayer_id,
480 status = %index_status,
481 tx_id = %tx_id,
482 "transaction body could not be reconciled from stale status index"
483 );
484 continue;
485 };
486
487 if FINAL_TRANSACTION_STATUSES.contains(&tx.status) && tx.status != *index_status {
488 let final_status_key =
489 self.relayer_status_sorted_key(&tx.relayer_id, &tx.status);
490 let final_status_score = self.status_sorted_score(tx);
491
492 pipe.zrem(&sorted_key, tx_id);
493 pipe.srem(&legacy_key, tx_id);
494 pipe.zadd(&final_status_key, tx_id, final_status_score);
495 repaired_on_page += 1;
496 } else if tx.status != *index_status {
497 warn!(
500 relayer_id = %relayer_id,
501 index_status = %index_status,
502 body_status = %tx.status,
503 tx_id = %tx_id,
504 "skipping non-final transaction status-index mismatch"
505 );
506 }
507 }
508
509 if repaired_on_page > 0 {
510 let mut conn = self
511 .get_connection(
512 self.connections.primary(),
513 "reconcile_stale_status_indexes_update",
514 )
515 .await?;
516 pipe.query_async::<()>(&mut conn).await.map_err(|e| {
517 self.map_redis_error(e, "reconcile_stale_status_indexes_update")
518 })?;
519 total_repaired += repaired_on_page;
520 } else {
521 offset += ids.len();
522 }
523 }
524
525 Ok(total_repaired)
526 }
527
528 async fn get_transactions_by_ids(
530 &self,
531 ids: &[String],
532 ) -> Result<BatchRetrievalResult<TransactionRepoModel>, RepositoryError> {
533 if ids.is_empty() {
534 debug!("no transaction IDs provided for batch fetch");
535 return Ok(BatchRetrievalResult {
536 results: vec![],
537 failed_ids: vec![],
538 });
539 }
540
541 let mut conn = self
542 .get_connection(self.connections.reader(), "batch_fetch_transactions")
543 .await?;
544
545 let reverse_keys: Vec<String> = ids.iter().map(|id| self.tx_to_relayer_key(id)).collect();
546
547 debug!(count = %ids.len(), "fetching relayer IDs for transactions");
548
549 let relayer_ids: Vec<Option<String>> = conn
550 .mget(&reverse_keys)
551 .await
552 .map_err(|e| self.map_redis_error(e, "batch_fetch_relayer_ids"))?;
553
554 let mut tx_keys = Vec::new();
555 let mut valid_ids = Vec::new();
556 let mut failed_ids = Vec::new();
557 for (i, relayer_id) in relayer_ids.into_iter().enumerate() {
558 match relayer_id {
559 Some(relayer_id) => {
560 tx_keys.push(self.tx_key(&relayer_id, &ids[i]));
561 valid_ids.push(ids[i].clone());
562 }
563 None => {
564 warn!(tx_id = %ids[i], "no relayer found for transaction");
565 failed_ids.push(ids[i].clone());
566 }
567 }
568 }
569
570 if tx_keys.is_empty() {
571 debug!("no valid transactions found for batch fetch");
572 return Ok(BatchRetrievalResult {
573 results: vec![],
574 failed_ids,
575 });
576 }
577
578 debug!(count = %tx_keys.len(), "batch fetching transaction data");
579
580 let values: Vec<Option<String>> = conn
581 .mget(&tx_keys)
582 .await
583 .map_err(|e| self.map_redis_error(e, "batch_fetch_transactions"))?;
584
585 let mut transactions = Vec::new();
586 let mut failed_count = 0;
587 for (i, value) in values.into_iter().enumerate() {
588 match value {
589 Some(json) => {
590 match self.deserialize_entity::<TransactionRepoModel>(
591 &json,
592 &valid_ids[i],
593 "transaction",
594 ) {
595 Ok(tx) => transactions.push(tx),
596 Err(e) => {
597 failed_count += 1;
598 error!(tx_id = %valid_ids[i], error = %e, "failed to deserialize transaction");
599 }
601 }
602 }
603 None => {
604 warn!(tx_id = %valid_ids[i], "transaction not found in batch fetch");
605 failed_ids.push(valid_ids[i].clone());
606 }
607 }
608 }
609
610 if failed_count > 0 {
611 warn!(failed_count = %failed_count, total_count = %valid_ids.len(), "failed to deserialize transactions in batch");
612 }
613
614 debug!(count = %transactions.len(), "successfully fetched transactions");
615 Ok(BatchRetrievalResult {
616 results: transactions,
617 failed_ids,
618 })
619 }
620
621 fn extract_nonce(&self, network_data: &NetworkTransactionData) -> Option<u64> {
623 match network_data.get_evm_transaction_data() {
624 Ok(tx_data) => tx_data.nonce,
625 Err(_) => {
626 debug!("no EVM transaction data available for nonce extraction");
627 None
628 }
629 }
630 }
631
632 async fn ensure_status_sorted_set(
649 &self,
650 relayer_id: &str,
651 status: &TransactionStatus,
652 ) -> Result<u64, RepositoryError> {
653 let sorted_key = self.relayer_status_sorted_key(relayer_id, status);
654 let legacy_key = self.relayer_status_key(relayer_id, status);
655
656 let legacy_ids = {
658 let mut conn = self
659 .get_connection(self.connections.primary(), "ensure_status_sorted_set_check")
660 .await?;
661
662 let legacy_count: u64 = conn
664 .scard(&legacy_key)
665 .await
666 .map_err(|e| self.map_redis_error(e, "ensure_status_sorted_set_scard"))?;
667
668 if legacy_count == 0 {
669 let sorted_count: u64 = conn
671 .zcard(&sorted_key)
672 .await
673 .map_err(|e| self.map_redis_error(e, "ensure_status_sorted_set_zcard"))?;
674 return Ok(sorted_count);
675 }
676
677 debug!(
679 relayer_id = %relayer_id,
680 status = %status,
681 legacy_count = %legacy_count,
682 "migrating status set to sorted set"
683 );
684
685 let ids: Vec<String> = conn
686 .smembers(&legacy_key)
687 .await
688 .map_err(|e| self.map_redis_error(e, "ensure_status_sorted_set_smembers"))?;
689
690 ids
691 };
693
694 if legacy_ids.is_empty() {
695 return Ok(0);
696 }
697
698 let transactions = self.get_transactions_by_ids(&legacy_ids).await?;
700
701 let mut conn = self
703 .get_connection(
704 self.connections.primary(),
705 "ensure_status_sorted_set_migrate",
706 )
707 .await?;
708
709 if transactions.results.is_empty() {
710 let _: () = conn
712 .del(&legacy_key)
713 .await
714 .map_err(|e| self.map_redis_error(e, "ensure_status_sorted_set_del_stale"))?;
715 return Ok(0);
716 }
717
718 let mut pipe = redis::pipe();
721 pipe.atomic();
722
723 for tx in &transactions.results {
724 let score = self.status_sorted_score(tx);
725 pipe.zadd(&sorted_key, &tx.id, score);
726 }
727
728 pipe.del(&legacy_key);
730
731 pipe.query_async::<()>(&mut conn)
732 .await
733 .map_err(|e| self.map_redis_error(e, "ensure_status_sorted_set_migrate"))?;
734
735 let migrated_count = transactions.results.len() as u64;
736 debug!(
737 relayer_id = %relayer_id,
738 status = %status,
739 migrated_count = %migrated_count,
740 "completed migration of status set to sorted set"
741 );
742
743 Ok(migrated_count)
744 }
745
746 async fn update_indexes(
753 &self,
754 tx: &TransactionRepoModel,
755 old_tx: Option<&TransactionRepoModel>,
756 include_status_indexes: bool,
757 ) -> Result<(), RepositoryError> {
758 let mut conn = self
759 .get_connection(self.connections.primary(), "update_indexes")
760 .await?;
761 let mut pipe = redis::pipe();
762 pipe.atomic();
763
764 debug!(tx_id = %tx.id, "updating indexes for transaction");
765
766 let relayer_list_key = self.relayer_list_key();
768 pipe.sadd(&relayer_list_key, &tx.relayer_id);
769
770 let created_at_score = self.timestamp_to_score(&tx.created_at);
772
773 if include_status_indexes {
774 let status_score = self.status_sorted_score(tx);
776
777 let new_status_sorted_key = self.relayer_status_sorted_key(&tx.relayer_id, &tx.status);
779 pipe.zadd(&new_status_sorted_key, &tx.id, status_score);
780 debug!(tx_id = %tx.id, status = %tx.status, score = %status_score, "added transaction to status sorted set");
781 }
782
783 if let Some(nonce) = self.extract_nonce(&tx.network_data) {
784 let nonce_key = self.relayer_nonce_key(&tx.relayer_id, nonce);
785 pipe.set(&nonce_key, &tx.id);
786 debug!(tx_id = %tx.id, nonce = %nonce, "added nonce index for transaction");
787 }
788
789 let relayer_sorted_key = self.relayer_tx_by_created_at_key(&tx.relayer_id);
791 pipe.zadd(&relayer_sorted_key, &tx.id, created_at_score);
792 debug!(tx_id = %tx.id, score = %created_at_score, "added transaction to sorted set by created_at");
793
794 if let Some(old) = old_tx {
796 if include_status_indexes && old.status != tx.status {
797 let old_status_sorted_key =
799 self.relayer_status_sorted_key(&old.relayer_id, &old.status);
800 pipe.zrem(&old_status_sorted_key, &tx.id);
801
802 let old_status_legacy_key = self.relayer_status_key(&old.relayer_id, &old.status);
804 pipe.srem(&old_status_legacy_key, &tx.id);
805
806 debug!(tx_id = %tx.id, old_status = %old.status, new_status = %tx.status, "removing old status indexes for transaction");
807 }
808
809 if let Some(old_nonce) = self.extract_nonce(&old.network_data) {
811 let new_nonce = self.extract_nonce(&tx.network_data);
812 if Some(old_nonce) != new_nonce {
813 let old_nonce_key = self.relayer_nonce_key(&old.relayer_id, old_nonce);
814 pipe.del(&old_nonce_key);
815 debug!(tx_id = %tx.id, old_nonce = %old_nonce, new_nonce = ?new_nonce, "removing old nonce index for transaction");
816 }
817 }
818 }
819
820 pipe.exec_async(&mut conn).await.map_err(|e| {
822 error!(tx_id = %tx.id, error = %e, "index update pipeline failed for transaction");
823 self.map_redis_error(e, &format!("update_indexes_for_tx_{}", tx.id))
824 })?;
825
826 debug!(tx_id = %tx.id, "successfully updated indexes for transaction");
827 Ok(())
828 }
829
830 async fn remove_all_indexes(&self, tx: &TransactionRepoModel) -> Result<(), RepositoryError> {
832 let mut conn = self
833 .get_connection(self.connections.primary(), "remove_all_indexes")
834 .await?;
835 let mut pipe = redis::pipe();
836 pipe.atomic();
837
838 debug!(tx_id = %tx.id, "removing all indexes for transaction");
839
840 for status in &[
844 TransactionStatus::Canceled,
845 TransactionStatus::Pending,
846 TransactionStatus::Sent,
847 TransactionStatus::Submitted,
848 TransactionStatus::Mined,
849 TransactionStatus::Confirmed,
850 TransactionStatus::Failed,
851 TransactionStatus::Expired,
852 ] {
853 let status_sorted_key = self.relayer_status_sorted_key(&tx.relayer_id, status);
855 pipe.zrem(&status_sorted_key, &tx.id);
856
857 let status_legacy_key = self.relayer_status_key(&tx.relayer_id, status);
859 pipe.srem(&status_legacy_key, &tx.id);
860 }
861
862 if let Some(nonce) = self.extract_nonce(&tx.network_data) {
864 let nonce_key = self.relayer_nonce_key(&tx.relayer_id, nonce);
865 pipe.del(&nonce_key);
866 debug!(tx_id = %tx.id, nonce = %nonce, "removing nonce index for transaction");
867 }
868
869 let relayer_sorted_key = self.relayer_tx_by_created_at_key(&tx.relayer_id);
871 pipe.zrem(&relayer_sorted_key, &tx.id);
872 debug!(tx_id = %tx.id, "removing transaction from sorted set by created_at");
873
874 let reverse_key = self.tx_to_relayer_key(&tx.id);
876 pipe.del(&reverse_key);
877
878 pipe.exec_async(&mut conn).await.map_err(|e| {
879 error!(tx_id = %tx.id, error = %e, "index removal failed for transaction");
880 self.map_redis_error(e, &format!("remove_indexes_for_tx_{}", tx.id))
881 })?;
882
883 debug!(tx_id = %tx.id, "successfully removed all indexes for transaction");
884 Ok(())
885 }
886
887 fn track_status_change_metrics(
889 &self,
890 _original_tx: &TransactionRepoModel,
891 updated_tx: &TransactionRepoModel,
892 old_status: &TransactionStatus,
893 new_status: &TransactionStatus,
894 ) {
895 let network_type = format!("{:?}", updated_tx.network_type).to_lowercase();
896 let relayer_id = updated_tx.relayer_id.as_str();
897
898 if *old_status != TransactionStatus::Submitted
900 && *new_status == TransactionStatus::Submitted
901 {
902 TRANSACTIONS_SUBMITTED
903 .with_label_values(&[relayer_id, &network_type])
904 .inc();
905
906 if let Ok(created_time) = chrono::DateTime::parse_from_rfc3339(&updated_tx.created_at) {
907 let processing_seconds =
908 (Utc::now() - created_time.with_timezone(&Utc)).num_seconds() as f64;
909 TRANSACTION_PROCESSING_TIME
910 .with_label_values(&[relayer_id, &network_type, "creation_to_submission"])
911 .observe(processing_seconds);
912 }
913 }
914
915 if old_status != new_status {
917 let old_status_str = format!("{old_status:?}").to_lowercase();
918 let old_status_gauge = TRANSACTIONS_BY_STATUS.with_label_values(&[
919 relayer_id,
920 &network_type,
921 &old_status_str,
922 ]);
923 let clamped_value = (old_status_gauge.get() - 1.0).max(0.0);
924 old_status_gauge.set(clamped_value);
925
926 let new_status_str = format!("{new_status:?}").to_lowercase();
927 TRANSACTIONS_BY_STATUS
928 .with_label_values(&[relayer_id, &network_type, &new_status_str])
929 .inc();
930 }
931
932 let was_final = is_final_state(old_status);
934 let is_final = is_final_state(new_status);
935
936 if !was_final && is_final {
937 let previous_status = format!("{old_status:?}").to_lowercase();
938 let meta = updated_tx.metadata.as_ref();
939 let had_insufficient_fee = meta.is_some_and(|m| m.insufficient_fee_retries > 0);
940 let had_try_again_later = meta.is_some_and(|m| m.try_again_later_retries > 0);
941
942 match new_status {
943 TransactionStatus::Confirmed => {
944 TRANSACTIONS_SUCCESS
945 .with_label_values(&[relayer_id, &network_type])
946 .inc();
947 if had_insufficient_fee {
948 TRANSACTIONS_INSUFFICIENT_FEE_SUCCESS
949 .with_label_values(&[relayer_id, &network_type])
950 .inc();
951 }
952 if had_try_again_later {
953 TRANSACTIONS_TRY_AGAIN_LATER_SUCCESS
954 .with_label_values(&[relayer_id, &network_type])
955 .inc();
956 }
957
958 if let (Some(sent_at_str), Some(confirmed_at_str)) =
959 (&updated_tx.sent_at, &updated_tx.confirmed_at)
960 {
961 if let (Ok(sent_time), Ok(confirmed_time)) = (
962 chrono::DateTime::parse_from_rfc3339(sent_at_str),
963 chrono::DateTime::parse_from_rfc3339(confirmed_at_str),
964 ) {
965 let processing_seconds = (confirmed_time.with_timezone(&Utc)
966 - sent_time.with_timezone(&Utc))
967 .num_seconds()
968 as f64;
969 TRANSACTION_PROCESSING_TIME
970 .with_label_values(&[
971 relayer_id,
972 &network_type,
973 "submission_to_confirmation",
974 ])
975 .observe(processing_seconds);
976 }
977 }
978
979 if let Ok(created_time) =
980 chrono::DateTime::parse_from_rfc3339(&updated_tx.created_at)
981 {
982 if let Some(confirmed_at_str) = &updated_tx.confirmed_at {
983 if let Ok(confirmed_time) =
984 chrono::DateTime::parse_from_rfc3339(confirmed_at_str)
985 {
986 let processing_seconds = (confirmed_time.with_timezone(&Utc)
987 - created_time.with_timezone(&Utc))
988 .num_seconds()
989 as f64;
990 TRANSACTION_PROCESSING_TIME
991 .with_label_values(&[
992 relayer_id,
993 &network_type,
994 "creation_to_confirmation",
995 ])
996 .observe(processing_seconds);
997 }
998 }
999 }
1000 }
1001 TransactionStatus::Failed => {
1002 let failure_reason = updated_tx
1003 .status_reason
1004 .as_deref()
1005 .map(|reason| {
1006 if reason.starts_with("Submission failed:") {
1007 "submission_failed"
1008 } else if reason.starts_with("Preparation failed:") {
1009 "preparation_failed"
1010 } else {
1011 "failed"
1012 }
1013 })
1014 .unwrap_or("failed");
1015 TRANSACTIONS_FAILED
1016 .with_label_values(&[
1017 relayer_id,
1018 &network_type,
1019 failure_reason,
1020 &previous_status,
1021 ])
1022 .inc();
1023 }
1024 TransactionStatus::Expired => {
1025 TRANSACTIONS_FAILED
1026 .with_label_values(&[
1027 relayer_id,
1028 &network_type,
1029 "expired",
1030 &previous_status,
1031 ])
1032 .inc();
1033 }
1034 TransactionStatus::Canceled => {
1035 TRANSACTIONS_FAILED
1036 .with_label_values(&[
1037 relayer_id,
1038 &network_type,
1039 "canceled",
1040 &previous_status,
1041 ])
1042 .inc();
1043 }
1044 _ => {}
1045 }
1046
1047 if *new_status != TransactionStatus::Confirmed {
1049 if had_insufficient_fee {
1050 TRANSACTIONS_INSUFFICIENT_FEE_FAILED
1051 .with_label_values(&[relayer_id, &network_type])
1052 .inc();
1053 }
1054 if had_try_again_later {
1055 TRANSACTIONS_TRY_AGAIN_LATER_FAILED
1056 .with_label_values(&[relayer_id, &network_type])
1057 .inc();
1058 }
1059 }
1060 }
1061 }
1062}
1063
1064impl fmt::Debug for RedisTransactionRepository {
1065 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1066 f.debug_struct("RedisTransactionRepository")
1067 .field("connections", &"<RedisConnections>")
1068 .field("key_prefix", &self.key_prefix)
1069 .finish()
1070 }
1071}
1072
1073#[async_trait]
1074impl Repository<TransactionRepoModel, String> for RedisTransactionRepository {
1075 async fn create(
1076 &self,
1077 entity: TransactionRepoModel,
1078 ) -> Result<TransactionRepoModel, RepositoryError> {
1079 if entity.id.is_empty() {
1080 return Err(RepositoryError::InvalidData(
1081 "Transaction ID cannot be empty".to_string(),
1082 ));
1083 }
1084
1085 let key = self.tx_key(&entity.relayer_id, &entity.id);
1086 let reverse_key = self.tx_to_relayer_key(&entity.id);
1087 let mut conn = self
1088 .get_connection(self.connections.primary(), "create")
1089 .await?;
1090
1091 debug!(tx_id = %entity.id, "creating transaction");
1092
1093 let value = self.serialize_entity(&entity, |t| &t.id, "transaction")?;
1094
1095 let existing: Option<String> = conn
1097 .get(&reverse_key)
1098 .await
1099 .map_err(|e| self.map_redis_error(e, "create_transaction_check"))?;
1100
1101 if existing.is_some() {
1102 return Err(RepositoryError::ConstraintViolation(format!(
1103 "Transaction with ID {} already exists",
1104 entity.id
1105 )));
1106 }
1107
1108 let mut pipe = redis::pipe();
1110 pipe.atomic();
1111 pipe.set(&key, &value);
1112 pipe.set(&reverse_key, &entity.relayer_id);
1113
1114 pipe.exec_async(&mut conn)
1115 .await
1116 .map_err(|e| self.map_redis_error(e, "create_transaction"))?;
1117
1118 if let Err(e) = self.update_indexes(&entity, None, true).await {
1120 error!(tx_id = %entity.id, error = %e, "failed to update indexes for new transaction");
1121 return Err(e);
1122 }
1123
1124 let network_type = format!("{:?}", entity.network_type).to_lowercase();
1126 let relayer_id = entity.relayer_id.as_str();
1127 TRANSACTIONS_CREATED
1128 .with_label_values(&[relayer_id, &network_type])
1129 .inc();
1130
1131 let status = &entity.status;
1133 let status_str = format!("{status:?}").to_lowercase();
1134 TRANSACTIONS_BY_STATUS
1135 .with_label_values(&[relayer_id, &network_type, &status_str])
1136 .inc();
1137
1138 debug!(tx_id = %entity.id, "successfully created transaction");
1139 Ok(entity)
1140 }
1141
1142 async fn get_by_id(&self, id: String) -> Result<TransactionRepoModel, RepositoryError> {
1143 if id.is_empty() {
1144 return Err(RepositoryError::InvalidData(
1145 "Transaction ID cannot be empty".to_string(),
1146 ));
1147 }
1148
1149 let mut conn = self
1150 .get_connection(self.connections.reader(), "get_by_id")
1151 .await?;
1152
1153 debug!(tx_id = %id, "fetching transaction");
1154
1155 let reverse_key = self.tx_to_relayer_key(&id);
1156 let relayer_id: Option<String> = conn
1157 .get(&reverse_key)
1158 .await
1159 .map_err(|e| self.map_redis_error(e, "get_transaction_reverse_lookup"))?;
1160
1161 let relayer_id = match relayer_id {
1162 Some(relayer_id) => relayer_id,
1163 None => {
1164 debug!(tx_id = %id, "transaction not found (no reverse lookup)");
1165 return Err(RepositoryError::NotFound(format!(
1166 "Transaction with ID {id} not found"
1167 )));
1168 }
1169 };
1170
1171 let key = self.tx_key(&relayer_id, &id);
1172 let value: Option<String> = conn
1173 .get(&key)
1174 .await
1175 .map_err(|e| self.map_redis_error(e, "get_transaction_by_id"))?;
1176
1177 match value {
1178 Some(json) => {
1179 let tx =
1180 self.deserialize_entity::<TransactionRepoModel>(&json, &id, "transaction")?;
1181 debug!(tx_id = %id, "successfully fetched transaction");
1182 Ok(tx)
1183 }
1184 None => {
1185 debug!(tx_id = %id, "transaction not found");
1186 Err(RepositoryError::NotFound(format!(
1187 "Transaction with ID {id} not found"
1188 )))
1189 }
1190 }
1191 }
1192
1193 async fn list_all(&self) -> Result<Vec<TransactionRepoModel>, RepositoryError> {
1195 let mut conn = self
1196 .get_connection(self.connections.reader(), "list_all")
1197 .await?;
1198
1199 debug!("fetching all transactions sorted by created_at (newest first)");
1200
1201 let relayer_list_key = self.relayer_list_key();
1203 let relayer_ids: Vec<String> = conn
1204 .smembers(&relayer_list_key)
1205 .await
1206 .map_err(|e| self.map_redis_error(e, "list_all_relayer_ids"))?;
1207
1208 debug!(count = %relayer_ids.len(), "found relayers");
1209
1210 let mut all_tx_ids = Vec::new();
1212 for relayer_id in relayer_ids {
1213 let relayer_sorted_key = self.relayer_tx_by_created_at_key(&relayer_id);
1214 let tx_ids: Vec<String> = redis::cmd("ZRANGE")
1215 .arg(&relayer_sorted_key)
1216 .arg(0)
1217 .arg(-1)
1218 .arg("REV")
1219 .query_async(&mut conn)
1220 .await
1221 .map_err(|e| self.map_redis_error(e, "list_all_relayer_sorted"))?;
1222
1223 all_tx_ids.extend(tx_ids);
1224 }
1225
1226 drop(conn);
1228
1229 let batch_result = self.get_transactions_by_ids(&all_tx_ids).await?;
1231 let mut all_transactions = batch_result.results;
1232
1233 all_transactions.sort_by(|a, b| b.created_at.cmp(&a.created_at));
1235
1236 debug!(count = %all_transactions.len(), "found transactions");
1237 Ok(all_transactions)
1238 }
1239
1240 async fn list_paginated(
1242 &self,
1243 query: PaginationQuery,
1244 ) -> Result<PaginatedResult<TransactionRepoModel>, RepositoryError> {
1245 if query.per_page == 0 {
1246 return Err(RepositoryError::InvalidData(
1247 "per_page must be greater than 0".to_string(),
1248 ));
1249 }
1250
1251 let mut conn = self
1252 .get_connection(self.connections.reader(), "list_paginated")
1253 .await?;
1254
1255 debug!(page = %query.page, per_page = %query.per_page, "fetching paginated transactions sorted by created_at (newest first)");
1256
1257 let relayer_list_key = self.relayer_list_key();
1259 let relayer_ids: Vec<String> = conn
1260 .smembers(&relayer_list_key)
1261 .await
1262 .map_err(|e| self.map_redis_error(e, "list_paginated_relayer_ids"))?;
1263
1264 let mut all_tx_ids = Vec::new();
1266 for relayer_id in relayer_ids {
1267 let relayer_sorted_key = self.relayer_tx_by_created_at_key(&relayer_id);
1268 let tx_ids: Vec<String> = redis::cmd("ZRANGE")
1269 .arg(&relayer_sorted_key)
1270 .arg(0)
1271 .arg(-1)
1272 .arg("REV")
1273 .query_async(&mut conn)
1274 .await
1275 .map_err(|e| self.map_redis_error(e, "list_paginated_relayer_sorted"))?;
1276
1277 all_tx_ids.extend(tx_ids);
1278 }
1279
1280 drop(conn);
1282
1283 let batch_result = self.get_transactions_by_ids(&all_tx_ids).await?;
1285 let mut all_transactions = batch_result.results;
1286
1287 all_transactions.sort_by(|a, b| b.created_at.cmp(&a.created_at));
1289
1290 let total = all_transactions.len() as u64;
1291 let start = ((query.page - 1) * query.per_page) as usize;
1292 let end = (start + query.per_page as usize).min(all_transactions.len());
1293
1294 if start >= all_transactions.len() {
1295 debug!(page = %query.page, total = %total, "page is beyond available data");
1296 return Ok(PaginatedResult {
1297 items: vec![],
1298 total,
1299 page: query.page,
1300 per_page: query.per_page,
1301 });
1302 }
1303
1304 let items = all_transactions[start..end].to_vec();
1305
1306 debug!(count = %items.len(), page = %query.page, "successfully fetched transactions for page");
1307
1308 Ok(PaginatedResult {
1309 items,
1310 total,
1311 page: query.page,
1312 per_page: query.per_page,
1313 })
1314 }
1315
1316 async fn update(
1317 &self,
1318 id: String,
1319 entity: TransactionRepoModel,
1320 ) -> Result<TransactionRepoModel, RepositoryError> {
1321 if id.is_empty() {
1322 return Err(RepositoryError::InvalidData(
1323 "Transaction ID cannot be empty".to_string(),
1324 ));
1325 }
1326
1327 debug!(tx_id = %id, "updating transaction");
1328
1329 let old_tx = self.get_by_id(id.clone()).await?;
1331
1332 let key = self.tx_key(&entity.relayer_id, &id);
1333 let mut conn = self
1334 .get_connection(self.connections.primary(), "update")
1335 .await?;
1336
1337 let value = self.serialize_entity(&entity, |t| &t.id, "transaction")?;
1338
1339 let _: () = conn
1341 .set(&key, value)
1342 .await
1343 .map_err(|e| self.map_redis_error(e, "update_transaction"))?;
1344
1345 self.update_indexes(&entity, Some(&old_tx), true).await?;
1347
1348 debug!(tx_id = %id, "successfully updated transaction");
1349 Ok(entity)
1350 }
1351
1352 async fn delete_by_id(&self, id: String) -> Result<(), RepositoryError> {
1353 if id.is_empty() {
1354 return Err(RepositoryError::InvalidData(
1355 "Transaction ID cannot be empty".to_string(),
1356 ));
1357 }
1358
1359 debug!(tx_id = %id, "deleting transaction");
1360
1361 let tx = self.get_by_id(id.clone()).await?;
1363
1364 let key = self.tx_key(&tx.relayer_id, &id);
1365 let reverse_key = self.tx_to_relayer_key(&id);
1366 let mut conn = self
1367 .get_connection(self.connections.primary(), "delete_by_id")
1368 .await?;
1369
1370 let mut pipe = redis::pipe();
1371 pipe.atomic();
1372 pipe.del(&key);
1373 pipe.del(&reverse_key);
1374
1375 pipe.exec_async(&mut conn)
1376 .await
1377 .map_err(|e| self.map_redis_error(e, "delete_transaction"))?;
1378
1379 if let Err(e) = self.remove_all_indexes(&tx).await {
1381 error!(tx_id = %id, error = %e, "failed to remove indexes for deleted transaction");
1382 }
1383
1384 debug!(tx_id = %id, "successfully deleted transaction");
1385 Ok(())
1386 }
1387
1388 async fn count(&self) -> Result<usize, RepositoryError> {
1390 let mut conn = self
1391 .get_connection(self.connections.reader(), "count")
1392 .await?;
1393
1394 debug!("counting transactions");
1395
1396 let relayer_list_key = self.relayer_list_key();
1398 let relayer_ids: Vec<String> = conn
1399 .smembers(&relayer_list_key)
1400 .await
1401 .map_err(|e| self.map_redis_error(e, "count_relayer_ids"))?;
1402
1403 let mut total_count = 0usize;
1404 for relayer_id in relayer_ids {
1405 let relayer_sorted_key = self.relayer_tx_by_created_at_key(&relayer_id);
1406 let count: usize = conn
1407 .zcard(&relayer_sorted_key)
1408 .await
1409 .map_err(|e| self.map_redis_error(e, "count_relayer_transactions"))?;
1410 total_count += count;
1411 }
1412
1413 debug!(count = %total_count, "transaction count");
1414 Ok(total_count)
1415 }
1416
1417 async fn has_entries(&self) -> Result<bool, RepositoryError> {
1418 let mut conn = self
1419 .get_connection(self.connections.reader(), "has_entries")
1420 .await?;
1421 let relayer_list_key = self.relayer_list_key();
1422
1423 debug!("checking if transaction entries exist");
1424
1425 let exists: bool = conn
1426 .exists(&relayer_list_key)
1427 .await
1428 .map_err(|e| self.map_redis_error(e, "has_entries_check"))?;
1429
1430 debug!(exists = %exists, "transaction entries exist");
1431 Ok(exists)
1432 }
1433
1434 async fn drop_all_entries(&self) -> Result<(), RepositoryError> {
1435 let mut conn = self
1436 .get_connection(self.connections.primary(), "drop_all_entries")
1437 .await?;
1438 let relayer_list_key = self.relayer_list_key();
1439
1440 debug!("dropping all transaction entries");
1441
1442 let relayer_ids: Vec<String> = conn
1444 .smembers(&relayer_list_key)
1445 .await
1446 .map_err(|e| self.map_redis_error(e, "drop_all_entries_get_relayer_ids"))?;
1447
1448 if relayer_ids.is_empty() {
1449 debug!("no transaction entries to drop");
1450 return Ok(());
1451 }
1452
1453 let mut pipe = redis::pipe();
1455 pipe.atomic();
1456
1457 for relayer_id in &relayer_ids {
1459 let pattern = format!(
1461 "{}:{}:{}:{}:*",
1462 self.key_prefix, RELAYER_PREFIX, relayer_id, TX_PREFIX
1463 );
1464 let mut cursor = 0;
1465 let mut tx_ids = Vec::new();
1466
1467 loop {
1468 let (next_cursor, keys): (u64, Vec<String>) = redis::cmd("SCAN")
1469 .cursor_arg(cursor)
1470 .arg("MATCH")
1471 .arg(&pattern)
1472 .query_async(&mut conn)
1473 .await
1474 .map_err(|e| self.map_redis_error(e, "drop_all_entries_scan"))?;
1475
1476 for key in keys {
1478 pipe.del(&key);
1479 if let Some(tx_id) = key.split(':').next_back() {
1480 tx_ids.push(tx_id.to_string());
1481 }
1482 }
1483
1484 cursor = next_cursor;
1485 if cursor == 0 {
1486 break;
1487 }
1488 }
1489
1490 for tx_id in tx_ids {
1492 let reverse_key = self.tx_to_relayer_key(&tx_id);
1493 pipe.del(&reverse_key);
1494
1495 for status in &[
1498 TransactionStatus::Canceled,
1499 TransactionStatus::Pending,
1500 TransactionStatus::Sent,
1501 TransactionStatus::Submitted,
1502 TransactionStatus::Mined,
1503 TransactionStatus::Confirmed,
1504 TransactionStatus::Failed,
1505 TransactionStatus::Expired,
1506 ] {
1507 let status_sorted_key = self.relayer_status_sorted_key(relayer_id, status);
1509 pipe.zrem(&status_sorted_key, &tx_id);
1510
1511 let status_key = self.relayer_status_key(relayer_id, status);
1513 pipe.srem(&status_key, &tx_id);
1514 }
1515 }
1516
1517 let relayer_sorted_key = self.relayer_tx_by_created_at_key(relayer_id);
1519 pipe.del(&relayer_sorted_key);
1520 }
1521
1522 pipe.del(&relayer_list_key);
1524
1525 pipe.exec_async(&mut conn)
1526 .await
1527 .map_err(|e| self.map_redis_error(e, "drop_all_entries_pipeline"))?;
1528
1529 debug!(count = %relayer_ids.len(), "dropped all transaction entries for relayers");
1530 Ok(())
1531 }
1532}
1533
1534#[async_trait]
1535impl TransactionRepository for RedisTransactionRepository {
1536 async fn find_by_relayer_id(
1537 &self,
1538 relayer_id: &str,
1539 query: PaginationQuery,
1540 ) -> Result<PaginatedResult<TransactionRepoModel>, RepositoryError> {
1541 let mut conn = self
1542 .get_connection(self.connections.reader(), "find_by_relayer_id")
1543 .await?;
1544
1545 debug!(relayer_id = %relayer_id, page = %query.page, per_page = %query.per_page, "fetching transactions for relayer sorted by created_at (newest first)");
1546
1547 let relayer_sorted_key = self.relayer_tx_by_created_at_key(relayer_id);
1548
1549 let sorted_set_count: u64 = conn
1551 .zcard(&relayer_sorted_key)
1552 .await
1553 .map_err(|e| self.map_redis_error(e, "find_by_relayer_id_count"))?;
1554
1555 if sorted_set_count == 0 {
1558 debug!(relayer_id = %relayer_id, "no transactions found for relayer (sorted set is empty)");
1559 return Ok(PaginatedResult {
1560 items: vec![],
1561 total: 0,
1562 page: query.page,
1563 per_page: query.per_page,
1564 });
1565 }
1566
1567 let total = sorted_set_count;
1568
1569 let start = ((query.page - 1) * query.per_page) as isize;
1571 let end = start + query.per_page as isize - 1;
1572
1573 if start as u64 >= total {
1574 debug!(relayer_id = %relayer_id, page = %query.page, total = %total, "page is beyond available data");
1575 return Ok(PaginatedResult {
1576 items: vec![],
1577 total,
1578 page: query.page,
1579 per_page: query.per_page,
1580 });
1581 }
1582
1583 let page_ids: Vec<String> = redis::cmd("ZRANGE")
1585 .arg(&relayer_sorted_key)
1586 .arg(start)
1587 .arg(end)
1588 .arg("REV")
1589 .query_async(&mut conn)
1590 .await
1591 .map_err(|e| self.map_redis_error(e, "find_by_relayer_id_sorted"))?;
1592
1593 drop(conn);
1595
1596 let items = self.get_transactions_by_ids(&page_ids).await?;
1597
1598 debug!(relayer_id = %relayer_id, count = %items.results.len(), page = %query.page, "successfully fetched transactions for relayer");
1599
1600 Ok(PaginatedResult {
1601 items: items.results,
1602 total,
1603 page: query.page,
1604 per_page: query.per_page,
1605 })
1606 }
1607
1608 async fn find_by_status(
1610 &self,
1611 relayer_id: &str,
1612 statuses: &[TransactionStatus],
1613 ) -> Result<Vec<TransactionRepoModel>, RepositoryError> {
1614 for status in statuses {
1616 self.ensure_status_sorted_set(relayer_id, status).await?;
1617 }
1618
1619 let mut conn = self
1621 .get_connection(self.connections.reader(), "find_by_status")
1622 .await?;
1623
1624 let mut all_ids: Vec<String> = Vec::new();
1625 for status in statuses {
1626 let sorted_key = self.relayer_status_sorted_key(relayer_id, status);
1628 let ids: Vec<String> = redis::cmd("ZRANGE")
1629 .arg(&sorted_key)
1630 .arg(0)
1631 .arg(-1)
1632 .arg("REV") .query_async(&mut conn)
1634 .await
1635 .map_err(|e| self.map_redis_error(e, "find_by_status"))?;
1636
1637 all_ids.extend(ids);
1638 }
1639
1640 drop(conn);
1642
1643 if all_ids.is_empty() {
1644 return Ok(vec![]);
1645 }
1646
1647 all_ids.sort();
1649 all_ids.dedup();
1650
1651 let mut transactions = self.get_transactions_by_ids(&all_ids).await?;
1653
1654 transactions
1656 .results
1657 .sort_by(|a, b| b.created_at.cmp(&a.created_at));
1658
1659 Ok(transactions.results)
1660 }
1661
1662 async fn find_by_status_paginated(
1663 &self,
1664 relayer_id: &str,
1665 statuses: &[TransactionStatus],
1666 query: PaginationQuery,
1667 oldest_first: bool,
1668 ) -> Result<PaginatedResult<TransactionRepoModel>, RepositoryError> {
1669 for status in statuses {
1671 self.ensure_status_sorted_set(relayer_id, status).await?;
1672 }
1673
1674 let mut conn = self
1675 .get_connection(self.connections.reader(), "find_by_status_paginated")
1676 .await?;
1677
1678 if statuses.len() == 1 {
1680 let sorted_key = self.relayer_status_sorted_key(relayer_id, &statuses[0]);
1681
1682 let total: u64 = conn
1684 .zcard(&sorted_key)
1685 .await
1686 .map_err(|e| self.map_redis_error(e, "find_by_status_paginated_count"))?;
1687
1688 if total == 0 {
1689 return Ok(PaginatedResult {
1690 items: vec![],
1691 total: 0,
1692 page: query.page,
1693 per_page: query.per_page,
1694 });
1695 }
1696
1697 let start = ((query.page.saturating_sub(1)) * query.per_page) as isize;
1699 let end = start + query.per_page as isize - 1;
1700
1701 let mut cmd = redis::cmd("ZRANGE");
1704 cmd.arg(&sorted_key).arg(start).arg(end);
1705 if !oldest_first {
1706 cmd.arg("REV");
1707 }
1708 let page_ids: Vec<String> = cmd
1709 .query_async(&mut conn)
1710 .await
1711 .map_err(|e| self.map_redis_error(e, "find_by_status_paginated"))?;
1712
1713 drop(conn);
1715
1716 let transactions = self.get_transactions_by_ids(&page_ids).await?;
1717
1718 debug!(
1719 relayer_id = %relayer_id,
1720 status = %statuses[0],
1721 total = %total,
1722 page = %query.page,
1723 page_size = %transactions.results.len(),
1724 "fetched paginated transactions by single status"
1725 );
1726
1727 return Ok(PaginatedResult {
1728 items: transactions.results,
1729 total,
1730 page: query.page,
1731 per_page: query.per_page,
1732 });
1733 }
1734
1735 let mut all_ids: Vec<(String, f64)> = Vec::new();
1737 for status in statuses {
1738 let sorted_key = self.relayer_status_sorted_key(relayer_id, status);
1739
1740 let ids_with_scores: Vec<(String, f64)> = redis::cmd("ZRANGE")
1742 .arg(&sorted_key)
1743 .arg(0)
1744 .arg(-1)
1745 .arg("WITHSCORES")
1746 .query_async(&mut conn)
1747 .await
1748 .map_err(|e| self.map_redis_error(e, "find_by_status_paginated_multi"))?;
1749
1750 all_ids.extend(ids_with_scores);
1751 }
1752
1753 drop(conn);
1755
1756 let mut id_map: std::collections::HashMap<String, f64> = std::collections::HashMap::new();
1758 for (id, score) in all_ids {
1759 id_map
1760 .entry(id)
1761 .and_modify(|s| {
1762 if oldest_first {
1764 if score < *s {
1765 *s = score
1766 }
1767 } else if score > *s {
1768 *s = score
1769 }
1770 })
1771 .or_insert(score);
1772 }
1773
1774 let mut sorted_ids: Vec<(String, f64)> = id_map.into_iter().collect();
1776 if oldest_first {
1777 sorted_ids.sort_by(|a, b| a.1.partial_cmp(&b.1).unwrap_or(std::cmp::Ordering::Equal));
1778 } else {
1779 sorted_ids.sort_by(|a, b| b.1.partial_cmp(&a.1).unwrap_or(std::cmp::Ordering::Equal));
1780 }
1781
1782 let total = sorted_ids.len() as u64;
1783
1784 if total == 0 {
1785 return Ok(PaginatedResult {
1786 items: vec![],
1787 total: 0,
1788 page: query.page,
1789 per_page: query.per_page,
1790 });
1791 }
1792
1793 let start = ((query.page.saturating_sub(1)) * query.per_page) as usize;
1795 let page_ids: Vec<String> = sorted_ids
1796 .into_iter()
1797 .skip(start)
1798 .take(query.per_page as usize)
1799 .map(|(id, _)| id)
1800 .collect();
1801
1802 let transactions = self.get_transactions_by_ids(&page_ids).await?;
1804
1805 debug!(
1806 relayer_id = %relayer_id,
1807 total = %total,
1808 page = %query.page,
1809 page_size = %transactions.results.len(),
1810 "fetched paginated transactions by status"
1811 );
1812
1813 Ok(PaginatedResult {
1814 items: transactions.results,
1815 total,
1816 page: query.page,
1817 per_page: query.per_page,
1818 })
1819 }
1820
1821 async fn find_by_status_paginated_filtered(
1822 &self,
1823 relayer_id: &str,
1824 statuses: &[TransactionStatus],
1825 query: PaginationQuery,
1826 oldest_first: bool,
1827 exclude_canceled: bool,
1828 ) -> Result<PaginatedResult<TransactionRepoModel>, RepositoryError> {
1829 if !exclude_canceled {
1830 return self
1831 .find_by_status_paginated(relayer_id, statuses, query, oldest_first)
1832 .await;
1833 }
1834
1835 for status in statuses {
1841 self.ensure_status_sorted_set(relayer_id, status).await?;
1842 }
1843
1844 let mut conn = self
1845 .get_connection(
1846 self.connections.reader(),
1847 "find_by_status_paginated_filtered",
1848 )
1849 .await?;
1850
1851 let mut all_ids: Vec<(String, f64)> = Vec::new();
1853 for status in statuses {
1854 let sorted_key = self.relayer_status_sorted_key(relayer_id, status);
1855 let ids_with_scores: Vec<(String, f64)> = redis::cmd("ZRANGE")
1856 .arg(&sorted_key)
1857 .arg(0)
1858 .arg(-1)
1859 .arg("WITHSCORES")
1860 .query_async(&mut conn)
1861 .await
1862 .map_err(|e| self.map_redis_error(e, "find_by_status_paginated_filtered"))?;
1863 all_ids.extend(ids_with_scores);
1864 }
1865
1866 drop(conn);
1867
1868 let mut id_map: std::collections::HashMap<String, f64> = std::collections::HashMap::new();
1870 for (id, score) in all_ids {
1871 id_map
1872 .entry(id)
1873 .and_modify(|s| {
1874 if oldest_first {
1875 if score < *s {
1876 *s = score
1877 }
1878 } else if score > *s {
1879 *s = score
1880 }
1881 })
1882 .or_insert(score);
1883 }
1884
1885 let mut sorted_ids: Vec<(String, f64)> = id_map.into_iter().collect();
1886 if oldest_first {
1887 sorted_ids.sort_by(|a, b| a.1.partial_cmp(&b.1).unwrap_or(std::cmp::Ordering::Equal));
1888 } else {
1889 sorted_ids.sort_by(|a, b| b.1.partial_cmp(&a.1).unwrap_or(std::cmp::Ordering::Equal));
1890 }
1891 let ordered_ids: Vec<String> = sorted_ids.into_iter().map(|(id, _)| id).collect();
1892
1893 let fetched = self.get_transactions_by_ids(&ordered_ids).await?;
1895 let mut by_id: std::collections::HashMap<String, TransactionRepoModel> = fetched
1896 .results
1897 .into_iter()
1898 .map(|t| (t.id.clone(), t))
1899 .collect();
1900
1901 let filtered: Vec<TransactionRepoModel> = ordered_ids
1902 .iter()
1903 .filter_map(|id| by_id.remove(id))
1904 .filter(|t| t.is_canceled != Some(true))
1905 .collect();
1906
1907 let total = filtered.len() as u64;
1908 let start = ((query.page.saturating_sub(1)) * query.per_page) as usize;
1909 let items: Vec<TransactionRepoModel> = filtered
1910 .into_iter()
1911 .skip(start)
1912 .take(query.per_page as usize)
1913 .collect();
1914
1915 Ok(PaginatedResult {
1916 items,
1917 total,
1918 page: query.page,
1919 per_page: query.per_page,
1920 })
1921 }
1922
1923 async fn find_by_nonce(
1924 &self,
1925 relayer_id: &str,
1926 nonce: u64,
1927 ) -> Result<Option<TransactionRepoModel>, RepositoryError> {
1928 let mut conn = self
1929 .get_connection(self.connections.reader(), "find_by_nonce")
1930 .await?;
1931 let nonce_key = self.relayer_nonce_key(relayer_id, nonce);
1932
1933 let tx_id: Option<String> = conn
1935 .get(nonce_key)
1936 .await
1937 .map_err(|e| self.map_redis_error(e, "find_by_nonce"))?;
1938
1939 match tx_id {
1940 Some(tx_id) => {
1941 match self.get_by_id(tx_id.clone()).await {
1942 Ok(tx) => Ok(Some(tx)),
1943 Err(RepositoryError::NotFound(_)) => {
1944 warn!(relayer_id = %relayer_id, nonce = %nonce, "stale nonce index found for relayer");
1946 Ok(None)
1947 }
1948 Err(e) => Err(e),
1949 }
1950 }
1951 None => Ok(None),
1952 }
1953 }
1954
1955 async fn get_nonce_occupancy(
1956 &self,
1957 relayer_id: &str,
1958 from_nonce: u64,
1959 to_nonce: u64,
1960 ) -> Result<Vec<(u64, Option<TransactionStatus>)>, RepositoryError> {
1961 if from_nonce >= to_nonce {
1962 return Ok(vec![]);
1963 }
1964
1965 let nonces: Vec<u64> = (from_nonce..to_nonce).collect();
1966 let nonce_keys: Vec<String> = nonces
1967 .iter()
1968 .map(|n| self.relayer_nonce_key(relayer_id, *n))
1969 .collect();
1970
1971 let mut conn = self
1974 .get_connection(self.connections.primary(), "get_nonce_occupancy")
1975 .await?;
1976 let tx_ids: Vec<Option<String>> = redis::cmd("MGET")
1977 .arg(&nonce_keys)
1978 .query_async(&mut conn)
1979 .await
1980 .map_err(|e| self.map_redis_error(e, "get_nonce_occupancy:mget_nonces"))?;
1981
1982 let mut tx_key_entries: Vec<(usize, String)> = Vec::new();
1985 for (i, tx_id) in tx_ids.iter().enumerate() {
1986 if let Some(id) = tx_id {
1987 tx_key_entries.push((i, self.tx_key(relayer_id, id)));
1988 }
1989 }
1990
1991 let tx_statuses: Vec<Option<TransactionStatus>> = if tx_key_entries.is_empty() {
1993 vec![]
1994 } else {
1995 let data_keys: Vec<&str> = tx_key_entries.iter().map(|(_, k)| k.as_str()).collect();
1996 let raw_values: Vec<Option<String>> = redis::cmd("MGET")
1997 .arg(&data_keys)
1998 .query_async(&mut conn)
1999 .await
2000 .map_err(|e| self.map_redis_error(e, "get_nonce_occupancy:mget_txs"))?;
2001
2002 raw_values
2003 .into_iter()
2004 .enumerate()
2005 .map(|(i, v)| {
2006 v.and_then(|json| {
2007 match serde_json::from_str::<TransactionRepoModel>(&json) {
2008 Ok(tx) => Some(tx.status),
2009 Err(e) => {
2010 let nonce = tx_key_entries.get(i).map(|(idx, _)| nonces[*idx]);
2011 warn!(
2012 relayer_id = %relayer_id,
2013 nonce = ?nonce,
2014 error = %e,
2015 "get_nonce_occupancy: failed to deserialize transaction, treating as empty"
2016 );
2017 None
2018 }
2019 }
2020 })
2021 })
2022 .collect()
2023 };
2024
2025 let mut results: Vec<(u64, Option<TransactionStatus>)> =
2027 nonces.iter().map(|n| (*n, None)).collect();
2028
2029 for (idx, (original_idx, _)) in tx_key_entries.iter().enumerate() {
2030 if let Some(status) = tx_statuses.get(idx).and_then(|s| s.clone()) {
2031 results[*original_idx].1 = Some(status);
2032 }
2033 }
2034
2035 Ok(results)
2036 }
2037
2038 async fn update_status(
2039 &self,
2040 tx_id: String,
2041 status: TransactionStatus,
2042 ) -> Result<TransactionRepoModel, RepositoryError> {
2043 let update = TransactionUpdateRequest {
2044 status: Some(status),
2045 ..Default::default()
2046 };
2047 self.partial_update(tx_id, update).await
2048 }
2049
2050 async fn partial_update(
2051 &self,
2052 tx_id: String,
2053 update: TransactionUpdateRequest,
2054 ) -> Result<TransactionRepoModel, RepositoryError> {
2055 let patch_json = serde_json::to_string(&update).map_err(|e| {
2057 RepositoryError::InvalidData(format!("Failed to serialize update patch: {e}"))
2058 })?;
2059
2060 let delete_at_value = if let Some(ref status) = update.status {
2063 if FINAL_TRANSACTION_STATUSES.contains(status) {
2064 let expiration_hours = ServerConfig::get_transaction_expiration_hours();
2065 let seconds = (expiration_hours * 3600.0) as i64;
2066 let delete_time = Utc::now() + chrono::Duration::seconds(seconds);
2067 Some(delete_time.to_rfc3339())
2068 } else {
2069 None
2070 }
2071 } else {
2072 None
2073 };
2074 let delete_at_arg = delete_at_value.as_deref().unwrap_or("");
2075 let index_metadata_json = if update.status.is_some() {
2078 self.partial_update_index_metadata(&tx_id, &update)?
2079 } else {
2080 String::new()
2081 };
2082
2083 let (lookup_key, key_prefix, key_suffix) = self.tx_key_parts(&tx_id);
2084
2085 let patch_script = Script::new(
2091 r#"
2092 local relayer_id = redis.call('GET', KEYS[1])
2093 if not relayer_id then return false end
2094
2095 local tx_key = ARGV[1] .. relayer_id .. ARGV[2]
2096 local current = redis.call('GET', tx_key)
2097 if not current then return false end
2098
2099 local tx = cjson.decode(current)
2100 local patch = cjson.decode(ARGV[3])
2101
2102 -- Guard: reject status changes on finalized transactions.
2103 -- A stale worker must not resurrect a tx that another worker
2104 -- already moved to a terminal state.
2105 local final_states = {confirmed=true, failed=true, expired=true, canceled=true}
2106 if final_states[tx["status"]] and patch["status"] then
2107 return {current, current}
2108 end
2109
2110 local old_snapshot = current
2111 local old_status = tx["status"]
2112
2113 -- lua-cjson cannot distinguish empty Lua tables from empty
2114 -- arrays, so a decode/encode round-trip turns [] into {}.
2115 -- Record which keys held [] in the stored doc and the patch
2116 -- so we can restore them after cjson.encode.
2117 -- NOTE: this relies on each array-typed field having a unique key
2118 -- name across the entire JSON document (including nested objects).
2119 -- If the model ever introduces duplicate key names at different
2120 -- nesting levels (e.g. metadata.hashes), the gsub below could
2121 -- restore the wrong occurrence.
2122 local empty_arrs = {}
2123 for k in string.gmatch(current, '"([^"]+)"%s*:%s*%[%s*%]') do
2124 empty_arrs[k] = true
2125 end
2126 for k in string.gmatch(ARGV[3], '"([^"]+)"%s*:%s*%[%s*%]') do
2127 empty_arrs[k] = true
2128 end
2129
2130 for k, v in pairs(patch) do
2131 tx[k] = v
2132 end
2133
2134 -- Apply delete_at if transitioning to a final state and not already set
2135 if ARGV[4] ~= '' and (not tx["delete_at"] or tx["delete_at"] == cjson.null) then
2136 tx["delete_at"] = ARGV[4]
2137 end
2138
2139 local updated = cjson.encode(tx)
2140 local new_status = tx["status"]
2141
2142 -- Restore empty arrays that cjson.encode converted to {}
2143 for k, _ in pairs(empty_arrs) do
2144 updated = string.gsub(
2145 updated, '"'..k..'"%s*:%s*{}', '"'..k..'":[]', 1
2146 )
2147 end
2148
2149 local index_meta = nil
2150 if patch["status"] and new_status ~= old_status then
2151 index_meta = cjson.decode(ARGV[5])
2152 end
2153
2154 redis.call('SET', tx_key, updated)
2155
2156 if index_meta then
2157 local tx_id = index_meta["tx_id"]
2158 local new_status_sorted_suffix = index_meta["sorted"][new_status]
2159 local old_status_sorted_suffix = index_meta["sorted"][old_status]
2160 local old_status_legacy_suffix = index_meta["legacy"][old_status]
2161
2162 if new_status_sorted_suffix then
2163 local score = nil
2164 if new_status == "confirmed" and index_meta["confirmed_score"] ~= "" then
2165 score = index_meta["confirmed_score"]
2166 else
2167 local created_key = ARGV[1] .. relayer_id .. index_meta["created_key_suffix"]
2168 score = redis.call('ZSCORE', created_key, tx_id)
2169 if not score then score = 0 end
2170 end
2171
2172 redis.call('ZADD', ARGV[1] .. relayer_id .. new_status_sorted_suffix, score, tx_id)
2173 end
2174
2175 if old_status_sorted_suffix then
2176 redis.call('ZREM', ARGV[1] .. relayer_id .. old_status_sorted_suffix, tx_id)
2177 end
2178 if old_status_legacy_suffix then
2179 redis.call('SREM', ARGV[1] .. relayer_id .. old_status_legacy_suffix, tx_id)
2180 end
2181
2182 if final_states[new_status] then
2183 for _, status in ipairs(index_meta["nonfinal"]) do
2184 local sorted_suffix = index_meta["sorted"][status]
2185 local legacy_suffix = index_meta["legacy"][status]
2186 if sorted_suffix then
2187 redis.call('ZREM', ARGV[1] .. relayer_id .. sorted_suffix, tx_id)
2188 end
2189 if legacy_suffix then
2190 redis.call('SREM', ARGV[1] .. relayer_id .. legacy_suffix, tx_id)
2191 end
2192 end
2193 end
2194 end
2195
2196 return {old_snapshot, updated}
2197 "#,
2198 );
2199
2200 let result: Option<Vec<String>> = self
2201 .run_script_with_retry_vec(
2202 &patch_script,
2203 &lookup_key,
2204 &key_prefix,
2205 &key_suffix,
2206 &[&patch_json, delete_at_arg, &index_metadata_json],
2207 "partial_update",
2208 )
2209 .await?;
2210
2211 let parts = result.ok_or_else(|| {
2212 RepositoryError::NotFound(format!("Transaction with ID {tx_id} not found"))
2213 })?;
2214
2215 if parts.len() != 2 {
2216 return Err(RepositoryError::UnexpectedError(format!(
2217 "partial_update script returned {} elements, expected 2",
2218 parts.len()
2219 )));
2220 }
2221
2222 let old_json = &parts[0];
2223 let new_json = &parts[1];
2224
2225 let original_tx =
2226 self.deserialize_entity::<TransactionRepoModel>(old_json, &tx_id, "transaction")?;
2227 let updated_tx =
2228 self.deserialize_entity::<TransactionRepoModel>(new_json, &tx_id, "transaction")?;
2229
2230 self.update_indexes(&updated_tx, Some(&original_tx), false)
2234 .await?;
2235
2236 debug!(tx_id = %tx_id, "successfully updated transaction via patch");
2237
2238 if original_tx.status != updated_tx.status {
2242 self.track_status_change_metrics(
2243 &original_tx,
2244 &updated_tx,
2245 &original_tx.status,
2246 &updated_tx.status,
2247 );
2248 }
2249
2250 Ok(updated_tx)
2251 }
2252
2253 async fn reconcile_stale_status_indexes(
2254 &self,
2255 relayer_id: &str,
2256 ) -> Result<usize, RepositoryError> {
2257 let cutoff_ms =
2258 (Utc::now() - chrono::Duration::seconds(STALE_INDEX_MIN_AGE_SECS)).timestamp_millis();
2259 let mut total_repaired = 0usize;
2260 let mut per_status_counts = Vec::new();
2261
2262 for status in NON_FINAL_STATUS_INDEXES {
2265 let repaired = self
2266 .reconcile_status_index(relayer_id, status, cutoff_ms, ReconcileMode::Promote)
2267 .await?;
2268 if repaired > 0 {
2269 per_status_counts.push((status.to_string(), repaired));
2270 total_repaired += repaired;
2271 }
2272 }
2273
2274 for status in FINAL_TRANSACTION_STATUSES {
2278 let repaired = self
2279 .reconcile_status_index(relayer_id, status, cutoff_ms, ReconcileMode::PurgeOrphans)
2280 .await?;
2281 if repaired > 0 {
2282 per_status_counts.push((status.to_string(), repaired));
2283 total_repaired += repaired;
2284 }
2285 }
2286
2287 if total_repaired > 0 {
2288 info!(
2289 relayer_id = %relayer_id,
2290 repaired_count = total_repaired,
2291 per_status_counts = ?per_status_counts,
2292 "reconciled stale transaction status indexes"
2293 );
2294 }
2295
2296 Ok(total_repaired)
2297 }
2298
2299 async fn update_network_data(
2300 &self,
2301 tx_id: String,
2302 network_data: NetworkTransactionData,
2303 ) -> Result<TransactionRepoModel, RepositoryError> {
2304 let update = TransactionUpdateRequest {
2305 network_data: Some(network_data),
2306 ..Default::default()
2307 };
2308 self.partial_update(tx_id, update).await
2309 }
2310
2311 async fn set_sent_at(
2312 &self,
2313 tx_id: String,
2314 sent_at: String,
2315 ) -> Result<TransactionRepoModel, RepositoryError> {
2316 let update = TransactionUpdateRequest {
2317 sent_at: Some(sent_at),
2318 ..Default::default()
2319 };
2320 self.partial_update(tx_id, update).await
2321 }
2322
2323 async fn increment_status_check_failures(
2324 &self,
2325 tx_id: String,
2326 ) -> Result<TransactionRepoModel, RepositoryError> {
2327 self.run_atomic_script(
2328 r#"
2329 local function set_obj(json, key, tbl)
2330 local enc = cjson.encode(tbl)
2331 local r, n = string.gsub(json, '"'..key..'"%s*:%s*%b{}', '"'..key..'":'..enc, 1)
2332 if n > 0 then return r end
2333 r, n = string.gsub(json, '"'..key..'"%s*:%s*null', '"'..key..'":'..enc, 1)
2334 if n > 0 then return r end
2335 return string.gsub(json, '}%s*$', ',"'..key..'":'..enc..'}', 1)
2336 end
2337
2338 local relayer_id = redis.call('GET', KEYS[1])
2339 if not relayer_id then return false end
2340
2341 local tx_key = ARGV[1] .. relayer_id .. ARGV[2]
2342 local current = redis.call('GET', tx_key)
2343 if not current then return false end
2344
2345 local tx = cjson.decode(current)
2346 local final_states = {confirmed=true, failed=true, expired=true, canceled=true}
2347 if final_states[tx["status"]] then return current end
2348
2349 local metadata = tx["metadata"]
2350 if type(metadata) ~= 'table' then metadata = {} end
2351 metadata["consecutive_failures"] = (metadata["consecutive_failures"] or 0) + 1
2352 metadata["total_failures"] = (metadata["total_failures"] or 0) + 1
2353
2354 local updated = set_obj(current, "metadata", metadata)
2355 redis.call('SET', tx_key, updated)
2356 return updated
2357 "#,
2358 &tx_id,
2359 &[],
2360 "increment_status_check_failures",
2361 )
2362 .await
2363 }
2364
2365 async fn reset_status_check_consecutive_failures(
2366 &self,
2367 tx_id: String,
2368 ) -> Result<TransactionRepoModel, RepositoryError> {
2369 self.run_atomic_script(
2370 r#"
2371 local function set_obj(json, key, tbl)
2372 local enc = cjson.encode(tbl)
2373 local r, n = string.gsub(json, '"'..key..'"%s*:%s*%b{}', '"'..key..'":'..enc, 1)
2374 if n > 0 then return r end
2375 r, n = string.gsub(json, '"'..key..'"%s*:%s*null', '"'..key..'":'..enc, 1)
2376 if n > 0 then return r end
2377 return string.gsub(json, '}%s*$', ',"'..key..'":'..enc..'}', 1)
2378 end
2379
2380 local relayer_id = redis.call('GET', KEYS[1])
2381 if not relayer_id then return false end
2382
2383 local tx_key = ARGV[1] .. relayer_id .. ARGV[2]
2384 local current = redis.call('GET', tx_key)
2385 if not current then return false end
2386
2387 local tx = cjson.decode(current)
2388 local final_states = {confirmed=true, failed=true, expired=true, canceled=true}
2389 if final_states[tx["status"]] then return current end
2390
2391 local metadata = tx["metadata"]
2392 if type(metadata) ~= 'table' then metadata = {} end
2393 metadata["consecutive_failures"] = 0
2394
2395 local updated = set_obj(current, "metadata", metadata)
2396 redis.call('SET', tx_key, updated)
2397 return updated
2398 "#,
2399 &tx_id,
2400 &[],
2401 "reset_status_check_consecutive_failures",
2402 )
2403 .await
2404 }
2405
2406 async fn record_stellar_insufficient_fee_retry(
2407 &self,
2408 tx_id: String,
2409 sent_at: String,
2410 ) -> Result<TransactionRepoModel, RepositoryError> {
2411 self.run_atomic_script(
2412 r#"
2413 local function set_str(json, key, val)
2414 local enc = cjson.encode(val)
2415 local r, n = string.gsub(json, '"'..key..'"%s*:%s*"[^"]*"', '"'..key..'":'..enc, 1)
2416 if n > 0 then return r end
2417 r, n = string.gsub(json, '"'..key..'"%s*:%s*null', '"'..key..'":'..enc, 1)
2418 if n > 0 then return r end
2419 return string.gsub(json, '}%s*$', ',"'..key..'":'..enc..'}', 1)
2420 end
2421 local function set_obj(json, key, tbl)
2422 local enc = cjson.encode(tbl)
2423 local r, n = string.gsub(json, '"'..key..'"%s*:%s*%b{}', '"'..key..'":'..enc, 1)
2424 if n > 0 then return r end
2425 r, n = string.gsub(json, '"'..key..'"%s*:%s*null', '"'..key..'":'..enc, 1)
2426 if n > 0 then return r end
2427 return string.gsub(json, '}%s*$', ',"'..key..'":'..enc..'}', 1)
2428 end
2429
2430 local relayer_id = redis.call('GET', KEYS[1])
2431 if not relayer_id then return false end
2432
2433 local tx_key = ARGV[1] .. relayer_id .. ARGV[2]
2434 local current = redis.call('GET', tx_key)
2435 if not current then return false end
2436
2437 local tx = cjson.decode(current)
2438 local final_states = {confirmed=true, failed=true, expired=true, canceled=true}
2439 if final_states[tx["status"]] then return current end
2440
2441 local metadata = tx["metadata"]
2442 if type(metadata) ~= 'table' then metadata = {} end
2443 metadata["insufficient_fee_retries"] = (metadata["insufficient_fee_retries"] or 0) + 1
2444
2445 local updated = set_str(current, "sent_at", ARGV[3])
2446 updated = set_obj(updated, "metadata", metadata)
2447 redis.call('SET', tx_key, updated)
2448 return updated
2449 "#,
2450 &tx_id,
2451 &[&sent_at],
2452 "record_stellar_insufficient_fee_retry",
2453 )
2454 .await
2455 }
2456
2457 async fn record_stellar_try_again_later_retry(
2458 &self,
2459 tx_id: String,
2460 sent_at: String,
2461 ) -> Result<TransactionRepoModel, RepositoryError> {
2462 self.run_atomic_script(
2463 r#"
2464 local function set_str(json, key, val)
2465 local enc = cjson.encode(val)
2466 local r, n = string.gsub(json, '"'..key..'"%s*:%s*"[^"]*"', '"'..key..'":'..enc, 1)
2467 if n > 0 then return r end
2468 r, n = string.gsub(json, '"'..key..'"%s*:%s*null', '"'..key..'":'..enc, 1)
2469 if n > 0 then return r end
2470 return string.gsub(json, '}%s*$', ',"'..key..'":'..enc..'}', 1)
2471 end
2472 local function set_obj(json, key, tbl)
2473 local enc = cjson.encode(tbl)
2474 local r, n = string.gsub(json, '"'..key..'"%s*:%s*%b{}', '"'..key..'":'..enc, 1)
2475 if n > 0 then return r end
2476 r, n = string.gsub(json, '"'..key..'"%s*:%s*null', '"'..key..'":'..enc, 1)
2477 if n > 0 then return r end
2478 return string.gsub(json, '}%s*$', ',"'..key..'":'..enc..'}', 1)
2479 end
2480
2481 local relayer_id = redis.call('GET', KEYS[1])
2482 if not relayer_id then return false end
2483
2484 local tx_key = ARGV[1] .. relayer_id .. ARGV[2]
2485 local current = redis.call('GET', tx_key)
2486 if not current then return false end
2487
2488 local tx = cjson.decode(current)
2489 local final_states = {confirmed=true, failed=true, expired=true, canceled=true}
2490 if final_states[tx["status"]] then return current end
2491
2492 local metadata = tx["metadata"]
2493 if type(metadata) ~= 'table' then metadata = {} end
2494 metadata["try_again_later_retries"] = (metadata["try_again_later_retries"] or 0) + 1
2495
2496 local updated = set_str(current, "sent_at", ARGV[3])
2497 updated = set_obj(updated, "metadata", metadata)
2498 redis.call('SET', tx_key, updated)
2499 return updated
2500 "#,
2501 &tx_id,
2502 &[&sent_at],
2503 "record_stellar_try_again_later_retry",
2504 )
2505 .await
2506 }
2507
2508 async fn set_confirmed_at(
2509 &self,
2510 tx_id: String,
2511 confirmed_at: String,
2512 ) -> Result<TransactionRepoModel, RepositoryError> {
2513 let update = TransactionUpdateRequest {
2514 confirmed_at: Some(confirmed_at),
2515 ..Default::default()
2516 };
2517 self.partial_update(tx_id, update).await
2518 }
2519
2520 async fn count_by_status(
2524 &self,
2525 relayer_id: &str,
2526 statuses: &[TransactionStatus],
2527 ) -> Result<u64, RepositoryError> {
2528 let mut conn = self
2529 .get_connection(self.connections.reader(), "count_by_status")
2530 .await?;
2531 let mut total_count: u64 = 0;
2532
2533 for status in statuses {
2534 self.ensure_status_sorted_set(relayer_id, status).await?;
2536
2537 let sorted_key = self.relayer_status_sorted_key(relayer_id, status);
2538 let count: u64 = conn
2539 .zcard(&sorted_key)
2540 .await
2541 .map_err(|e| self.map_redis_error(e, "count_by_status"))?;
2542 total_count += count;
2543 }
2544
2545 debug!(relayer_id = %relayer_id, count = %total_count, "counted transactions by status");
2546 Ok(total_count)
2547 }
2548
2549 async fn delete_by_ids(&self, ids: Vec<String>) -> Result<BatchDeleteResult, RepositoryError> {
2550 if ids.is_empty() {
2551 debug!("no transaction IDs provided for batch delete");
2552 return Ok(BatchDeleteResult::default());
2553 }
2554
2555 debug!(count = %ids.len(), "batch deleting transactions by IDs (with fetch)");
2556
2557 let batch_result = self.get_transactions_by_ids(&ids).await?;
2559
2560 let requests: Vec<TransactionDeleteRequest> = batch_result
2562 .results
2563 .iter()
2564 .map(|tx| TransactionDeleteRequest {
2565 id: tx.id.clone(),
2566 relayer_id: tx.relayer_id.clone(),
2567 nonce: self.extract_nonce(&tx.network_data),
2568 })
2569 .collect();
2570
2571 let mut result = self.delete_by_requests(requests).await?;
2573
2574 for id in batch_result.failed_ids {
2576 result
2577 .failed
2578 .push((id.clone(), format!("Transaction with ID {id} not found")));
2579 }
2580
2581 Ok(result)
2582 }
2583
2584 async fn delete_by_requests(
2585 &self,
2586 requests: Vec<TransactionDeleteRequest>,
2587 ) -> Result<BatchDeleteResult, RepositoryError> {
2588 if requests.is_empty() {
2589 debug!("no delete requests provided for batch delete");
2590 return Ok(BatchDeleteResult::default());
2591 }
2592
2593 debug!(count = %requests.len(), "batch deleting transactions by requests (no fetch)");
2594 let mut conn = self
2595 .get_connection(self.connections.primary(), "batch_delete_no_fetch")
2596 .await?;
2597 let mut pipe = redis::pipe();
2598 pipe.atomic();
2599
2600 let all_statuses = [
2602 TransactionStatus::Canceled,
2603 TransactionStatus::Pending,
2604 TransactionStatus::Sent,
2605 TransactionStatus::Submitted,
2606 TransactionStatus::Mined,
2607 TransactionStatus::Confirmed,
2608 TransactionStatus::Failed,
2609 TransactionStatus::Expired,
2610 ];
2611
2612 for req in &requests {
2614 let tx_key = self.tx_key(&req.relayer_id, &req.id);
2616 pipe.del(&tx_key);
2617
2618 let reverse_key = self.tx_to_relayer_key(&req.id);
2620 pipe.del(&reverse_key);
2621
2622 for status in &all_statuses {
2624 let status_sorted_key = self.relayer_status_sorted_key(&req.relayer_id, status);
2625 pipe.zrem(&status_sorted_key, &req.id);
2626
2627 let status_legacy_key = self.relayer_status_key(&req.relayer_id, status);
2628 pipe.srem(&status_legacy_key, &req.id);
2629 }
2630
2631 if let Some(nonce) = req.nonce {
2633 let nonce_key = self.relayer_nonce_key(&req.relayer_id, nonce);
2634 pipe.del(&nonce_key);
2635 }
2636
2637 let relayer_sorted_key = self.relayer_tx_by_created_at_key(&req.relayer_id);
2639 pipe.zrem(&relayer_sorted_key, &req.id);
2640 }
2641
2642 match pipe.exec_async(&mut conn).await {
2644 Ok(_) => {
2645 let deleted_count = requests.len();
2646 debug!(
2647 deleted_count = %deleted_count,
2648 "batch delete completed"
2649 );
2650 Ok(BatchDeleteResult {
2651 deleted_count,
2652 failed: vec![],
2653 })
2654 }
2655 Err(e) => {
2656 error!(error = %e, "batch delete pipeline failed");
2657 let failed: Vec<(String, String)> = requests
2659 .iter()
2660 .map(|req| (req.id.clone(), format!("Redis pipeline error: {e}")))
2661 .collect();
2662 Ok(BatchDeleteResult {
2663 deleted_count: 0,
2664 failed,
2665 })
2666 }
2667 }
2668 }
2669}
2670
2671#[cfg(test)]
2672mod tests {
2673 use super::*;
2674 use crate::models::{
2675 evm::Speed, EvmTransactionData, MemoSpec, NetworkType, StellarTransactionData,
2676 TransactionInput,
2677 };
2678 use alloy::primitives::U256;
2679 use deadpool_redis::{Config, Runtime};
2680 use lazy_static::lazy_static;
2681 use std::str::FromStr;
2682 use tokio;
2683 use uuid::Uuid;
2684
2685 use tokio::sync::Mutex;
2686
2687 lazy_static! {
2689 static ref ENV_MUTEX: Mutex<()> = Mutex::new(());
2690 }
2691
2692 fn create_test_transaction(id: &str) -> TransactionRepoModel {
2694 TransactionRepoModel {
2695 id: id.to_string(),
2696 relayer_id: "relayer-1".to_string(),
2697 status: TransactionStatus::Pending,
2698 status_reason: None,
2699 created_at: "2025-01-27T15:31:10.777083+00:00".to_string(),
2700 sent_at: Some("2025-01-27T15:31:10.777083+00:00".to_string()),
2701 confirmed_at: Some("2025-01-27T15:31:10.777083+00:00".to_string()),
2702 valid_until: None,
2703 delete_at: None,
2704 network_type: NetworkType::Evm,
2705 priced_at: None,
2706 hashes: vec![],
2707 network_data: NetworkTransactionData::Evm(EvmTransactionData {
2708 gas_price: Some(1000000000),
2709 gas_limit: Some(21000),
2710 nonce: Some(1),
2711 value: U256::from_str("1000000000000000000").unwrap(),
2712 data: Some("0x".to_string()),
2713 from: "0xSender".to_string(),
2714 to: Some("0xRecipient".to_string()),
2715 chain_id: 1,
2716 signature: None,
2717 hash: Some(format!("0x{id}")),
2718 speed: Some(Speed::Fast),
2719 max_fee_per_gas: None,
2720 max_priority_fee_per_gas: None,
2721 raw: None,
2722 }),
2723 noop_count: None,
2724 is_canceled: Some(false),
2725 metadata: None,
2726 }
2727 }
2728
2729 fn create_test_transaction_with_relayer(id: &str, relayer_id: &str) -> TransactionRepoModel {
2730 let mut tx = create_test_transaction(id);
2731 tx.relayer_id = relayer_id.to_string();
2732 tx
2733 }
2734
2735 fn create_stellar_test_transaction(
2736 id: &str,
2737 relayer_id: &str,
2738 sequence_number: i64,
2739 max_fee: i64,
2740 memo_id: u64,
2741 ) -> TransactionRepoModel {
2742 let mut tx = create_test_transaction_with_relayer(id, relayer_id);
2743 tx.network_type = NetworkType::Stellar;
2744 tx.network_data = NetworkTransactionData::Stellar(StellarTransactionData {
2745 source_account: "GAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAWHF".to_string(),
2746 fee: Some(100),
2747 sequence_number: Some(sequence_number),
2748 memo: Some(MemoSpec::Id { value: memo_id }),
2749 valid_until: None,
2750 network_passphrase: "Test SDF Network ; September 2015".to_string(),
2751 signatures: vec![],
2752 hash: None,
2753 simulation_transaction_data: None,
2754 transaction_input: TransactionInput::SignedXdr {
2755 xdr: "AAAA".to_string(),
2756 max_fee,
2757 },
2758 signed_envelope_xdr: None,
2759 transaction_result_xdr: None,
2760 });
2761 tx
2762 }
2763
2764 fn create_test_transaction_with_status(
2765 id: &str,
2766 relayer_id: &str,
2767 status: TransactionStatus,
2768 ) -> TransactionRepoModel {
2769 let mut tx = create_test_transaction_with_relayer(id, relayer_id);
2770 tx.status = status;
2771 tx
2772 }
2773
2774 fn create_test_transaction_with_nonce(
2775 id: &str,
2776 nonce: u64,
2777 relayer_id: &str,
2778 ) -> TransactionRepoModel {
2779 let mut tx = create_test_transaction_with_relayer(id, relayer_id);
2780 if let NetworkTransactionData::Evm(ref mut evm_data) = tx.network_data {
2781 evm_data.nonce = Some(nonce);
2782 }
2783 tx
2784 }
2785
2786 async fn setup_test_repo() -> RedisTransactionRepository {
2787 let redis_url = std::env::var("REDIS_TEST_URL")
2789 .unwrap_or_else(|_| "redis://127.0.0.1:6379".to_string());
2790
2791 let cfg = Config::from_url(&redis_url);
2792 let pool = Arc::new(
2793 cfg.builder()
2794 .expect("Failed to create pool builder")
2795 .max_size(16)
2796 .runtime(Runtime::Tokio1)
2797 .build()
2798 .expect("Failed to build Redis pool"),
2799 );
2800
2801 let connections = Arc::new(RedisConnections::new_single_pool(pool));
2803
2804 let random_id = Uuid::new_v4().to_string();
2805 let key_prefix = format!("test_prefix:{random_id}");
2806
2807 RedisTransactionRepository::new(connections, key_prefix)
2808 .expect("Failed to create RedisTransactionRepository")
2809 }
2810
2811 fn old_status_score() -> f64 {
2812 (Utc::now() - chrono::Duration::hours(2)).timestamp_millis() as f64
2813 }
2814
2815 async fn zadd_status_member(
2816 repo: &RedisTransactionRepository,
2817 relayer_id: &str,
2818 status: &TransactionStatus,
2819 tx_id: &str,
2820 score: f64,
2821 ) {
2822 let key = repo.relayer_status_sorted_key(relayer_id, status);
2823 let mut conn = repo
2824 .get_connection(repo.connections.primary(), "test_zadd_status_member")
2825 .await
2826 .unwrap();
2827 redis::cmd("ZADD")
2828 .arg(&key)
2829 .arg(score)
2830 .arg(tx_id)
2831 .query_async::<()>(&mut conn)
2832 .await
2833 .unwrap();
2834 }
2835
2836 async fn zrem_status_member(
2837 repo: &RedisTransactionRepository,
2838 relayer_id: &str,
2839 status: &TransactionStatus,
2840 tx_id: &str,
2841 ) {
2842 let key = repo.relayer_status_sorted_key(relayer_id, status);
2843 let mut conn = repo
2844 .get_connection(repo.connections.primary(), "test_zrem_status_member")
2845 .await
2846 .unwrap();
2847 redis::cmd("ZREM")
2848 .arg(&key)
2849 .arg(tx_id)
2850 .query_async::<()>(&mut conn)
2851 .await
2852 .unwrap();
2853 }
2854
2855 async fn status_member_exists(
2856 repo: &RedisTransactionRepository,
2857 relayer_id: &str,
2858 status: &TransactionStatus,
2859 tx_id: &str,
2860 ) -> bool {
2861 let key = repo.relayer_status_sorted_key(relayer_id, status);
2862 let mut conn = repo
2863 .get_connection(repo.connections.primary(), "test_status_member_exists")
2864 .await
2865 .unwrap();
2866 let score: Option<f64> = redis::cmd("ZSCORE")
2867 .arg(&key)
2868 .arg(tx_id)
2869 .query_async(&mut conn)
2870 .await
2871 .unwrap();
2872 score.is_some()
2873 }
2874
2875 #[tokio::test]
2876 #[ignore = "Requires active Redis instance"]
2877 async fn test_new_repository_creation() {
2878 let repo = setup_test_repo().await;
2879 assert!(repo.key_prefix.contains("test_prefix"));
2880 }
2881
2882 #[tokio::test]
2883 #[ignore = "Requires active Redis instance"]
2884 async fn test_new_repository_empty_prefix_fails() {
2885 let redis_url = std::env::var("REDIS_TEST_URL")
2886 .unwrap_or_else(|_| "redis://127.0.0.1:6379".to_string());
2887 let cfg = Config::from_url(&redis_url);
2888 let pool = Arc::new(
2889 cfg.builder()
2890 .expect("Failed to create pool builder")
2891 .max_size(16)
2892 .runtime(Runtime::Tokio1)
2893 .build()
2894 .expect("Failed to build Redis pool"),
2895 );
2896 let connections = Arc::new(RedisConnections::new_single_pool(pool));
2897
2898 let result = RedisTransactionRepository::new(connections, "".to_string());
2899 assert!(matches!(result, Err(RepositoryError::InvalidData(_))));
2900 }
2901
2902 #[tokio::test]
2903 #[ignore = "Requires active Redis instance"]
2904 async fn test_key_generation() {
2905 let repo = setup_test_repo().await;
2906
2907 assert!(repo
2908 .tx_key("relayer-1", "test-id")
2909 .contains(":relayer:relayer-1:tx:test-id"));
2910 assert!(repo
2911 .tx_to_relayer_key("test-id")
2912 .contains(":relayer:tx_to_relayer:test-id"));
2913 assert!(repo.relayer_list_key().contains(":relayer_list"));
2914 assert!(repo
2915 .relayer_status_key("relayer-1", &TransactionStatus::Pending)
2916 .contains(":relayer:relayer-1:status:Pending"));
2917 assert!(repo
2918 .relayer_nonce_key("relayer-1", 42)
2919 .contains(":relayer:relayer-1:nonce:42"));
2920 }
2921
2922 #[tokio::test]
2923 #[ignore = "Requires active Redis instance"]
2924 async fn test_serialize_deserialize_transaction() {
2925 let repo = setup_test_repo().await;
2926 let tx = create_test_transaction("test-1");
2927
2928 let serialized = repo
2929 .serialize_entity(&tx, |t| &t.id, "transaction")
2930 .expect("Serialization should succeed");
2931 let deserialized: TransactionRepoModel = repo
2932 .deserialize_entity(&serialized, "test-1", "transaction")
2933 .expect("Deserialization should succeed");
2934
2935 assert_eq!(tx.id, deserialized.id);
2936 assert_eq!(tx.relayer_id, deserialized.relayer_id);
2937 assert_eq!(tx.status, deserialized.status);
2938 }
2939
2940 #[tokio::test]
2941 #[ignore = "Requires active Redis instance"]
2942 async fn test_extract_nonce() {
2943 let repo = setup_test_repo().await;
2944 let random_id = Uuid::new_v4().to_string();
2945 let relayer_id = Uuid::new_v4().to_string();
2946 let tx_with_nonce = create_test_transaction_with_nonce(&random_id, 42, &relayer_id);
2947
2948 let nonce = repo.extract_nonce(&tx_with_nonce.network_data);
2949 assert_eq!(nonce, Some(42));
2950 }
2951
2952 #[tokio::test]
2953 #[ignore = "Requires active Redis instance"]
2954 async fn test_create_transaction() {
2955 let repo = setup_test_repo().await;
2956 let random_id = Uuid::new_v4().to_string();
2957 let tx = create_test_transaction(&random_id);
2958
2959 let result = repo.create(tx.clone()).await.unwrap();
2960 assert_eq!(result.id, tx.id);
2961 }
2962
2963 #[tokio::test]
2964 #[ignore = "Requires active Redis instance"]
2965 async fn test_get_transaction() {
2966 let repo = setup_test_repo().await;
2967 let random_id = Uuid::new_v4().to_string();
2968 let tx = create_test_transaction(&random_id);
2969
2970 repo.create(tx.clone()).await.unwrap();
2971 let stored = repo.get_by_id(random_id.to_string()).await.unwrap();
2972 assert_eq!(stored.id, tx.id);
2973 assert_eq!(stored.relayer_id, tx.relayer_id);
2974 }
2975
2976 #[tokio::test]
2977 #[ignore = "Requires active Redis instance"]
2978 async fn test_update_transaction() {
2979 let repo = setup_test_repo().await;
2980 let random_id = Uuid::new_v4().to_string();
2981 let mut tx = create_test_transaction(&random_id);
2982
2983 repo.create(tx.clone()).await.unwrap();
2984 tx.status = TransactionStatus::Confirmed;
2985
2986 let updated = repo.update(random_id.to_string(), tx).await.unwrap();
2987 assert!(matches!(updated.status, TransactionStatus::Confirmed));
2988 }
2989
2990 #[tokio::test]
2991 #[ignore = "Requires active Redis instance"]
2992 async fn test_delete_transaction() {
2993 let repo = setup_test_repo().await;
2994 let random_id = Uuid::new_v4().to_string();
2995 let tx = create_test_transaction(&random_id);
2996
2997 repo.create(tx).await.unwrap();
2998 repo.delete_by_id(random_id.to_string()).await.unwrap();
2999
3000 let result = repo.get_by_id(random_id.to_string()).await;
3001 assert!(matches!(result, Err(RepositoryError::NotFound(_))));
3002 }
3003
3004 #[tokio::test]
3005 #[ignore = "Requires active Redis instance"]
3006 async fn test_list_all_transactions() {
3007 let repo = setup_test_repo().await;
3008 let random_id = Uuid::new_v4().to_string();
3009 let random_id2 = Uuid::new_v4().to_string();
3010
3011 let tx1 = create_test_transaction(&random_id);
3012 let tx2 = create_test_transaction(&random_id2);
3013
3014 repo.create(tx1).await.unwrap();
3015 repo.create(tx2).await.unwrap();
3016
3017 let transactions = repo.list_all().await.unwrap();
3018 assert!(transactions.len() >= 2);
3019 }
3020
3021 #[tokio::test]
3022 #[ignore = "Requires active Redis instance"]
3023 async fn test_count_transactions() {
3024 let repo = setup_test_repo().await;
3025 let random_id = Uuid::new_v4().to_string();
3026 let tx = create_test_transaction(&random_id);
3027
3028 let count = repo.count().await.unwrap();
3029 repo.create(tx).await.unwrap();
3030 assert!(repo.count().await.unwrap() > count);
3031 }
3032
3033 #[tokio::test]
3034 #[ignore = "Requires active Redis instance"]
3035 async fn test_get_nonexistent_transaction() {
3036 let repo = setup_test_repo().await;
3037 let result = repo.get_by_id("nonexistent".to_string()).await;
3038 assert!(matches!(result, Err(RepositoryError::NotFound(_))));
3039 }
3040
3041 #[tokio::test]
3042 #[ignore = "Requires active Redis instance"]
3043 async fn test_duplicate_transaction_creation() {
3044 let repo = setup_test_repo().await;
3045 let random_id = Uuid::new_v4().to_string();
3046
3047 let tx = create_test_transaction(&random_id);
3048
3049 repo.create(tx.clone()).await.unwrap();
3050 let result = repo.create(tx).await;
3051
3052 assert!(matches!(
3053 result,
3054 Err(RepositoryError::ConstraintViolation(_))
3055 ));
3056 }
3057
3058 #[tokio::test]
3059 #[ignore = "Requires active Redis instance"]
3060 async fn test_update_nonexistent_transaction() {
3061 let repo = setup_test_repo().await;
3062 let tx = create_test_transaction("test-1");
3063
3064 let result = repo.update("nonexistent".to_string(), tx).await;
3065 assert!(matches!(result, Err(RepositoryError::NotFound(_))));
3066 }
3067
3068 #[tokio::test]
3069 #[ignore = "Requires active Redis instance"]
3070 async fn test_list_paginated() {
3071 let repo = setup_test_repo().await;
3072
3073 for _ in 1..=10 {
3075 let random_id = Uuid::new_v4().to_string();
3076 let tx = create_test_transaction(&random_id);
3077 repo.create(tx).await.unwrap();
3078 }
3079
3080 let query = PaginationQuery {
3082 page: 1,
3083 per_page: 3,
3084 };
3085 let result = repo.list_paginated(query).await.unwrap();
3086 assert_eq!(result.items.len(), 3);
3087 assert!(result.total >= 10);
3088 assert_eq!(result.page, 1);
3089 assert_eq!(result.per_page, 3);
3090
3091 let query = PaginationQuery {
3093 page: 1000,
3094 per_page: 3,
3095 };
3096 let result = repo.list_paginated(query).await.unwrap();
3097 assert_eq!(result.items.len(), 0);
3098 }
3099
3100 #[tokio::test]
3101 #[ignore = "Requires active Redis instance"]
3102 async fn test_find_by_relayer_id() {
3103 let repo = setup_test_repo().await;
3104 let random_id = Uuid::new_v4().to_string();
3105 let random_id2 = Uuid::new_v4().to_string();
3106 let random_id3 = Uuid::new_v4().to_string();
3107
3108 let tx1 = create_test_transaction_with_relayer(&random_id, "relayer-1");
3109 let tx2 = create_test_transaction_with_relayer(&random_id2, "relayer-1");
3110 let tx3 = create_test_transaction_with_relayer(&random_id3, "relayer-2");
3111
3112 repo.create(tx1).await.unwrap();
3113 repo.create(tx2).await.unwrap();
3114 repo.create(tx3).await.unwrap();
3115
3116 let query = PaginationQuery {
3118 page: 1,
3119 per_page: 10,
3120 };
3121 let result = repo
3122 .find_by_relayer_id("relayer-1", query.clone())
3123 .await
3124 .unwrap();
3125 assert!(result.total >= 2);
3126 assert!(result.items.len() >= 2);
3127 assert!(result.items.iter().all(|tx| tx.relayer_id == "relayer-1"));
3128
3129 let result = repo
3131 .find_by_relayer_id("relayer-2", query.clone())
3132 .await
3133 .unwrap();
3134 assert!(result.total >= 1);
3135 assert!(!result.items.is_empty());
3136 assert!(result.items.iter().all(|tx| tx.relayer_id == "relayer-2"));
3137
3138 let result = repo
3140 .find_by_relayer_id("non-existent", query.clone())
3141 .await
3142 .unwrap();
3143 assert_eq!(result.total, 0);
3144 assert_eq!(result.items.len(), 0);
3145 }
3146
3147 #[tokio::test]
3148 #[ignore = "Requires active Redis instance"]
3149 async fn test_find_by_relayer_id_sorted_by_created_at_newest_first() {
3150 let repo = setup_test_repo().await;
3151 let relayer_id = Uuid::new_v4().to_string();
3152
3153 let mut tx1 = create_test_transaction_with_relayer("test-1", &relayer_id);
3155 tx1.created_at = "2025-01-27T10:00:00.000000+00:00".to_string(); let mut tx2 = create_test_transaction_with_relayer("test-2", &relayer_id);
3158 tx2.created_at = "2025-01-27T12:00:00.000000+00:00".to_string(); let mut tx3 = create_test_transaction_with_relayer("test-3", &relayer_id);
3161 tx3.created_at = "2025-01-27T14:00:00.000000+00:00".to_string(); repo.create(tx2.clone()).await.unwrap(); repo.create(tx1.clone()).await.unwrap(); repo.create(tx3.clone()).await.unwrap(); let query = PaginationQuery {
3169 page: 1,
3170 per_page: 10,
3171 };
3172 let result = repo.find_by_relayer_id(&relayer_id, query).await.unwrap();
3173
3174 assert_eq!(result.total, 3);
3175 assert_eq!(result.items.len(), 3);
3176
3177 assert_eq!(
3179 result.items[0].id, "test-3",
3180 "First item should be newest (test-3)"
3181 );
3182 assert_eq!(
3183 result.items[0].created_at,
3184 "2025-01-27T14:00:00.000000+00:00"
3185 );
3186
3187 assert_eq!(
3188 result.items[1].id, "test-2",
3189 "Second item should be middle (test-2)"
3190 );
3191 assert_eq!(
3192 result.items[1].created_at,
3193 "2025-01-27T12:00:00.000000+00:00"
3194 );
3195
3196 assert_eq!(
3197 result.items[2].id, "test-1",
3198 "Third item should be oldest (test-1)"
3199 );
3200 assert_eq!(
3201 result.items[2].created_at,
3202 "2025-01-27T10:00:00.000000+00:00"
3203 );
3204 }
3205
3206 #[tokio::test]
3207 #[ignore = "Requires active Redis instance"]
3208 async fn test_find_by_status() {
3209 let repo = setup_test_repo().await;
3210 let random_id = Uuid::new_v4().to_string();
3211 let random_id2 = Uuid::new_v4().to_string();
3212 let random_id3 = Uuid::new_v4().to_string();
3213 let relayer_id = Uuid::new_v4().to_string();
3214 let tx1 = create_test_transaction_with_status(
3215 &random_id,
3216 &relayer_id,
3217 TransactionStatus::Pending,
3218 );
3219 let tx2 =
3220 create_test_transaction_with_status(&random_id2, &relayer_id, TransactionStatus::Sent);
3221 let tx3 = create_test_transaction_with_status(
3222 &random_id3,
3223 &relayer_id,
3224 TransactionStatus::Confirmed,
3225 );
3226
3227 repo.create(tx1).await.unwrap();
3228 repo.create(tx2).await.unwrap();
3229 repo.create(tx3).await.unwrap();
3230
3231 let result = repo
3233 .find_by_status(&relayer_id, &[TransactionStatus::Pending])
3234 .await
3235 .unwrap();
3236 assert_eq!(result.len(), 1);
3237 assert_eq!(result[0].status, TransactionStatus::Pending);
3238
3239 let result = repo
3241 .find_by_status(
3242 &relayer_id,
3243 &[TransactionStatus::Pending, TransactionStatus::Sent],
3244 )
3245 .await
3246 .unwrap();
3247 assert_eq!(result.len(), 2);
3248
3249 let result = repo
3251 .find_by_status(&relayer_id, &[TransactionStatus::Failed])
3252 .await
3253 .unwrap();
3254 assert_eq!(result.len(), 0);
3255 }
3256
3257 #[tokio::test]
3258 #[ignore = "Requires active Redis instance"]
3259 async fn test_find_by_status_paginated() {
3260 let repo = setup_test_repo().await;
3261 let relayer_id = Uuid::new_v4().to_string();
3262
3263 for i in 1..=5 {
3265 let tx_id = Uuid::new_v4().to_string();
3266 let mut tx = create_test_transaction_with_status(
3267 &tx_id,
3268 &relayer_id,
3269 TransactionStatus::Pending,
3270 );
3271 tx.created_at = format!("2025-01-27T{:02}:00:00.000000+00:00", 10 + i);
3272 repo.create(tx).await.unwrap();
3273 }
3274
3275 for i in 6..=7 {
3277 let tx_id = Uuid::new_v4().to_string();
3278 let mut tx = create_test_transaction_with_status(
3279 &tx_id,
3280 &relayer_id,
3281 TransactionStatus::Confirmed,
3282 );
3283 tx.created_at = format!("2025-01-27T{:02}:00:00.000000+00:00", 10 + i);
3284 repo.create(tx).await.unwrap();
3285 }
3286
3287 let query = PaginationQuery {
3289 page: 1,
3290 per_page: 2,
3291 };
3292 let result = repo
3293 .find_by_status_paginated(&relayer_id, &[TransactionStatus::Pending], query, false)
3294 .await
3295 .unwrap();
3296
3297 assert_eq!(result.total, 5);
3298 assert_eq!(result.items.len(), 2);
3299 assert_eq!(result.page, 1);
3300 assert_eq!(result.per_page, 2);
3301
3302 let query = PaginationQuery {
3304 page: 2,
3305 per_page: 2,
3306 };
3307 let result = repo
3308 .find_by_status_paginated(&relayer_id, &[TransactionStatus::Pending], query, false)
3309 .await
3310 .unwrap();
3311
3312 assert_eq!(result.total, 5);
3313 assert_eq!(result.items.len(), 2);
3314 assert_eq!(result.page, 2);
3315
3316 let query = PaginationQuery {
3318 page: 3,
3319 per_page: 2,
3320 };
3321 let result = repo
3322 .find_by_status_paginated(&relayer_id, &[TransactionStatus::Pending], query, false)
3323 .await
3324 .unwrap();
3325
3326 assert_eq!(result.total, 5);
3327 assert_eq!(result.items.len(), 1);
3328
3329 let query = PaginationQuery {
3331 page: 1,
3332 per_page: 10,
3333 };
3334 let result = repo
3335 .find_by_status_paginated(
3336 &relayer_id,
3337 &[TransactionStatus::Pending, TransactionStatus::Confirmed],
3338 query,
3339 false,
3340 )
3341 .await
3342 .unwrap();
3343
3344 assert_eq!(result.total, 7);
3345 assert_eq!(result.items.len(), 7);
3346
3347 let query = PaginationQuery {
3349 page: 1,
3350 per_page: 10,
3351 };
3352 let result = repo
3353 .find_by_status_paginated(&relayer_id, &[TransactionStatus::Failed], query, false)
3354 .await
3355 .unwrap();
3356
3357 assert_eq!(result.total, 0);
3358 assert_eq!(result.items.len(), 0);
3359 }
3360
3361 #[tokio::test]
3362 #[ignore = "Requires active Redis instance"]
3363 async fn test_find_by_status_paginated_oldest_first() {
3364 let repo = setup_test_repo().await;
3365 let relayer_id = Uuid::new_v4().to_string();
3366
3367 for i in 1..=5 {
3369 let tx_id = format!("tx{}-{}", i, Uuid::new_v4());
3370 let mut tx = create_test_transaction(&tx_id);
3371 tx.relayer_id = relayer_id.clone();
3372 tx.status = TransactionStatus::Pending;
3373 tx.created_at = format!("2025-01-27T{:02}:00:00.000000+00:00", 10 + i);
3374 repo.create(tx).await.unwrap();
3375 }
3376
3377 let query = PaginationQuery {
3379 page: 1,
3380 per_page: 3,
3381 };
3382 let result = repo
3383 .find_by_status_paginated(
3384 &relayer_id,
3385 &[TransactionStatus::Pending],
3386 query.clone(),
3387 true,
3388 )
3389 .await
3390 .unwrap();
3391
3392 assert_eq!(result.total, 5);
3393 assert_eq!(result.items.len(), 3);
3394 assert!(
3396 result.items[0].created_at < result.items[1].created_at,
3397 "First item should be older than second"
3398 );
3399 assert!(
3400 result.items[1].created_at < result.items[2].created_at,
3401 "Second item should be older than third"
3402 );
3403
3404 let result_newest = repo
3406 .find_by_status_paginated(&relayer_id, &[TransactionStatus::Pending], query, false)
3407 .await
3408 .unwrap();
3409
3410 assert_eq!(result_newest.items.len(), 3);
3411 assert!(
3413 result_newest.items[0].created_at > result_newest.items[1].created_at,
3414 "First item should be newer than second"
3415 );
3416 assert!(
3417 result_newest.items[1].created_at > result_newest.items[2].created_at,
3418 "Second item should be newer than third"
3419 );
3420 }
3421
3422 #[tokio::test]
3423 #[ignore = "Requires active Redis instance"]
3424 async fn test_find_by_status_paginated_oldest_first_single_item() {
3425 let repo = setup_test_repo().await;
3426 let relayer_id = Uuid::new_v4().to_string();
3427
3428 let timestamps = [
3430 "2025-01-27T08:00:00.000000+00:00", "2025-01-27T10:00:00.000000+00:00", "2025-01-27T12:00:00.000000+00:00", ];
3434
3435 let mut oldest_id = String::new();
3436 let mut newest_id = String::new();
3437
3438 for (i, timestamp) in timestamps.iter().enumerate() {
3439 let tx_id = format!("tx-{}-{}", i, Uuid::new_v4());
3440 if i == 0 {
3441 oldest_id = tx_id.clone();
3442 }
3443 if i == 2 {
3444 newest_id = tx_id.clone();
3445 }
3446 let mut tx = create_test_transaction(&tx_id);
3447 tx.relayer_id = relayer_id.clone();
3448 tx.status = TransactionStatus::Pending;
3449 tx.created_at = timestamp.to_string();
3450 repo.create(tx).await.unwrap();
3451 }
3452
3453 let query = PaginationQuery {
3455 page: 1,
3456 per_page: 1,
3457 };
3458 let result = repo
3459 .find_by_status_paginated(
3460 &relayer_id,
3461 &[TransactionStatus::Pending],
3462 query.clone(),
3463 true,
3464 )
3465 .await
3466 .unwrap();
3467
3468 assert_eq!(result.total, 3);
3469 assert_eq!(result.items.len(), 1);
3470 assert_eq!(
3471 result.items[0].id, oldest_id,
3472 "With oldest_first=true and per_page=1, should return the oldest transaction"
3473 );
3474
3475 let result = repo
3477 .find_by_status_paginated(&relayer_id, &[TransactionStatus::Pending], query, false)
3478 .await
3479 .unwrap();
3480
3481 assert_eq!(result.items.len(), 1);
3482 assert_eq!(
3483 result.items[0].id, newest_id,
3484 "With oldest_first=false and per_page=1, should return the newest transaction"
3485 );
3486 }
3487
3488 #[tokio::test]
3489 #[ignore = "Requires active Redis instance"]
3490 async fn test_find_by_nonce() {
3491 let repo = setup_test_repo().await;
3492 let random_id = Uuid::new_v4().to_string();
3493 let random_id2 = Uuid::new_v4().to_string();
3494 let relayer_id = Uuid::new_v4().to_string();
3495
3496 let tx1 = create_test_transaction_with_nonce(&random_id, 42, &relayer_id);
3497 let tx2 = create_test_transaction_with_nonce(&random_id2, 43, &relayer_id);
3498
3499 repo.create(tx1.clone()).await.unwrap();
3500 repo.create(tx2).await.unwrap();
3501
3502 let result = repo.find_by_nonce(&relayer_id, 42).await.unwrap();
3504 assert!(result.is_some());
3505 assert_eq!(result.unwrap().id, random_id);
3506
3507 let result = repo.find_by_nonce(&relayer_id, 99).await.unwrap();
3509 assert!(result.is_none());
3510
3511 let result = repo.find_by_nonce("non-existent", 42).await.unwrap();
3513 assert!(result.is_none());
3514 }
3515
3516 #[tokio::test]
3517 #[ignore = "Requires active Redis instance"]
3518 async fn test_get_nonce_occupancy_mixed_slots() {
3519 let repo = setup_test_repo().await;
3520 let relayer_id = Uuid::new_v4().to_string();
3521
3522 let tx1 = create_test_transaction_with_nonce(&Uuid::new_v4().to_string(), 10, &relayer_id);
3524 repo.create(tx1).await.unwrap();
3525
3526 let mut tx2 =
3527 create_test_transaction_with_nonce(&Uuid::new_v4().to_string(), 11, &relayer_id);
3528 tx2.status = TransactionStatus::Failed;
3529 repo.create(tx2).await.unwrap();
3530
3531 let result = repo.get_nonce_occupancy(&relayer_id, 10, 13).await.unwrap();
3532
3533 assert_eq!(result.len(), 3);
3534 assert_eq!(result[0], (10, Some(TransactionStatus::Pending)));
3535 assert_eq!(result[1], (11, Some(TransactionStatus::Failed)));
3536 assert_eq!(result[2], (12, None));
3537 }
3538
3539 #[tokio::test]
3540 #[ignore = "Requires active Redis instance"]
3541 async fn test_get_nonce_occupancy_empty_range() {
3542 let repo = setup_test_repo().await;
3543
3544 let result = repo.get_nonce_occupancy("any-relayer", 5, 5).await.unwrap();
3545 assert!(result.is_empty());
3546
3547 let result = repo
3548 .get_nonce_occupancy("any-relayer", 10, 5)
3549 .await
3550 .unwrap();
3551 assert!(result.is_empty());
3552 }
3553
3554 #[tokio::test]
3555 #[ignore = "Requires active Redis instance"]
3556 async fn test_get_nonce_occupancy_all_empty() {
3557 let repo = setup_test_repo().await;
3558 let relayer_id = Uuid::new_v4().to_string();
3559
3560 let result = repo
3562 .get_nonce_occupancy(&relayer_id, 100, 103)
3563 .await
3564 .unwrap();
3565
3566 assert_eq!(result.len(), 3);
3567 assert!(result.iter().all(|(_, status)| status.is_none()));
3568 }
3569
3570 #[tokio::test]
3571 #[ignore = "Requires active Redis instance"]
3572 async fn test_update_status() {
3573 let repo = setup_test_repo().await;
3574 let random_id = Uuid::new_v4().to_string();
3575 let tx = create_test_transaction(&random_id);
3576
3577 repo.create(tx).await.unwrap();
3578 let updated = repo
3579 .update_status(random_id.to_string(), TransactionStatus::Confirmed)
3580 .await
3581 .unwrap();
3582 assert_eq!(updated.status, TransactionStatus::Confirmed);
3583 }
3584
3585 #[tokio::test]
3586 #[ignore = "Requires active Redis instance"]
3587 async fn test_partial_update() {
3588 let repo = setup_test_repo().await;
3589 let random_id = Uuid::new_v4().to_string();
3590 let tx = create_test_transaction(&random_id);
3591
3592 repo.create(tx).await.unwrap();
3593
3594 let update = TransactionUpdateRequest {
3595 status: Some(TransactionStatus::Sent),
3596 status_reason: Some("Transaction sent".to_string()),
3597 sent_at: Some("2025-01-27T16:00:00.000000+00:00".to_string()),
3598 confirmed_at: None,
3599 network_data: None,
3600 hashes: None,
3601 is_canceled: None,
3602 priced_at: None,
3603 noop_count: None,
3604 delete_at: None,
3605 metadata: None,
3606 };
3607
3608 let updated = repo
3609 .partial_update(random_id.to_string(), update)
3610 .await
3611 .unwrap();
3612 assert_eq!(updated.status, TransactionStatus::Sent);
3613 assert_eq!(updated.status_reason, Some("Transaction sent".to_string()));
3614 assert_eq!(
3615 updated.sent_at,
3616 Some("2025-01-27T16:00:00.000000+00:00".to_string())
3617 );
3618 }
3619
3620 #[tokio::test]
3621 #[ignore = "Requires active Redis instance"]
3622 async fn test_partial_update_confirmed_updates_status_indexes() {
3623 let repo = setup_test_repo().await;
3624 let tx_id = Uuid::new_v4().to_string();
3625 let relayer_id = Uuid::new_v4().to_string();
3626 let mut tx =
3627 create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Submitted);
3628 tx.confirmed_at = None;
3629
3630 repo.create(tx).await.unwrap();
3631
3632 let update = TransactionUpdateRequest {
3633 status: Some(TransactionStatus::Confirmed),
3634 confirmed_at: Some("2025-01-27T16:00:00.000000+00:00".to_string()),
3635 ..Default::default()
3636 };
3637
3638 repo.partial_update(tx_id.clone(), update).await.unwrap();
3639
3640 assert!(
3641 !status_member_exists(&repo, &relayer_id, &TransactionStatus::Submitted, &tx_id).await
3642 );
3643 assert!(
3644 status_member_exists(&repo, &relayer_id, &TransactionStatus::Confirmed, &tx_id).await
3645 );
3646 }
3647
3648 #[tokio::test]
3649 #[ignore = "Requires active Redis instance"]
3650 async fn test_partial_update_finalization_purges_stale_nonfinal_indexes() {
3651 let repo = setup_test_repo().await;
3652 let tx_id = Uuid::new_v4().to_string();
3653 let relayer_id = Uuid::new_v4().to_string();
3654 let mut tx =
3655 create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Submitted);
3656 tx.confirmed_at = None;
3657
3658 repo.create(tx).await.unwrap();
3659 zadd_status_member(
3660 &repo,
3661 &relayer_id,
3662 &TransactionStatus::Sent,
3663 &tx_id,
3664 old_status_score(),
3665 )
3666 .await;
3667
3668 let update = TransactionUpdateRequest {
3669 status: Some(TransactionStatus::Confirmed),
3670 confirmed_at: Some("2025-01-27T16:00:00.000000+00:00".to_string()),
3671 ..Default::default()
3672 };
3673
3674 repo.partial_update(tx_id.clone(), update).await.unwrap();
3675
3676 for status in NON_FINAL_STATUS_INDEXES {
3677 assert!(
3678 !status_member_exists(&repo, &relayer_id, status, &tx_id).await,
3679 "transaction should not remain in {status:?}"
3680 );
3681 }
3682 assert!(
3683 status_member_exists(&repo, &relayer_id, &TransactionStatus::Confirmed, &tx_id).await
3684 );
3685 }
3686
3687 #[tokio::test]
3688 #[ignore = "Requires active Redis instance"]
3689 async fn test_reconcile_stale_status_indexes_heals_final_status_ghost() {
3690 let repo = setup_test_repo().await;
3691 let tx_id = Uuid::new_v4().to_string();
3692 let relayer_id = Uuid::new_v4().to_string();
3693 let tx =
3694 create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Submitted);
3695
3696 repo.create(tx).await.unwrap();
3697 let update = TransactionUpdateRequest {
3698 status: Some(TransactionStatus::Confirmed),
3699 confirmed_at: Some("2025-01-27T16:00:00.000000+00:00".to_string()),
3700 ..Default::default()
3701 };
3702 repo.partial_update(tx_id.clone(), update).await.unwrap();
3703
3704 zadd_status_member(
3705 &repo,
3706 &relayer_id,
3707 &TransactionStatus::Submitted,
3708 &tx_id,
3709 old_status_score(),
3710 )
3711 .await;
3712 zrem_status_member(&repo, &relayer_id, &TransactionStatus::Confirmed, &tx_id).await;
3713
3714 let repaired = repo
3715 .reconcile_stale_status_indexes(&relayer_id)
3716 .await
3717 .unwrap();
3718
3719 assert_eq!(repaired, 1);
3720 assert!(
3721 !status_member_exists(&repo, &relayer_id, &TransactionStatus::Submitted, &tx_id).await
3722 );
3723 assert!(
3724 status_member_exists(&repo, &relayer_id, &TransactionStatus::Confirmed, &tx_id).await
3725 );
3726 }
3727
3728 #[tokio::test]
3729 #[ignore = "Requires active Redis instance"]
3730 async fn test_reconcile_stale_status_indexes_removes_dangling_entry() {
3731 let repo = setup_test_repo().await;
3732 let tx_id = Uuid::new_v4().to_string();
3733 let relayer_id = Uuid::new_v4().to_string();
3734
3735 zadd_status_member(
3736 &repo,
3737 &relayer_id,
3738 &TransactionStatus::Sent,
3739 &tx_id,
3740 old_status_score(),
3741 )
3742 .await;
3743
3744 let repaired = repo
3745 .reconcile_stale_status_indexes(&relayer_id)
3746 .await
3747 .unwrap();
3748
3749 assert_eq!(repaired, 1);
3750 assert!(!status_member_exists(&repo, &relayer_id, &TransactionStatus::Sent, &tx_id).await);
3751 }
3752
3753 #[tokio::test]
3754 #[ignore = "Requires active Redis instance"]
3755 async fn test_reconcile_stale_status_indexes_skips_young_entry() {
3756 let repo = setup_test_repo().await;
3757 let tx_id = Uuid::new_v4().to_string();
3758 let relayer_id = Uuid::new_v4().to_string();
3759
3760 zadd_status_member(
3761 &repo,
3762 &relayer_id,
3763 &TransactionStatus::Submitted,
3764 &tx_id,
3765 Utc::now().timestamp_millis() as f64,
3766 )
3767 .await;
3768
3769 let repaired = repo
3770 .reconcile_stale_status_indexes(&relayer_id)
3771 .await
3772 .unwrap();
3773
3774 assert_eq!(repaired, 0);
3775 assert!(
3776 status_member_exists(&repo, &relayer_id, &TransactionStatus::Submitted, &tx_id).await
3777 );
3778 }
3779
3780 #[tokio::test]
3781 #[ignore = "Requires active Redis instance"]
3782 async fn test_reconcile_stale_status_indexes_skips_legit_inflight_entry() {
3783 let repo = setup_test_repo().await;
3784 let tx_id = Uuid::new_v4().to_string();
3785 let relayer_id = Uuid::new_v4().to_string();
3786 let tx =
3787 create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Submitted);
3788
3789 repo.create(tx).await.unwrap();
3790
3791 let repaired = repo
3792 .reconcile_stale_status_indexes(&relayer_id)
3793 .await
3794 .unwrap();
3795
3796 assert_eq!(repaired, 0);
3797 assert!(
3798 status_member_exists(&repo, &relayer_id, &TransactionStatus::Submitted, &tx_id).await
3799 );
3800 }
3801
3802 #[tokio::test]
3803 #[ignore = "Requires active Redis instance"]
3804 async fn test_reconcile_removes_final_index_orphan_with_missing_body() {
3805 let repo = setup_test_repo().await;
3806 let tx_id = Uuid::new_v4().to_string();
3807 let relayer_id = Uuid::new_v4().to_string();
3808
3809 zadd_status_member(
3812 &repo,
3813 &relayer_id,
3814 &TransactionStatus::Confirmed,
3815 &tx_id,
3816 old_status_score(),
3817 )
3818 .await;
3819
3820 let repaired = repo
3821 .reconcile_stale_status_indexes(&relayer_id)
3822 .await
3823 .unwrap();
3824
3825 assert_eq!(repaired, 1);
3826 assert!(
3827 !status_member_exists(&repo, &relayer_id, &TransactionStatus::Confirmed, &tx_id).await
3828 );
3829 }
3830
3831 #[tokio::test]
3832 #[ignore = "Requires active Redis instance"]
3833 async fn test_reconcile_keeps_live_final_index_entry() {
3834 let repo = setup_test_repo().await;
3835 let tx_id = Uuid::new_v4().to_string();
3836 let relayer_id = Uuid::new_v4().to_string();
3837 let tx =
3838 create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Submitted);
3839
3840 repo.create(tx).await.unwrap();
3841 let update = TransactionUpdateRequest {
3844 status: Some(TransactionStatus::Confirmed),
3845 confirmed_at: Some("2025-01-27T16:00:00.000000+00:00".to_string()),
3846 ..Default::default()
3847 };
3848 repo.partial_update(tx_id.clone(), update).await.unwrap();
3849
3850 let repaired = repo
3851 .reconcile_stale_status_indexes(&relayer_id)
3852 .await
3853 .unwrap();
3854
3855 assert_eq!(repaired, 0);
3856 assert!(
3857 status_member_exists(&repo, &relayer_id, &TransactionStatus::Confirmed, &tx_id).await
3858 );
3859 let body = repo.get_by_id(tx_id.clone()).await.unwrap();
3860 assert_eq!(body.status, TransactionStatus::Confirmed);
3861 }
3862
3863 #[tokio::test]
3864 #[ignore = "Requires active Redis instance"]
3865 async fn test_reconcile_respects_cutoff_for_final_orphan() {
3866 let repo = setup_test_repo().await;
3867 let tx_id = Uuid::new_v4().to_string();
3868 let relayer_id = Uuid::new_v4().to_string();
3869
3870 zadd_status_member(
3873 &repo,
3874 &relayer_id,
3875 &TransactionStatus::Confirmed,
3876 &tx_id,
3877 Utc::now().timestamp_millis() as f64,
3878 )
3879 .await;
3880
3881 let repaired = repo
3882 .reconcile_stale_status_indexes(&relayer_id)
3883 .await
3884 .unwrap();
3885
3886 assert_eq!(repaired, 0);
3887 assert!(
3888 status_member_exists(&repo, &relayer_id, &TransactionStatus::Confirmed, &tx_id).await
3889 );
3890 }
3891
3892 #[tokio::test]
3893 #[ignore = "Requires active Redis instance"]
3894 async fn test_reconcile_final_index_never_zadds() {
3895 let repo = setup_test_repo().await;
3896 let tx_id = Uuid::new_v4().to_string();
3897 let relayer_id = Uuid::new_v4().to_string();
3898
3899 zadd_status_member(
3900 &repo,
3901 &relayer_id,
3902 &TransactionStatus::Confirmed,
3903 &tx_id,
3904 old_status_score(),
3905 )
3906 .await;
3907
3908 let repaired = repo
3909 .reconcile_stale_status_indexes(&relayer_id)
3910 .await
3911 .unwrap();
3912
3913 assert_eq!(repaired, 1);
3914 for status in ALL_TRANSACTION_STATUSES {
3917 assert!(
3918 !status_member_exists(&repo, &relayer_id, status, &tx_id).await,
3919 "orphan should not have been re-added to {status:?}"
3920 );
3921 }
3922 }
3923
3924 #[tokio::test]
3925 #[ignore = "Requires active Redis instance"]
3926 async fn test_reconcile_leaves_present_mismatched_final_body() {
3927 let repo = setup_test_repo().await;
3928 let tx_id = Uuid::new_v4().to_string();
3929 let relayer_id = Uuid::new_v4().to_string();
3930 let tx =
3931 create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Submitted);
3932
3933 repo.create(tx).await.unwrap();
3934 let update = TransactionUpdateRequest {
3935 status: Some(TransactionStatus::Failed),
3936 ..Default::default()
3937 };
3938 repo.partial_update(tx_id.clone(), update).await.unwrap();
3939
3940 zadd_status_member(
3944 &repo,
3945 &relayer_id,
3946 &TransactionStatus::Confirmed,
3947 &tx_id,
3948 old_status_score(),
3949 )
3950 .await;
3951
3952 let repaired = repo
3953 .reconcile_stale_status_indexes(&relayer_id)
3954 .await
3955 .unwrap();
3956
3957 assert_eq!(repaired, 0);
3958 assert!(
3959 status_member_exists(&repo, &relayer_id, &TransactionStatus::Confirmed, &tx_id).await
3960 );
3961 }
3962
3963 #[tokio::test]
3964 #[ignore = "Requires active Redis instance"]
3965 async fn test_partial_update_preserves_large_stellar_i64_fields() {
3966 let repo = setup_test_repo().await;
3973 let relayer_id = Uuid::new_v4().to_string();
3974 let tx_id = Uuid::new_v4().to_string();
3975 let seq = 643918676885760_i64;
3978 let max_fee = 549755813888_i64;
3979 let memo_id = 1099511627776_u64;
3980
3981 let tx = create_stellar_test_transaction(&tx_id, &relayer_id, seq, max_fee, memo_id);
3982 repo.create(tx).await.unwrap();
3983
3984 let update = TransactionUpdateRequest {
3985 status: Some(TransactionStatus::Sent),
3986 sent_at: Some("2025-01-27T16:00:00.000000+00:00".to_string()),
3987 ..Default::default()
3988 };
3989
3990 let updated = repo.partial_update(tx_id.clone(), update).await.unwrap();
3993 assert_large_stellar_fields(&updated, seq, max_fee, memo_id);
3994
3995 let reloaded = repo.get_by_id(tx_id).await.unwrap();
3998 assert_large_stellar_fields(&reloaded, seq, max_fee, memo_id);
3999 }
4000
4001 fn assert_large_stellar_fields(
4002 tx: &TransactionRepoModel,
4003 seq: i64,
4004 max_fee: i64,
4005 memo_id: u64,
4006 ) {
4007 let stellar = tx.network_data.get_stellar_transaction_data().unwrap();
4008 assert_eq!(stellar.sequence_number, Some(seq));
4009 assert_eq!(stellar.memo, Some(MemoSpec::Id { value: memo_id }));
4010 match stellar.transaction_input {
4011 TransactionInput::SignedXdr { max_fee: mf, .. } => assert_eq!(mf, max_fee),
4012 other => panic!("expected SignedXdr, got {other:?}"),
4013 }
4014 }
4015
4016 #[tokio::test]
4017 #[ignore = "Requires active Redis instance"]
4018 async fn test_set_sent_at() {
4019 let repo = setup_test_repo().await;
4020 let random_id = Uuid::new_v4().to_string();
4021 let tx = create_test_transaction(&random_id);
4022
4023 repo.create(tx).await.unwrap();
4024 let updated = repo
4025 .set_sent_at(
4026 random_id.to_string(),
4027 "2025-01-27T16:00:00.000000+00:00".to_string(),
4028 )
4029 .await
4030 .unwrap();
4031 assert_eq!(
4032 updated.sent_at,
4033 Some("2025-01-27T16:00:00.000000+00:00".to_string())
4034 );
4035 }
4036
4037 #[tokio::test]
4038 #[ignore = "Requires active Redis instance"]
4039 async fn test_set_confirmed_at() {
4040 let repo = setup_test_repo().await;
4041 let random_id = Uuid::new_v4().to_string();
4042 let tx = create_test_transaction(&random_id);
4043
4044 repo.create(tx).await.unwrap();
4045 let updated = repo
4046 .set_confirmed_at(
4047 random_id.to_string(),
4048 "2025-01-27T16:00:00.000000+00:00".to_string(),
4049 )
4050 .await
4051 .unwrap();
4052 assert_eq!(
4053 updated.confirmed_at,
4054 Some("2025-01-27T16:00:00.000000+00:00".to_string())
4055 );
4056 }
4057
4058 #[tokio::test]
4059 #[ignore = "Requires active Redis instance"]
4060 async fn test_update_network_data() {
4061 let repo = setup_test_repo().await;
4062 let random_id = Uuid::new_v4().to_string();
4063 let tx = create_test_transaction(&random_id);
4064
4065 repo.create(tx).await.unwrap();
4066
4067 let new_network_data = NetworkTransactionData::Evm(EvmTransactionData {
4068 gas_price: Some(2000000000),
4069 gas_limit: Some(42000),
4070 nonce: Some(2),
4071 value: U256::from_str("2000000000000000000").unwrap(),
4072 data: Some("0x1234".to_string()),
4073 from: "0xNewSender".to_string(),
4074 to: Some("0xNewRecipient".to_string()),
4075 chain_id: 1,
4076 signature: None,
4077 hash: Some("0xnewhash".to_string()),
4078 speed: Some(Speed::SafeLow),
4079 max_fee_per_gas: None,
4080 max_priority_fee_per_gas: None,
4081 raw: None,
4082 });
4083
4084 let updated = repo
4085 .update_network_data(random_id.to_string(), new_network_data.clone())
4086 .await
4087 .unwrap();
4088 assert_eq!(
4089 updated
4090 .network_data
4091 .get_evm_transaction_data()
4092 .unwrap()
4093 .hash,
4094 new_network_data.get_evm_transaction_data().unwrap().hash
4095 );
4096 }
4097
4098 #[tokio::test]
4099 #[ignore = "Requires active Redis instance"]
4100 async fn test_debug_implementation() {
4101 let repo = setup_test_repo().await;
4102 let debug_str = format!("{repo:?}");
4103 assert!(debug_str.contains("RedisTransactionRepository"));
4104 assert!(debug_str.contains("test_prefix"));
4105 }
4106
4107 #[tokio::test]
4108 #[ignore = "Requires active Redis instance"]
4109 async fn test_error_handling_empty_id() {
4110 let repo = setup_test_repo().await;
4111
4112 let result = repo.get_by_id("".to_string()).await;
4113 assert!(matches!(result, Err(RepositoryError::InvalidData(_))));
4114
4115 let result = repo
4116 .update("".to_string(), create_test_transaction("test"))
4117 .await;
4118 assert!(matches!(result, Err(RepositoryError::InvalidData(_))));
4119
4120 let result = repo.delete_by_id("".to_string()).await;
4121 assert!(matches!(result, Err(RepositoryError::InvalidData(_))));
4122 }
4123
4124 #[tokio::test]
4125 #[ignore = "Requires active Redis instance"]
4126 async fn test_pagination_validation() {
4127 let repo = setup_test_repo().await;
4128
4129 let query = PaginationQuery {
4130 page: 1,
4131 per_page: 0,
4132 };
4133 let result = repo.list_paginated(query).await;
4134 assert!(matches!(result, Err(RepositoryError::InvalidData(_))));
4135 }
4136
4137 #[tokio::test]
4138 #[ignore = "Requires active Redis instance"]
4139 async fn test_index_consistency() {
4140 let repo = setup_test_repo().await;
4141 let random_id = Uuid::new_v4().to_string();
4142 let relayer_id = Uuid::new_v4().to_string();
4143 let tx = create_test_transaction_with_nonce(&random_id, 42, &relayer_id);
4144
4145 repo.create(tx.clone()).await.unwrap();
4147
4148 let found = repo.find_by_nonce(&relayer_id, 42).await.unwrap();
4150 assert!(found.is_some());
4151
4152 let mut updated_tx = tx.clone();
4154 if let NetworkTransactionData::Evm(ref mut evm_data) = updated_tx.network_data {
4155 evm_data.nonce = Some(43);
4156 }
4157
4158 repo.update(random_id.to_string(), updated_tx)
4159 .await
4160 .unwrap();
4161
4162 let old_nonce_result = repo.find_by_nonce(&relayer_id, 42).await.unwrap();
4164 assert!(old_nonce_result.is_none());
4165
4166 let new_nonce_result = repo.find_by_nonce(&relayer_id, 43).await.unwrap();
4168 assert!(new_nonce_result.is_some());
4169 }
4170
4171 #[tokio::test]
4172 #[ignore = "Requires active Redis instance"]
4173 async fn test_has_entries() {
4174 let repo = setup_test_repo().await;
4175 assert!(!repo.has_entries().await.unwrap());
4176
4177 let tx_id = uuid::Uuid::new_v4().to_string();
4178 let tx = create_test_transaction(&tx_id);
4179 repo.create(tx.clone()).await.unwrap();
4180
4181 assert!(repo.has_entries().await.unwrap());
4182 }
4183
4184 #[tokio::test]
4185 #[ignore = "Requires active Redis instance"]
4186 async fn test_drop_all_entries() {
4187 let repo = setup_test_repo().await;
4188 let tx_id = uuid::Uuid::new_v4().to_string();
4189 let tx = create_test_transaction(&tx_id);
4190 repo.create(tx.clone()).await.unwrap();
4191 assert!(repo.has_entries().await.unwrap());
4192
4193 repo.drop_all_entries().await.unwrap();
4194 assert!(!repo.has_entries().await.unwrap());
4195 }
4196
4197 #[tokio::test]
4199 #[ignore = "Requires active Redis instance"]
4200 async fn test_update_status_sets_delete_at_for_final_statuses() {
4201 let _lock = ENV_MUTEX.lock().await;
4202
4203 use chrono::{DateTime, Duration, Utc};
4204 use std::env;
4205
4206 env::set_var("TRANSACTION_EXPIRATION_HOURS", "6");
4208
4209 let repo = setup_test_repo().await;
4210
4211 let final_statuses = [
4212 TransactionStatus::Canceled,
4213 TransactionStatus::Confirmed,
4214 TransactionStatus::Failed,
4215 TransactionStatus::Expired,
4216 ];
4217
4218 for (i, status) in final_statuses.iter().enumerate() {
4219 let tx_id = format!("test-final-{}-{}", i, Uuid::new_v4());
4220 let mut tx = create_test_transaction(&tx_id);
4221
4222 tx.delete_at = None;
4224 tx.status = TransactionStatus::Pending;
4225
4226 repo.create(tx).await.unwrap();
4227
4228 let before_update = Utc::now();
4229
4230 let updated = repo
4232 .update_status(tx_id.clone(), status.clone())
4233 .await
4234 .unwrap();
4235
4236 assert!(
4238 updated.delete_at.is_some(),
4239 "delete_at should be set for status: {status:?}"
4240 );
4241
4242 let delete_at_str = updated.delete_at.unwrap();
4244 let delete_at = DateTime::parse_from_rfc3339(&delete_at_str)
4245 .expect("delete_at should be valid RFC3339")
4246 .with_timezone(&Utc);
4247
4248 let duration_from_before = delete_at.signed_duration_since(before_update);
4249 let expected_duration = Duration::hours(6);
4250 let tolerance = Duration::minutes(5);
4251
4252 assert!(
4253 duration_from_before >= expected_duration - tolerance
4254 && duration_from_before <= expected_duration + tolerance,
4255 "delete_at should be approximately 6 hours from now for status: {status:?}. Duration: {duration_from_before:?}"
4256 );
4257 }
4258
4259 env::remove_var("TRANSACTION_EXPIRATION_HOURS");
4261 }
4262
4263 #[tokio::test]
4264 #[ignore = "Requires active Redis instance"]
4265 async fn test_update_status_does_not_set_delete_at_for_non_final_statuses() {
4266 let _lock = ENV_MUTEX.lock().await;
4267
4268 use std::env;
4269
4270 env::set_var("TRANSACTION_EXPIRATION_HOURS", "4");
4271
4272 let repo = setup_test_repo().await;
4273
4274 let non_final_statuses = [
4275 TransactionStatus::Pending,
4276 TransactionStatus::Sent,
4277 TransactionStatus::Submitted,
4278 TransactionStatus::Mined,
4279 ];
4280
4281 for (i, status) in non_final_statuses.iter().enumerate() {
4282 let tx_id = format!("test-non-final-{}-{}", i, Uuid::new_v4());
4283 let mut tx = create_test_transaction(&tx_id);
4284 tx.delete_at = None;
4285 tx.status = TransactionStatus::Pending;
4286
4287 repo.create(tx).await.unwrap();
4288
4289 let updated = repo
4291 .update_status(tx_id.clone(), status.clone())
4292 .await
4293 .unwrap();
4294
4295 assert!(
4297 updated.delete_at.is_none(),
4298 "delete_at should NOT be set for status: {status:?}"
4299 );
4300 }
4301
4302 env::remove_var("TRANSACTION_EXPIRATION_HOURS");
4304 }
4305
4306 #[tokio::test]
4307 #[ignore = "Requires active Redis instance"]
4308 async fn test_partial_update_sets_delete_at_for_final_statuses() {
4309 let _lock = ENV_MUTEX.lock().await;
4310
4311 use chrono::{DateTime, Duration, Utc};
4312 use std::env;
4313
4314 env::set_var("TRANSACTION_EXPIRATION_HOURS", "8");
4315
4316 let repo = setup_test_repo().await;
4317 let tx_id = format!("test-partial-final-{}", Uuid::new_v4());
4318 let mut tx = create_test_transaction(&tx_id);
4319 tx.delete_at = None;
4320 tx.status = TransactionStatus::Pending;
4321
4322 repo.create(tx).await.unwrap();
4323
4324 let before_update = Utc::now();
4325
4326 let update = TransactionUpdateRequest {
4328 status: Some(TransactionStatus::Confirmed),
4329 status_reason: Some("Transaction completed".to_string()),
4330 confirmed_at: Some("2023-01-01T12:05:00Z".to_string()),
4331 ..Default::default()
4332 };
4333
4334 let updated = repo.partial_update(tx_id.clone(), update).await.unwrap();
4335
4336 assert!(
4338 updated.delete_at.is_some(),
4339 "delete_at should be set when updating to Confirmed status"
4340 );
4341
4342 let delete_at_str = updated.delete_at.unwrap();
4344 let delete_at = DateTime::parse_from_rfc3339(&delete_at_str)
4345 .expect("delete_at should be valid RFC3339")
4346 .with_timezone(&Utc);
4347
4348 let duration_from_before = delete_at.signed_duration_since(before_update);
4349 let expected_duration = Duration::hours(8);
4350 let tolerance = Duration::minutes(5);
4351
4352 assert!(
4353 duration_from_before >= expected_duration - tolerance
4354 && duration_from_before <= expected_duration + tolerance,
4355 "delete_at should be approximately 8 hours from now. Duration: {duration_from_before:?}"
4356 );
4357
4358 assert_eq!(updated.status, TransactionStatus::Confirmed);
4360 assert_eq!(
4361 updated.status_reason,
4362 Some("Transaction completed".to_string())
4363 );
4364 assert_eq!(
4365 updated.confirmed_at,
4366 Some("2023-01-01T12:05:00Z".to_string())
4367 );
4368
4369 env::remove_var("TRANSACTION_EXPIRATION_HOURS");
4371 }
4372
4373 #[tokio::test]
4374 #[ignore = "Requires active Redis instance"]
4375 async fn test_update_status_preserves_existing_delete_at() {
4376 let _lock = ENV_MUTEX.lock().await;
4377
4378 use std::env;
4379
4380 env::set_var("TRANSACTION_EXPIRATION_HOURS", "2");
4381
4382 let repo = setup_test_repo().await;
4383 let tx_id = format!("test-preserve-delete-at-{}", Uuid::new_v4());
4384 let mut tx = create_test_transaction(&tx_id);
4385
4386 let existing_delete_at = "2025-01-01T12:00:00Z".to_string();
4388 tx.delete_at = Some(existing_delete_at.clone());
4389 tx.status = TransactionStatus::Pending;
4390
4391 repo.create(tx).await.unwrap();
4392
4393 let updated = repo
4395 .update_status(tx_id.clone(), TransactionStatus::Confirmed)
4396 .await
4397 .unwrap();
4398
4399 assert_eq!(
4401 updated.delete_at,
4402 Some(existing_delete_at),
4403 "Existing delete_at should be preserved when updating to final status"
4404 );
4405
4406 env::remove_var("TRANSACTION_EXPIRATION_HOURS");
4408 }
4409 #[tokio::test]
4410 #[ignore = "Requires active Redis instance"]
4411 async fn test_partial_update_without_status_change_preserves_delete_at() {
4412 let _lock = ENV_MUTEX.lock().await;
4413
4414 use std::env;
4415
4416 env::set_var("TRANSACTION_EXPIRATION_HOURS", "3");
4417
4418 let repo = setup_test_repo().await;
4419 let tx_id = format!("test-preserve-no-status-{}", Uuid::new_v4());
4420 let mut tx = create_test_transaction(&tx_id);
4421 tx.delete_at = None;
4422 tx.status = TransactionStatus::Pending;
4423
4424 repo.create(tx).await.unwrap();
4425
4426 let updated1 = repo
4428 .update_status(tx_id.clone(), TransactionStatus::Confirmed)
4429 .await
4430 .unwrap();
4431
4432 assert!(updated1.delete_at.is_some());
4433 let original_delete_at = updated1.delete_at.clone();
4434
4435 let update = TransactionUpdateRequest {
4437 status: None, status_reason: Some("Updated reason".to_string()),
4439 confirmed_at: Some("2023-01-01T12:10:00Z".to_string()),
4440 ..Default::default()
4441 };
4442
4443 let updated2 = repo.partial_update(tx_id.clone(), update).await.unwrap();
4444
4445 assert_eq!(
4447 updated2.delete_at, original_delete_at,
4448 "delete_at should be preserved when status is not updated"
4449 );
4450
4451 assert_eq!(updated2.status, TransactionStatus::Confirmed); assert_eq!(updated2.status_reason, Some("Updated reason".to_string()));
4454 assert_eq!(
4455 updated2.confirmed_at,
4456 Some("2023-01-01T12:10:00Z".to_string())
4457 );
4458
4459 env::remove_var("TRANSACTION_EXPIRATION_HOURS");
4461 }
4462
4463 #[tokio::test]
4466 #[ignore = "Requires active Redis instance"]
4467 async fn test_delete_by_ids_empty_list() {
4468 let repo = setup_test_repo().await;
4469 let tx_id = format!("test-empty-{}", Uuid::new_v4());
4470
4471 let tx = create_test_transaction(&tx_id);
4473 repo.create(tx).await.unwrap();
4474
4475 let result = repo.delete_by_ids(vec![]).await.unwrap();
4477
4478 assert_eq!(result.deleted_count, 0);
4479 assert!(result.failed.is_empty());
4480
4481 assert!(repo.get_by_id(tx_id).await.is_ok());
4483 }
4484
4485 #[tokio::test]
4486 #[ignore = "Requires active Redis instance"]
4487 async fn test_delete_by_ids_single_transaction() {
4488 let repo = setup_test_repo().await;
4489 let tx_id = format!("test-single-{}", Uuid::new_v4());
4490
4491 let tx = create_test_transaction(&tx_id);
4492 repo.create(tx).await.unwrap();
4493
4494 let result = repo.delete_by_ids(vec![tx_id.clone()]).await.unwrap();
4495
4496 assert_eq!(result.deleted_count, 1);
4497 assert!(result.failed.is_empty());
4498
4499 assert!(repo.get_by_id(tx_id).await.is_err());
4501 }
4502
4503 #[tokio::test]
4504 #[ignore = "Requires active Redis instance"]
4505 async fn test_delete_by_ids_multiple_transactions() {
4506 let repo = setup_test_repo().await;
4507 let base_id = Uuid::new_v4();
4508
4509 let mut created_ids = Vec::new();
4511 for i in 1..=5 {
4512 let tx_id = format!("test-multi-{base_id}-{i}");
4513 let tx = create_test_transaction(&tx_id);
4514 repo.create(tx).await.unwrap();
4515 created_ids.push(tx_id);
4516 }
4517
4518 let ids_to_delete = vec![
4520 created_ids[0].clone(),
4521 created_ids[2].clone(),
4522 created_ids[4].clone(),
4523 ];
4524 let result = repo.delete_by_ids(ids_to_delete).await.unwrap();
4525
4526 assert_eq!(result.deleted_count, 3);
4527 assert!(result.failed.is_empty());
4528
4529 assert!(repo.get_by_id(created_ids[0].clone()).await.is_err());
4531 assert!(repo.get_by_id(created_ids[1].clone()).await.is_ok()); assert!(repo.get_by_id(created_ids[2].clone()).await.is_err());
4533 assert!(repo.get_by_id(created_ids[3].clone()).await.is_ok()); assert!(repo.get_by_id(created_ids[4].clone()).await.is_err());
4535 }
4536
4537 #[tokio::test]
4538 #[ignore = "Requires active Redis instance"]
4539 async fn test_delete_by_ids_nonexistent_transactions() {
4540 let repo = setup_test_repo().await;
4541 let base_id = Uuid::new_v4();
4542
4543 let ids_to_delete = vec![
4545 format!("nonexistent-{}-1", base_id),
4546 format!("nonexistent-{}-2", base_id),
4547 ];
4548 let result = repo.delete_by_ids(ids_to_delete.clone()).await.unwrap();
4549
4550 assert_eq!(result.deleted_count, 0);
4551 assert_eq!(result.failed.len(), 2);
4552
4553 let failed_ids: Vec<&String> = result.failed.iter().map(|(id, _)| id).collect();
4555 assert!(failed_ids.contains(&&ids_to_delete[0]));
4556 assert!(failed_ids.contains(&&ids_to_delete[1]));
4557 }
4558
4559 #[tokio::test]
4560 #[ignore = "Requires active Redis instance"]
4561 async fn test_delete_by_ids_mixed_existing_and_nonexistent() {
4562 let repo = setup_test_repo().await;
4563 let base_id = Uuid::new_v4();
4564
4565 let existing_ids: Vec<String> = (1..=3)
4567 .map(|i| format!("test-mixed-existing-{base_id}-{i}"))
4568 .collect();
4569
4570 for id in &existing_ids {
4571 let tx = create_test_transaction(id);
4572 repo.create(tx).await.unwrap();
4573 }
4574
4575 let nonexistent_ids: Vec<String> = (1..=2)
4576 .map(|i| format!("test-mixed-nonexistent-{base_id}-{i}"))
4577 .collect();
4578
4579 let ids_to_delete = vec![
4581 existing_ids[0].clone(),
4582 nonexistent_ids[0].clone(),
4583 existing_ids[1].clone(),
4584 nonexistent_ids[1].clone(),
4585 ];
4586 let result = repo.delete_by_ids(ids_to_delete).await.unwrap();
4587
4588 assert_eq!(result.deleted_count, 2);
4589 assert_eq!(result.failed.len(), 2);
4590
4591 assert!(repo.get_by_id(existing_ids[0].clone()).await.is_err());
4593 assert!(repo.get_by_id(existing_ids[1].clone()).await.is_err());
4594
4595 assert!(repo.get_by_id(existing_ids[2].clone()).await.is_ok());
4597 }
4598
4599 #[tokio::test]
4600 #[ignore = "Requires active Redis instance"]
4601 async fn test_delete_by_ids_removes_all_indexes() {
4602 let repo = setup_test_repo().await;
4603 let relayer_id = format!("relayer-{}", Uuid::new_v4());
4604 let tx_id = format!("test-indexes-{}", Uuid::new_v4());
4605
4606 let mut tx = create_test_transaction(&tx_id);
4608 tx.relayer_id = relayer_id.clone();
4609 tx.status = TransactionStatus::Confirmed;
4610 repo.create(tx).await.unwrap();
4611
4612 let found = repo
4614 .find_by_status(&relayer_id, &[TransactionStatus::Confirmed])
4615 .await
4616 .unwrap();
4617 assert!(found.iter().any(|t| t.id == tx_id));
4618
4619 let result = repo.delete_by_ids(vec![tx_id.clone()]).await.unwrap();
4621 assert_eq!(result.deleted_count, 1);
4622
4623 let found_after = repo
4625 .find_by_status(&relayer_id, &[TransactionStatus::Confirmed])
4626 .await
4627 .unwrap();
4628 assert!(!found_after.iter().any(|t| t.id == tx_id));
4629
4630 assert!(repo.get_by_id(tx_id).await.is_err());
4632 }
4633
4634 #[tokio::test]
4635 #[ignore = "Requires active Redis instance"]
4636 async fn test_delete_by_ids_removes_nonce_index() {
4637 let repo = setup_test_repo().await;
4638 let relayer_id = format!("relayer-{}", Uuid::new_v4());
4639 let tx_id = format!("test-nonce-{}", Uuid::new_v4());
4640 let nonce = 12345u64;
4641
4642 let tx = create_test_transaction_with_nonce(&tx_id, nonce, &relayer_id);
4644 repo.create(tx).await.unwrap();
4645
4646 let found = repo.find_by_nonce(&relayer_id, nonce).await.unwrap();
4648 assert!(found.is_some());
4649 assert_eq!(found.unwrap().id, tx_id);
4650
4651 let result = repo.delete_by_ids(vec![tx_id.clone()]).await.unwrap();
4653 assert_eq!(result.deleted_count, 1);
4654
4655 let found_after = repo.find_by_nonce(&relayer_id, nonce).await.unwrap();
4657 assert!(found_after.is_none());
4658 }
4659
4660 #[tokio::test]
4661 #[ignore = "Requires active Redis instance"]
4662 async fn test_delete_by_ids_large_batch() {
4663 let repo = setup_test_repo().await;
4664 let base_id = Uuid::new_v4();
4665
4666 let count = 50;
4668 let mut created_ids = Vec::new();
4669
4670 for i in 0..count {
4671 let tx_id = format!("test-large-{base_id}-{i}");
4672 let tx = create_test_transaction(&tx_id);
4673 repo.create(tx).await.unwrap();
4674 created_ids.push(tx_id);
4675 }
4676
4677 let result = repo.delete_by_ids(created_ids.clone()).await.unwrap();
4679
4680 assert_eq!(result.deleted_count, count);
4681 assert!(result.failed.is_empty());
4682
4683 for id in created_ids {
4685 assert!(repo.get_by_id(id).await.is_err());
4686 }
4687 }
4688
4689 #[tokio::test]
4690 #[ignore = "Requires active Redis instance"]
4691 async fn test_delete_by_ids_preserves_other_relayer_transactions() {
4692 let repo = setup_test_repo().await;
4693 let relayer_1 = format!("relayer-1-{}", Uuid::new_v4());
4694 let relayer_2 = format!("relayer-2-{}", Uuid::new_v4());
4695 let tx_id_1 = format!("tx-relayer-1-{}", Uuid::new_v4());
4696 let tx_id_2 = format!("tx-relayer-2-{}", Uuid::new_v4());
4697
4698 let tx1 = create_test_transaction_with_relayer(&tx_id_1, &relayer_1);
4700 let tx2 = create_test_transaction_with_relayer(&tx_id_2, &relayer_2);
4701
4702 repo.create(tx1).await.unwrap();
4703 repo.create(tx2).await.unwrap();
4704
4705 let result = repo.delete_by_ids(vec![tx_id_1.clone()]).await.unwrap();
4707
4708 assert_eq!(result.deleted_count, 1);
4709
4710 assert!(repo.get_by_id(tx_id_1).await.is_err());
4712
4713 let remaining = repo.get_by_id(tx_id_2).await.unwrap();
4715 assert_eq!(remaining.relayer_id, relayer_2);
4716 }
4717
4718 #[tokio::test]
4721 #[ignore = "Requires active Redis instance"]
4722 async fn test_increment_status_check_failures_no_prior_metadata() {
4723 let _lock = ENV_MUTEX.lock().await;
4724 let repo = setup_test_repo().await;
4725 let relayer_id = Uuid::new_v4().to_string();
4726 let tx_id = Uuid::new_v4().to_string();
4727 let mut tx =
4728 create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Sent);
4729 tx.metadata = None;
4730 repo.create(tx).await.unwrap();
4731
4732 let updated = repo.increment_status_check_failures(tx_id).await.unwrap();
4733
4734 let meta = updated.metadata.expect("metadata should be set");
4735 assert_eq!(meta.consecutive_failures, 1);
4736 assert_eq!(meta.total_failures, 1);
4737 assert_eq!(meta.insufficient_fee_retries, 0);
4738 }
4739
4740 #[tokio::test]
4741 #[ignore = "Requires active Redis instance"]
4742 async fn test_increment_status_check_failures_accumulates() {
4743 let _lock = ENV_MUTEX.lock().await;
4744 let repo = setup_test_repo().await;
4745 let relayer_id = Uuid::new_v4().to_string();
4746 let tx_id = Uuid::new_v4().to_string();
4747 let tx = create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Sent);
4748 repo.create(tx).await.unwrap();
4749
4750 repo.increment_status_check_failures(tx_id.clone())
4751 .await
4752 .unwrap();
4753 repo.increment_status_check_failures(tx_id.clone())
4754 .await
4755 .unwrap();
4756 let updated = repo.increment_status_check_failures(tx_id).await.unwrap();
4757
4758 let meta = updated.metadata.unwrap();
4759 assert_eq!(meta.consecutive_failures, 3);
4760 assert_eq!(meta.total_failures, 3);
4761 }
4762
4763 #[tokio::test]
4764 #[ignore = "Requires active Redis instance"]
4765 async fn test_increment_status_check_failures_noop_on_final_state() {
4766 let _lock = ENV_MUTEX.lock().await;
4767 let repo = setup_test_repo().await;
4768 let relayer_id = Uuid::new_v4().to_string();
4769 let tx_id = Uuid::new_v4().to_string();
4770 let tx =
4771 create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Confirmed);
4772 repo.create(tx).await.unwrap();
4773
4774 let result = repo.increment_status_check_failures(tx_id).await.unwrap();
4775
4776 assert!(result.metadata.is_none());
4778 assert_eq!(result.status, TransactionStatus::Confirmed);
4779 }
4780
4781 #[tokio::test]
4782 #[ignore = "Requires active Redis instance"]
4783 async fn test_increment_status_check_failures_not_found() {
4784 let _lock = ENV_MUTEX.lock().await;
4785 let repo = setup_test_repo().await;
4786
4787 let result = repo
4788 .increment_status_check_failures("nonexistent".to_string())
4789 .await;
4790
4791 assert!(matches!(result, Err(RepositoryError::NotFound(_))));
4792 }
4793
4794 #[tokio::test]
4797 #[ignore = "Requires active Redis instance"]
4798 async fn test_reset_consecutive_failures() {
4799 let _lock = ENV_MUTEX.lock().await;
4800 let repo = setup_test_repo().await;
4801 let relayer_id = Uuid::new_v4().to_string();
4802 let tx_id = Uuid::new_v4().to_string();
4803 let tx = create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Sent);
4804 repo.create(tx).await.unwrap();
4805
4806 repo.increment_status_check_failures(tx_id.clone())
4808 .await
4809 .unwrap();
4810 repo.increment_status_check_failures(tx_id.clone())
4811 .await
4812 .unwrap();
4813
4814 let updated = repo
4815 .reset_status_check_consecutive_failures(tx_id)
4816 .await
4817 .unwrap();
4818
4819 let meta = updated.metadata.unwrap();
4820 assert_eq!(meta.consecutive_failures, 0);
4821 assert_eq!(meta.total_failures, 2);
4823 }
4824
4825 #[tokio::test]
4826 #[ignore = "Requires active Redis instance"]
4827 async fn test_reset_consecutive_failures_noop_on_final_state() {
4828 let _lock = ENV_MUTEX.lock().await;
4829 let repo = setup_test_repo().await;
4830 let relayer_id = Uuid::new_v4().to_string();
4831 let tx_id = Uuid::new_v4().to_string();
4832 let mut tx =
4833 create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Failed);
4834 tx.metadata = Some(crate::models::TransactionMetadata {
4835 consecutive_failures: 5,
4836 total_failures: 10,
4837 insufficient_fee_retries: 0,
4838 try_again_later_retries: 0,
4839 nonce_too_high_retries: 0,
4840 });
4841 repo.create(tx).await.unwrap();
4842
4843 let result = repo
4844 .reset_status_check_consecutive_failures(tx_id)
4845 .await
4846 .unwrap();
4847
4848 let meta = result.metadata.unwrap();
4850 assert_eq!(meta.consecutive_failures, 5);
4851 }
4852
4853 #[tokio::test]
4854 #[ignore = "Requires active Redis instance"]
4855 async fn test_reset_consecutive_failures_not_found() {
4856 let _lock = ENV_MUTEX.lock().await;
4857 let repo = setup_test_repo().await;
4858
4859 let result = repo
4860 .reset_status_check_consecutive_failures("nonexistent".to_string())
4861 .await;
4862
4863 assert!(matches!(result, Err(RepositoryError::NotFound(_))));
4864 }
4865
4866 #[tokio::test]
4869 #[ignore = "Requires active Redis instance"]
4870 async fn test_record_insufficient_fee_retry() {
4871 let _lock = ENV_MUTEX.lock().await;
4872 let repo = setup_test_repo().await;
4873 let relayer_id = Uuid::new_v4().to_string();
4874 let tx_id = Uuid::new_v4().to_string();
4875 let mut tx =
4876 create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Sent);
4877 tx.sent_at = None;
4878 repo.create(tx).await.unwrap();
4879
4880 let updated = repo
4881 .record_stellar_insufficient_fee_retry(tx_id, "2025-03-18T10:00:00Z".to_string())
4882 .await
4883 .unwrap();
4884
4885 assert_eq!(updated.sent_at.as_deref(), Some("2025-03-18T10:00:00Z"));
4886 let meta = updated.metadata.unwrap();
4887 assert_eq!(meta.insufficient_fee_retries, 1);
4888 assert_eq!(meta.consecutive_failures, 0);
4889 assert_eq!(meta.total_failures, 0);
4890 }
4891
4892 #[tokio::test]
4893 #[ignore = "Requires active Redis instance"]
4894 async fn test_record_insufficient_fee_retry_accumulates() {
4895 let _lock = ENV_MUTEX.lock().await;
4896 let repo = setup_test_repo().await;
4897 let relayer_id = Uuid::new_v4().to_string();
4898 let tx_id = Uuid::new_v4().to_string();
4899 let tx = create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Sent);
4900 repo.create(tx).await.unwrap();
4901
4902 repo.record_stellar_insufficient_fee_retry(
4903 tx_id.clone(),
4904 "2025-03-18T10:00:00Z".to_string(),
4905 )
4906 .await
4907 .unwrap();
4908
4909 let updated = repo
4910 .record_stellar_insufficient_fee_retry(tx_id, "2025-03-18T10:01:00Z".to_string())
4911 .await
4912 .unwrap();
4913
4914 assert_eq!(updated.sent_at.as_deref(), Some("2025-03-18T10:01:00Z"));
4915 let meta = updated.metadata.unwrap();
4916 assert_eq!(meta.insufficient_fee_retries, 2);
4917 }
4918
4919 #[tokio::test]
4920 #[ignore = "Requires active Redis instance"]
4921 async fn test_record_insufficient_fee_retry_noop_on_final_state() {
4922 let _lock = ENV_MUTEX.lock().await;
4923 let repo = setup_test_repo().await;
4924 let relayer_id = Uuid::new_v4().to_string();
4925 let tx_id = Uuid::new_v4().to_string();
4926 let mut tx =
4927 create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Confirmed);
4928 tx.sent_at = Some("old-time".to_string());
4929 repo.create(tx).await.unwrap();
4930
4931 let result = repo
4932 .record_stellar_insufficient_fee_retry(tx_id, "new-time".to_string())
4933 .await
4934 .unwrap();
4935
4936 assert_eq!(result.sent_at.as_deref(), Some("old-time"));
4938 assert!(result.metadata.is_none());
4939 }
4940
4941 #[tokio::test]
4942 #[ignore = "Requires active Redis instance"]
4943 async fn test_record_insufficient_fee_retry_not_found() {
4944 let _lock = ENV_MUTEX.lock().await;
4945 let repo = setup_test_repo().await;
4946
4947 let result = repo
4948 .record_stellar_insufficient_fee_retry(
4949 "nonexistent".to_string(),
4950 "2025-03-18T10:00:00Z".to_string(),
4951 )
4952 .await;
4953
4954 assert!(matches!(result, Err(RepositoryError::NotFound(_))));
4955 }
4956
4957 #[tokio::test]
4960 #[ignore = "Requires active Redis instance"]
4961 async fn test_record_try_again_later_retry() {
4962 let _lock = ENV_MUTEX.lock().await;
4963 let repo = setup_test_repo().await;
4964 let relayer_id = Uuid::new_v4().to_string();
4965 let tx_id = Uuid::new_v4().to_string();
4966 let mut tx =
4967 create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Sent);
4968 tx.sent_at = None;
4969 repo.create(tx).await.unwrap();
4970
4971 let updated = repo
4972 .record_stellar_try_again_later_retry(tx_id, "2025-03-18T10:00:00Z".to_string())
4973 .await
4974 .unwrap();
4975
4976 assert_eq!(updated.sent_at.as_deref(), Some("2025-03-18T10:00:00Z"));
4977 let meta = updated.metadata.unwrap();
4978 assert_eq!(meta.try_again_later_retries, 1);
4979 assert_eq!(meta.consecutive_failures, 0);
4980 assert_eq!(meta.total_failures, 0);
4981 }
4982
4983 #[tokio::test]
4984 #[ignore = "Requires active Redis instance"]
4985 async fn test_record_try_again_later_retry_accumulates() {
4986 let _lock = ENV_MUTEX.lock().await;
4987 let repo = setup_test_repo().await;
4988 let relayer_id = Uuid::new_v4().to_string();
4989 let tx_id = Uuid::new_v4().to_string();
4990 let tx = create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Sent);
4991 repo.create(tx).await.unwrap();
4992
4993 repo.record_stellar_try_again_later_retry(
4994 tx_id.clone(),
4995 "2025-03-18T10:00:00Z".to_string(),
4996 )
4997 .await
4998 .unwrap();
4999
5000 let updated = repo
5001 .record_stellar_try_again_later_retry(tx_id, "2025-03-18T10:01:00Z".to_string())
5002 .await
5003 .unwrap();
5004
5005 assert_eq!(updated.sent_at.as_deref(), Some("2025-03-18T10:01:00Z"));
5006 let meta = updated.metadata.unwrap();
5007 assert_eq!(meta.try_again_later_retries, 2);
5008 }
5009
5010 #[tokio::test]
5011 #[ignore = "Requires active Redis instance"]
5012 async fn test_record_try_again_later_retry_noop_on_final_state() {
5013 let _lock = ENV_MUTEX.lock().await;
5014 let repo = setup_test_repo().await;
5015 let relayer_id = Uuid::new_v4().to_string();
5016 let tx_id = Uuid::new_v4().to_string();
5017 let mut tx =
5018 create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Confirmed);
5019 tx.sent_at = Some("old-time".to_string());
5020 repo.create(tx).await.unwrap();
5021
5022 let result = repo
5023 .record_stellar_try_again_later_retry(tx_id, "new-time".to_string())
5024 .await
5025 .unwrap();
5026
5027 assert_eq!(result.sent_at.as_deref(), Some("old-time"));
5029 assert!(result.metadata.is_none());
5030 }
5031
5032 #[tokio::test]
5033 #[ignore = "Requires active Redis instance"]
5034 async fn test_record_try_again_later_retry_not_found() {
5035 let _lock = ENV_MUTEX.lock().await;
5036 let repo = setup_test_repo().await;
5037
5038 let result = repo
5039 .record_stellar_try_again_later_retry(
5040 "nonexistent".to_string(),
5041 "2025-03-18T10:00:00Z".to_string(),
5042 )
5043 .await;
5044
5045 assert!(matches!(result, Err(RepositoryError::NotFound(_))));
5046 }
5047
5048 #[tokio::test]
5051 #[ignore = "Requires active Redis instance"]
5052 async fn test_increment_failures_preserves_try_again_later_retries() {
5053 let _lock = ENV_MUTEX.lock().await;
5054 let repo = setup_test_repo().await;
5055 let relayer_id = Uuid::new_v4().to_string();
5056 let tx_id = Uuid::new_v4().to_string();
5057 let tx = create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Sent);
5058 repo.create(tx).await.unwrap();
5059
5060 repo.record_stellar_try_again_later_retry(
5062 tx_id.clone(),
5063 "2025-03-18T10:00:00Z".to_string(),
5064 )
5065 .await
5066 .unwrap();
5067
5068 let updated = repo.increment_status_check_failures(tx_id).await.unwrap();
5070
5071 let meta = updated.metadata.unwrap();
5072 assert_eq!(
5073 meta.try_again_later_retries, 1,
5074 "try_again_later_retries must survive increment_status_check_failures"
5075 );
5076 assert_eq!(meta.consecutive_failures, 1);
5077 assert_eq!(meta.total_failures, 1);
5078 }
5079
5080 #[tokio::test]
5081 #[ignore = "Requires active Redis instance"]
5082 async fn test_increment_failures_preserves_insufficient_fee_retries() {
5083 let _lock = ENV_MUTEX.lock().await;
5084 let repo = setup_test_repo().await;
5085 let relayer_id = Uuid::new_v4().to_string();
5086 let tx_id = Uuid::new_v4().to_string();
5087 let tx = create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Sent);
5088 repo.create(tx).await.unwrap();
5089
5090 repo.record_stellar_insufficient_fee_retry(
5092 tx_id.clone(),
5093 "2025-03-18T10:00:00Z".to_string(),
5094 )
5095 .await
5096 .unwrap();
5097
5098 let updated = repo.increment_status_check_failures(tx_id).await.unwrap();
5100
5101 let meta = updated.metadata.unwrap();
5102 assert_eq!(
5103 meta.insufficient_fee_retries, 1,
5104 "insufficient_fee_retries must survive increment_status_check_failures"
5105 );
5106 assert_eq!(meta.consecutive_failures, 1);
5107 }
5108
5109 #[tokio::test]
5110 #[ignore = "Requires active Redis instance"]
5111 async fn test_reset_failures_preserves_retry_counters() {
5112 let _lock = ENV_MUTEX.lock().await;
5113 let repo = setup_test_repo().await;
5114 let relayer_id = Uuid::new_v4().to_string();
5115 let tx_id = Uuid::new_v4().to_string();
5116 let tx = create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Sent);
5117 repo.create(tx).await.unwrap();
5118
5119 repo.record_stellar_try_again_later_retry(
5121 tx_id.clone(),
5122 "2025-03-18T10:00:00Z".to_string(),
5123 )
5124 .await
5125 .unwrap();
5126 repo.record_stellar_insufficient_fee_retry(
5127 tx_id.clone(),
5128 "2025-03-18T10:01:00Z".to_string(),
5129 )
5130 .await
5131 .unwrap();
5132
5133 repo.increment_status_check_failures(tx_id.clone())
5135 .await
5136 .unwrap();
5137 let updated = repo
5138 .reset_status_check_consecutive_failures(tx_id)
5139 .await
5140 .unwrap();
5141
5142 let meta = updated.metadata.unwrap();
5143 assert_eq!(meta.consecutive_failures, 0);
5144 assert_eq!(meta.total_failures, 1);
5145 assert_eq!(
5146 meta.try_again_later_retries, 1,
5147 "try_again_later_retries must survive reset"
5148 );
5149 assert_eq!(
5150 meta.insufficient_fee_retries, 1,
5151 "insufficient_fee_retries must survive reset"
5152 );
5153 }
5154
5155 #[tokio::test]
5156 #[ignore = "Requires active Redis instance"]
5157 async fn test_fee_and_try_again_later_retries_independent() {
5158 let _lock = ENV_MUTEX.lock().await;
5159 let repo = setup_test_repo().await;
5160 let relayer_id = Uuid::new_v4().to_string();
5161 let tx_id = Uuid::new_v4().to_string();
5162 let tx = create_test_transaction_with_status(&tx_id, &relayer_id, TransactionStatus::Sent);
5163 repo.create(tx).await.unwrap();
5164
5165 repo.record_stellar_try_again_later_retry(
5167 tx_id.clone(),
5168 "2025-03-18T10:00:00Z".to_string(),
5169 )
5170 .await
5171 .unwrap();
5172 repo.record_stellar_try_again_later_retry(
5173 tx_id.clone(),
5174 "2025-03-18T10:01:00Z".to_string(),
5175 )
5176 .await
5177 .unwrap();
5178
5179 let updated = repo
5181 .record_stellar_insufficient_fee_retry(tx_id, "2025-03-18T10:02:00Z".to_string())
5182 .await
5183 .unwrap();
5184
5185 let meta = updated.metadata.unwrap();
5186 assert_eq!(
5187 meta.try_again_later_retries, 2,
5188 "try_again_later_retries must survive insufficient_fee_retry"
5189 );
5190 assert_eq!(meta.insufficient_fee_retries, 1);
5191 }
5192}