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