1use crate::{
8 models::{
9 NetworkTransactionData, TransactionRepoModel, TransactionStatus, TransactionUpdateRequest,
10 },
11 repositories::*,
12};
13use async_trait::async_trait;
14use eyre::Result;
15use itertools::Itertools;
16use std::collections::HashMap;
17use tokio::sync::{Mutex, MutexGuard};
18
19#[derive(Debug)]
20pub struct InMemoryTransactionRepository {
21 store: Mutex<HashMap<String, TransactionRepoModel>>,
22}
23
24impl Clone for InMemoryTransactionRepository {
25 fn clone(&self) -> Self {
26 let data = self
28 .store
29 .try_lock()
30 .map(|guard| guard.clone())
31 .unwrap_or_else(|_| HashMap::new());
32
33 Self {
34 store: Mutex::new(data),
35 }
36 }
37}
38
39impl InMemoryTransactionRepository {
40 pub fn new() -> Self {
41 Self {
42 store: Mutex::new(HashMap::new()),
43 }
44 }
45
46 async fn acquire_lock<T>(lock: &Mutex<T>) -> Result<MutexGuard<'_, T>, RepositoryError> {
47 Ok(lock.lock().await)
48 }
49
50 fn get_sort_key(tx: &TransactionRepoModel) -> (&str, bool) {
56 if tx.status == TransactionStatus::Confirmed {
57 if let Some(ref confirmed_at) = tx.confirmed_at {
58 return (confirmed_at, true);
59 }
60 }
62 (&tx.created_at, false)
63 }
64
65 fn compare_for_sort(a: &TransactionRepoModel, b: &TransactionRepoModel) -> std::cmp::Ordering {
68 let (a_key, _) = Self::get_sort_key(a);
69 let (b_key, _) = Self::get_sort_key(b);
70 b_key
71 .cmp(a_key) .then_with(|| b.id.cmp(&a.id)) }
74
75 fn is_final_state(status: &TransactionStatus) -> bool {
76 matches!(
77 status,
78 TransactionStatus::Confirmed
79 | TransactionStatus::Failed
80 | TransactionStatus::Expired
81 | TransactionStatus::Canceled
82 )
83 }
84}
85
86#[async_trait]
89impl Repository<TransactionRepoModel, String> for InMemoryTransactionRepository {
90 async fn create(
91 &self,
92 tx: TransactionRepoModel,
93 ) -> Result<TransactionRepoModel, RepositoryError> {
94 let mut store = Self::acquire_lock(&self.store).await?;
95 if store.contains_key(&tx.id) {
96 return Err(RepositoryError::ConstraintViolation(format!(
97 "Transaction with ID {} already exists",
98 tx.id
99 )));
100 }
101 store.insert(tx.id.clone(), tx.clone());
102 Ok(tx)
103 }
104
105 async fn get_by_id(&self, id: String) -> Result<TransactionRepoModel, RepositoryError> {
106 let store = Self::acquire_lock(&self.store).await?;
107 store
108 .get(&id)
109 .cloned()
110 .ok_or_else(|| RepositoryError::NotFound(format!("Transaction with ID {id} not found")))
111 }
112
113 #[allow(clippy::map_entry)]
114 async fn update(
115 &self,
116 id: String,
117 tx: TransactionRepoModel,
118 ) -> Result<TransactionRepoModel, RepositoryError> {
119 let mut store = Self::acquire_lock(&self.store).await?;
120 if store.contains_key(&id) {
121 let mut updated_tx = tx;
122 updated_tx.id = id.clone();
123 store.insert(id, updated_tx.clone());
124 Ok(updated_tx)
125 } else {
126 Err(RepositoryError::NotFound(format!(
127 "Transaction with ID {id} not found"
128 )))
129 }
130 }
131
132 async fn delete_by_id(&self, id: String) -> Result<(), RepositoryError> {
133 let mut store = Self::acquire_lock(&self.store).await?;
134 if store.remove(&id).is_some() {
135 Ok(())
136 } else {
137 Err(RepositoryError::NotFound(format!(
138 "Transaction with ID {id} not found"
139 )))
140 }
141 }
142
143 async fn list_all(&self) -> Result<Vec<TransactionRepoModel>, RepositoryError> {
144 let store = Self::acquire_lock(&self.store).await?;
145 Ok(store.values().cloned().collect())
146 }
147
148 async fn list_paginated(
149 &self,
150 query: PaginationQuery,
151 ) -> Result<PaginatedResult<TransactionRepoModel>, RepositoryError> {
152 let total = self.count().await?;
153 let start = ((query.page - 1) * query.per_page) as usize;
154 let store = Self::acquire_lock(&self.store).await?;
155 let items: Vec<TransactionRepoModel> = store
156 .values()
157 .skip(start)
158 .take(query.per_page as usize)
159 .cloned()
160 .collect();
161
162 Ok(PaginatedResult {
163 items,
164 total: total as u64,
165 page: query.page,
166 per_page: query.per_page,
167 })
168 }
169
170 async fn count(&self) -> Result<usize, RepositoryError> {
171 let store = Self::acquire_lock(&self.store).await?;
172 Ok(store.len())
173 }
174
175 async fn has_entries(&self) -> Result<bool, RepositoryError> {
176 let store = Self::acquire_lock(&self.store).await?;
177 Ok(!store.is_empty())
178 }
179
180 async fn drop_all_entries(&self) -> Result<(), RepositoryError> {
181 let mut store = Self::acquire_lock(&self.store).await?;
182 store.clear();
183 Ok(())
184 }
185}
186
187#[async_trait]
188impl TransactionRepository for InMemoryTransactionRepository {
189 async fn find_by_relayer_id(
190 &self,
191 relayer_id: &str,
192 query: PaginationQuery,
193 ) -> Result<PaginatedResult<TransactionRepoModel>, RepositoryError> {
194 let store = Self::acquire_lock(&self.store).await?;
195 let filtered: Vec<TransactionRepoModel> = store
196 .values()
197 .filter(|tx| tx.relayer_id == relayer_id)
198 .cloned()
199 .collect();
200
201 let total = filtered.len() as u64;
202
203 if total == 0 {
204 return Ok(PaginatedResult::<TransactionRepoModel> {
205 items: vec![],
206 total: 0,
207 page: query.page,
208 per_page: query.per_page,
209 });
210 }
211
212 let start = ((query.page - 1) * query.per_page) as usize;
213
214 let items = filtered
216 .into_iter()
217 .sorted_by(|a, b| b.created_at.cmp(&a.created_at)) .skip(start)
219 .take(query.per_page as usize)
220 .collect();
221
222 Ok(PaginatedResult {
223 items,
224 total,
225 page: query.page,
226 per_page: query.per_page,
227 })
228 }
229
230 async fn find_by_status(
231 &self,
232 relayer_id: &str,
233 statuses: &[TransactionStatus],
234 ) -> Result<Vec<TransactionRepoModel>, RepositoryError> {
235 let store = Self::acquire_lock(&self.store).await?;
236 let filtered: Vec<TransactionRepoModel> = store
237 .values()
238 .filter(|tx| tx.relayer_id == relayer_id && statuses.contains(&tx.status))
239 .cloned()
240 .collect();
241
242 let sorted = filtered
244 .into_iter()
245 .sorted_by(|a, b| b.created_at.cmp(&a.created_at))
246 .collect();
247
248 Ok(sorted)
249 }
250
251 async fn find_by_status_paginated(
252 &self,
253 relayer_id: &str,
254 statuses: &[TransactionStatus],
255 query: PaginationQuery,
256 oldest_first: bool,
257 ) -> Result<PaginatedResult<TransactionRepoModel>, RepositoryError> {
258 let store = Self::acquire_lock(&self.store).await?;
259
260 let filtered: Vec<TransactionRepoModel> = store
262 .values()
263 .filter(|tx| tx.relayer_id == relayer_id && statuses.contains(&tx.status))
264 .cloned()
265 .collect();
266
267 let total = filtered.len() as u64;
268 let start = ((query.page.saturating_sub(1)) * query.per_page) as usize;
269
270 let items: Vec<TransactionRepoModel> = if oldest_first {
273 filtered
274 .into_iter()
275 .sorted_by(|a, b| {
276 let (a_key, _) = Self::get_sort_key(a);
277 let (b_key, _) = Self::get_sort_key(b);
278 a_key
279 .cmp(b_key) .then_with(|| a.id.cmp(&b.id)) })
282 .skip(start)
283 .take(query.per_page as usize)
284 .collect()
285 } else {
286 filtered
287 .into_iter()
288 .sorted_by(Self::compare_for_sort) .skip(start)
290 .take(query.per_page as usize)
291 .collect()
292 };
293
294 Ok(PaginatedResult {
295 items,
296 total,
297 page: query.page,
298 per_page: query.per_page,
299 })
300 }
301
302 async fn find_by_status_paginated_filtered(
303 &self,
304 relayer_id: &str,
305 statuses: &[TransactionStatus],
306 query: PaginationQuery,
307 oldest_first: bool,
308 exclude_canceled: bool,
309 ) -> Result<PaginatedResult<TransactionRepoModel>, RepositoryError> {
310 if !exclude_canceled {
311 return self
312 .find_by_status_paginated(relayer_id, statuses, query, oldest_first)
313 .await;
314 }
315
316 let store = Self::acquire_lock(&self.store).await?;
317
318 let filtered: Vec<TransactionRepoModel> = store
321 .values()
322 .filter(|tx| {
323 tx.relayer_id == relayer_id
324 && statuses.contains(&tx.status)
325 && tx.is_canceled != Some(true)
326 })
327 .cloned()
328 .collect();
329
330 let total = filtered.len() as u64;
331 let start = ((query.page.saturating_sub(1)) * query.per_page) as usize;
332
333 let items: Vec<TransactionRepoModel> = if oldest_first {
334 filtered
335 .into_iter()
336 .sorted_by(|a, b| {
337 let (a_key, _) = Self::get_sort_key(a);
338 let (b_key, _) = Self::get_sort_key(b);
339 a_key.cmp(b_key).then_with(|| a.id.cmp(&b.id))
340 })
341 .skip(start)
342 .take(query.per_page as usize)
343 .collect()
344 } else {
345 filtered
346 .into_iter()
347 .sorted_by(Self::compare_for_sort)
348 .skip(start)
349 .take(query.per_page as usize)
350 .collect()
351 };
352
353 Ok(PaginatedResult {
354 items,
355 total,
356 page: query.page,
357 per_page: query.per_page,
358 })
359 }
360
361 async fn find_by_nonce(
362 &self,
363 relayer_id: &str,
364 nonce: u64,
365 ) -> Result<Option<TransactionRepoModel>, RepositoryError> {
366 let store = Self::acquire_lock(&self.store).await?;
367 let filtered: Vec<TransactionRepoModel> = store
368 .values()
369 .filter(|tx| {
370 tx.relayer_id == relayer_id
371 && match &tx.network_data {
372 NetworkTransactionData::Evm(data) => data.nonce == Some(nonce),
373 _ => false,
374 }
375 })
376 .cloned()
377 .collect();
378
379 Ok(filtered.into_iter().next())
380 }
381
382 async fn get_nonce_occupancy(
383 &self,
384 relayer_id: &str,
385 from_nonce: u64,
386 to_nonce: u64,
387 ) -> Result<Vec<(u64, Option<TransactionStatus>)>, RepositoryError> {
388 let mut results = Vec::new();
389 for nonce in from_nonce..to_nonce {
390 let tx = self.find_by_nonce(relayer_id, nonce).await?;
391 results.push((nonce, tx.map(|t| t.status)));
392 }
393 Ok(results)
394 }
395
396 async fn update_status(
397 &self,
398 tx_id: String,
399 status: TransactionStatus,
400 ) -> Result<TransactionRepoModel, RepositoryError> {
401 let update = TransactionUpdateRequest {
402 status: Some(status),
403 ..Default::default()
404 };
405 self.partial_update(tx_id, update).await
406 }
407
408 async fn partial_update(
409 &self,
410 tx_id: String,
411 update: TransactionUpdateRequest,
412 ) -> Result<TransactionRepoModel, RepositoryError> {
413 let mut store = Self::acquire_lock(&self.store).await?;
414
415 if let Some(tx) = store.get_mut(&tx_id) {
416 tx.apply_partial_update(update);
418 Ok(tx.clone())
419 } else {
420 Err(RepositoryError::NotFound(format!(
421 "Transaction with ID {tx_id} not found"
422 )))
423 }
424 }
425
426 async fn update_network_data(
427 &self,
428 tx_id: String,
429 network_data: NetworkTransactionData,
430 ) -> Result<TransactionRepoModel, RepositoryError> {
431 let mut tx = self.get_by_id(tx_id.clone()).await?;
432 tx.network_data = network_data;
433 self.update(tx_id, tx).await
434 }
435
436 async fn set_sent_at(
437 &self,
438 tx_id: String,
439 sent_at: String,
440 ) -> Result<TransactionRepoModel, RepositoryError> {
441 let update = TransactionUpdateRequest {
442 sent_at: Some(sent_at),
443 ..Default::default()
444 };
445 self.partial_update(tx_id, update).await
446 }
447
448 async fn increment_status_check_failures(
449 &self,
450 tx_id: String,
451 ) -> Result<TransactionRepoModel, RepositoryError> {
452 let mut store = Self::acquire_lock(&self.store).await?;
453
454 if let Some(tx) = store.get_mut(&tx_id) {
455 if Self::is_final_state(&tx.status) {
456 return Ok(tx.clone());
457 }
458 let mut metadata = tx.metadata.clone().unwrap_or_default();
459 metadata.consecutive_failures = metadata.consecutive_failures.saturating_add(1);
460 metadata.total_failures = metadata.total_failures.saturating_add(1);
461 tx.metadata = Some(metadata);
462 Ok(tx.clone())
463 } else {
464 Err(RepositoryError::NotFound(format!(
465 "Transaction with ID {tx_id} not found"
466 )))
467 }
468 }
469
470 async fn reset_status_check_consecutive_failures(
471 &self,
472 tx_id: String,
473 ) -> Result<TransactionRepoModel, RepositoryError> {
474 let mut store = Self::acquire_lock(&self.store).await?;
475
476 if let Some(tx) = store.get_mut(&tx_id) {
477 if Self::is_final_state(&tx.status) {
478 return Ok(tx.clone());
479 }
480 let mut metadata = tx.metadata.clone().unwrap_or_default();
481 metadata.consecutive_failures = 0;
482 tx.metadata = Some(metadata);
483 Ok(tx.clone())
484 } else {
485 Err(RepositoryError::NotFound(format!(
486 "Transaction with ID {tx_id} not found"
487 )))
488 }
489 }
490
491 async fn record_stellar_insufficient_fee_retry(
492 &self,
493 tx_id: String,
494 sent_at: String,
495 ) -> Result<TransactionRepoModel, RepositoryError> {
496 let mut store = Self::acquire_lock(&self.store).await?;
497
498 if let Some(tx) = store.get_mut(&tx_id) {
499 if Self::is_final_state(&tx.status) {
500 return Ok(tx.clone());
501 }
502 let mut metadata = tx.metadata.clone().unwrap_or_default();
503 metadata.insufficient_fee_retries = metadata.insufficient_fee_retries.saturating_add(1);
504 tx.metadata = Some(metadata);
505 tx.sent_at = Some(sent_at);
506 Ok(tx.clone())
507 } else {
508 Err(RepositoryError::NotFound(format!(
509 "Transaction with ID {tx_id} not found"
510 )))
511 }
512 }
513
514 async fn record_stellar_try_again_later_retry(
515 &self,
516 tx_id: String,
517 sent_at: String,
518 ) -> Result<TransactionRepoModel, RepositoryError> {
519 let mut store = Self::acquire_lock(&self.store).await?;
520
521 if let Some(tx) = store.get_mut(&tx_id) {
522 if Self::is_final_state(&tx.status) {
523 return Ok(tx.clone());
524 }
525 let mut metadata = tx.metadata.clone().unwrap_or_default();
526 metadata.try_again_later_retries = metadata.try_again_later_retries.saturating_add(1);
527 tx.metadata = Some(metadata);
528 tx.sent_at = Some(sent_at);
529 Ok(tx.clone())
530 } else {
531 Err(RepositoryError::NotFound(format!(
532 "Transaction with ID {tx_id} not found"
533 )))
534 }
535 }
536
537 async fn set_confirmed_at(
538 &self,
539 tx_id: String,
540 confirmed_at: String,
541 ) -> Result<TransactionRepoModel, RepositoryError> {
542 let mut tx = self.get_by_id(tx_id.clone()).await?;
543 tx.confirmed_at = Some(confirmed_at);
544 self.update(tx_id, tx).await
545 }
546
547 async fn count_by_status(
548 &self,
549 relayer_id: &str,
550 statuses: &[TransactionStatus],
551 ) -> Result<u64, RepositoryError> {
552 let store = Self::acquire_lock(&self.store).await?;
553 let count = store
554 .values()
555 .filter(|tx| tx.relayer_id == relayer_id && statuses.contains(&tx.status))
556 .count() as u64;
557 Ok(count)
558 }
559
560 async fn delete_by_ids(&self, ids: Vec<String>) -> Result<BatchDeleteResult, RepositoryError> {
561 if ids.is_empty() {
562 return Ok(BatchDeleteResult::default());
563 }
564
565 let mut store = Self::acquire_lock(&self.store).await?;
566 let mut deleted_count = 0;
567 let mut failed = Vec::new();
568
569 for id in ids {
570 if store.remove(&id).is_some() {
571 deleted_count += 1;
572 } else {
573 failed.push((id.clone(), format!("Transaction with ID {id} not found")));
574 }
575 }
576
577 Ok(BatchDeleteResult {
578 deleted_count,
579 failed,
580 })
581 }
582
583 async fn delete_by_requests(
584 &self,
585 requests: Vec<TransactionDeleteRequest>,
586 ) -> Result<BatchDeleteResult, RepositoryError> {
587 if requests.is_empty() {
588 return Ok(BatchDeleteResult::default());
589 }
590
591 let ids: Vec<String> = requests.into_iter().map(|r| r.id).collect();
593 self.delete_by_ids(ids).await
594 }
595}
596
597impl Default for InMemoryTransactionRepository {
598 fn default() -> Self {
599 Self::new()
600 }
601}
602
603#[cfg(test)]
604mod tests {
605 use crate::models::{evm::Speed, EvmTransactionData, NetworkType};
606 use lazy_static::lazy_static;
607 use std::str::FromStr;
608
609 use crate::models::U256;
610
611 use super::*;
612
613 use tokio::sync::Mutex;
614
615 lazy_static! {
616 static ref ENV_MUTEX: Mutex<()> = Mutex::new(());
617 }
618 fn create_test_transaction(id: &str) -> TransactionRepoModel {
620 TransactionRepoModel {
621 id: id.to_string(),
622 relayer_id: "relayer-1".to_string(),
623 status: TransactionStatus::Pending,
624 status_reason: None,
625 created_at: "2025-01-27T15:31:10.777083+00:00".to_string(),
626 sent_at: Some("2025-01-27T15:31:10.777083+00:00".to_string()),
627 confirmed_at: Some("2025-01-27T15:31:10.777083+00:00".to_string()),
628 valid_until: None,
629 delete_at: None,
630 network_type: NetworkType::Evm,
631 priced_at: None,
632 hashes: vec![],
633 network_data: NetworkTransactionData::Evm(EvmTransactionData {
634 gas_price: Some(1000000000),
635 gas_limit: Some(21000),
636 nonce: Some(1),
637 value: U256::from_str("1000000000000000000").unwrap(),
638 data: Some("0x".to_string()),
639 from: "0xSender".to_string(),
640 to: Some("0xRecipient".to_string()),
641 chain_id: 1,
642 signature: None,
643 hash: Some(format!("0x{id}")),
644 speed: Some(Speed::Fast),
645 max_fee_per_gas: None,
646 max_priority_fee_per_gas: None,
647 raw: None,
648 }),
649 noop_count: None,
650 is_canceled: Some(false),
651 metadata: None,
652 }
653 }
654
655 fn create_test_transaction_pending_state(id: &str) -> TransactionRepoModel {
656 TransactionRepoModel {
657 id: id.to_string(),
658 relayer_id: "relayer-1".to_string(),
659 status: TransactionStatus::Pending,
660 status_reason: None,
661 created_at: "2025-01-27T15:31:10.777083+00:00".to_string(),
662 sent_at: None,
663 confirmed_at: None,
664 valid_until: None,
665 delete_at: None,
666 network_type: NetworkType::Evm,
667 priced_at: None,
668 hashes: vec![],
669 network_data: NetworkTransactionData::Evm(EvmTransactionData {
670 gas_price: Some(1000000000),
671 gas_limit: Some(21000),
672 nonce: Some(1),
673 value: U256::from_str("1000000000000000000").unwrap(),
674 data: Some("0x".to_string()),
675 from: "0xSender".to_string(),
676 to: Some("0xRecipient".to_string()),
677 chain_id: 1,
678 signature: None,
679 hash: Some(format!("0x{id}")),
680 speed: Some(Speed::Fast),
681 max_fee_per_gas: None,
682 max_priority_fee_per_gas: None,
683 raw: None,
684 }),
685 noop_count: None,
686 is_canceled: Some(false),
687 metadata: None,
688 }
689 }
690
691 #[tokio::test]
692 async fn test_create_transaction() {
693 let repo = InMemoryTransactionRepository::new();
694 let tx = create_test_transaction("test-1");
695
696 let result = repo.create(tx.clone()).await.unwrap();
697 assert_eq!(result.id, tx.id);
698 assert_eq!(repo.count().await.unwrap(), 1);
699 }
700
701 #[tokio::test]
702 async fn test_get_transaction() {
703 let repo = InMemoryTransactionRepository::new();
704 let tx = create_test_transaction("test-1");
705
706 repo.create(tx.clone()).await.unwrap();
707 let stored = repo.get_by_id("test-1".to_string()).await.unwrap();
708 if let NetworkTransactionData::Evm(stored_data) = &stored.network_data {
709 if let NetworkTransactionData::Evm(tx_data) = &tx.network_data {
710 assert_eq!(stored_data.hash, tx_data.hash);
711 }
712 }
713 }
714
715 #[tokio::test]
716 async fn test_update_transaction() {
717 let repo = InMemoryTransactionRepository::new();
718 let mut tx = create_test_transaction("test-1");
719
720 repo.create(tx.clone()).await.unwrap();
721 tx.status = TransactionStatus::Confirmed;
722
723 let updated = repo.update("test-1".to_string(), tx).await.unwrap();
724 assert!(matches!(updated.status, TransactionStatus::Confirmed));
725 }
726
727 #[tokio::test]
728 async fn test_delete_transaction() {
729 let repo = InMemoryTransactionRepository::new();
730 let tx = create_test_transaction("test-1");
731
732 repo.create(tx).await.unwrap();
733 repo.delete_by_id("test-1".to_string()).await.unwrap();
734
735 let result = repo.get_by_id("test-1".to_string()).await;
736 assert!(result.is_err());
737 }
738
739 #[tokio::test]
740 async fn test_list_all_transactions() {
741 let repo = InMemoryTransactionRepository::new();
742 let tx1 = create_test_transaction("test-1");
743 let tx2 = create_test_transaction("test-2");
744
745 repo.create(tx1).await.unwrap();
746 repo.create(tx2).await.unwrap();
747
748 let transactions = repo.list_all().await.unwrap();
749 assert_eq!(transactions.len(), 2);
750 }
751
752 #[tokio::test]
753 async fn test_count_transactions() {
754 let repo = InMemoryTransactionRepository::new();
755 let tx = create_test_transaction("test-1");
756
757 assert_eq!(repo.count().await.unwrap(), 0);
758 repo.create(tx).await.unwrap();
759 assert_eq!(repo.count().await.unwrap(), 1);
760 }
761
762 #[tokio::test]
763 async fn test_get_nonexistent_transaction() {
764 let repo = InMemoryTransactionRepository::new();
765 let result = repo.get_by_id("nonexistent".to_string()).await;
766 assert!(matches!(result, Err(RepositoryError::NotFound(_))));
767 }
768
769 #[tokio::test]
770 async fn test_duplicate_transaction_creation() {
771 let repo = InMemoryTransactionRepository::new();
772 let tx = create_test_transaction("test-1");
773
774 repo.create(tx.clone()).await.unwrap();
775 let result = repo.create(tx).await;
776
777 assert!(matches!(
778 result,
779 Err(RepositoryError::ConstraintViolation(_))
780 ));
781 }
782
783 #[tokio::test]
784 async fn test_update_nonexistent_transaction() {
785 let repo = InMemoryTransactionRepository::new();
786 let tx = create_test_transaction("test-1");
787
788 let result = repo.update("nonexistent".to_string(), tx).await;
789 assert!(matches!(result, Err(RepositoryError::NotFound(_))));
790 }
791
792 #[tokio::test]
793 async fn test_partial_update() {
794 let repo = InMemoryTransactionRepository::new();
795 let tx = create_test_transaction_pending_state("test-tx-id");
796 repo.create(tx.clone()).await.unwrap();
797
798 let update1 = TransactionUpdateRequest {
800 status: Some(TransactionStatus::Sent),
801 status_reason: None,
802 sent_at: None,
803 confirmed_at: None,
804 network_data: None,
805 hashes: None,
806 priced_at: None,
807 noop_count: None,
808 is_canceled: None,
809 delete_at: None,
810 metadata: None,
811 };
812 let updated_tx1 = repo
813 .partial_update("test-tx-id".to_string(), update1)
814 .await
815 .unwrap();
816 assert_eq!(updated_tx1.status, TransactionStatus::Sent);
817 assert_eq!(updated_tx1.sent_at, None);
818
819 let update2 = TransactionUpdateRequest {
821 status: Some(TransactionStatus::Confirmed),
822 status_reason: None,
823 sent_at: Some("2023-01-01T12:00:00Z".to_string()),
824 confirmed_at: Some("2023-01-01T12:05:00Z".to_string()),
825 network_data: None,
826 hashes: None,
827 priced_at: None,
828 noop_count: None,
829 is_canceled: None,
830 delete_at: None,
831 metadata: None,
832 };
833 let updated_tx2 = repo
834 .partial_update("test-tx-id".to_string(), update2)
835 .await
836 .unwrap();
837 assert_eq!(updated_tx2.status, TransactionStatus::Confirmed);
838 assert_eq!(
839 updated_tx2.sent_at,
840 Some("2023-01-01T12:00:00Z".to_string())
841 );
842 assert_eq!(
843 updated_tx2.confirmed_at,
844 Some("2023-01-01T12:05:00Z".to_string())
845 );
846
847 let update3 = TransactionUpdateRequest {
849 status: Some(TransactionStatus::Failed),
850 status_reason: None,
851 sent_at: None,
852 confirmed_at: None,
853 network_data: None,
854 hashes: None,
855 priced_at: None,
856 noop_count: None,
857 is_canceled: None,
858 delete_at: None,
859 metadata: None,
860 };
861 let result = repo
862 .partial_update("non-existent-id".to_string(), update3)
863 .await;
864 assert!(result.is_err());
865 assert!(matches!(result.unwrap_err(), RepositoryError::NotFound(_)));
866 }
867
868 #[tokio::test]
869 async fn test_update_status() {
870 let repo = InMemoryTransactionRepository::new();
871 let tx = create_test_transaction("test-1");
872
873 repo.create(tx).await.unwrap();
874
875 let updated = repo
877 .update_status("test-1".to_string(), TransactionStatus::Confirmed)
878 .await
879 .unwrap();
880
881 assert_eq!(updated.status, TransactionStatus::Confirmed);
883
884 let stored = repo.get_by_id("test-1".to_string()).await.unwrap();
886 assert_eq!(stored.status, TransactionStatus::Confirmed);
887
888 let updated = repo
890 .update_status("test-1".to_string(), TransactionStatus::Failed)
891 .await
892 .unwrap();
893
894 assert_eq!(updated.status, TransactionStatus::Failed);
896
897 let result = repo
899 .update_status("non-existent".to_string(), TransactionStatus::Confirmed)
900 .await;
901 assert!(matches!(result, Err(RepositoryError::NotFound(_))));
902 }
903
904 #[tokio::test]
905 async fn test_list_paginated() {
906 let repo = InMemoryTransactionRepository::new();
907
908 for i in 1..=10 {
910 let tx = create_test_transaction(&format!("test-{i}"));
911 repo.create(tx).await.unwrap();
912 }
913
914 let query = PaginationQuery {
916 page: 1,
917 per_page: 3,
918 };
919 let result = repo.list_paginated(query).await.unwrap();
920 assert_eq!(result.items.len(), 3);
921 assert_eq!(result.total, 10);
922 assert_eq!(result.page, 1);
923 assert_eq!(result.per_page, 3);
924
925 let query = PaginationQuery {
927 page: 2,
928 per_page: 3,
929 };
930 let result = repo.list_paginated(query).await.unwrap();
931 assert_eq!(result.items.len(), 3);
932 assert_eq!(result.total, 10);
933 assert_eq!(result.page, 2);
934 assert_eq!(result.per_page, 3);
935
936 let query = PaginationQuery {
938 page: 4,
939 per_page: 3,
940 };
941 let result = repo.list_paginated(query).await.unwrap();
942 assert_eq!(result.items.len(), 1);
943 assert_eq!(result.total, 10);
944 assert_eq!(result.page, 4);
945 assert_eq!(result.per_page, 3);
946
947 let query = PaginationQuery {
949 page: 5,
950 per_page: 3,
951 };
952 let result = repo.list_paginated(query).await.unwrap();
953 assert_eq!(result.items.len(), 0);
954 assert_eq!(result.total, 10);
955 }
956
957 #[tokio::test]
958 async fn test_find_by_nonce() {
959 let repo = InMemoryTransactionRepository::new();
960
961 let tx1 = create_test_transaction("test-1");
963
964 let mut tx2 = create_test_transaction("test-2");
965 if let NetworkTransactionData::Evm(ref mut data) = tx2.network_data {
966 data.nonce = Some(2);
967 }
968
969 let mut tx3 = create_test_transaction("test-3");
970 tx3.relayer_id = "relayer-2".to_string();
971 if let NetworkTransactionData::Evm(ref mut data) = tx3.network_data {
972 data.nonce = Some(1);
973 }
974
975 repo.create(tx1).await.unwrap();
976 repo.create(tx2).await.unwrap();
977 repo.create(tx3).await.unwrap();
978
979 let result = repo.find_by_nonce("relayer-1", 1).await.unwrap();
981 assert!(result.is_some());
982 assert_eq!(result.as_ref().unwrap().id, "test-1");
983
984 let result = repo.find_by_nonce("relayer-1", 2).await.unwrap();
986 assert!(result.is_some());
987 assert_eq!(result.as_ref().unwrap().id, "test-2");
988
989 let result = repo.find_by_nonce("relayer-2", 1).await.unwrap();
991 assert!(result.is_some());
992 assert_eq!(result.as_ref().unwrap().id, "test-3");
993
994 let result = repo.find_by_nonce("relayer-1", 99).await.unwrap();
996 assert!(result.is_none());
997 }
998
999 #[tokio::test]
1000 async fn test_get_nonce_occupancy_mixed_slots() {
1001 let repo = InMemoryTransactionRepository::new();
1002
1003 let tx1 = create_test_transaction("tx-1"); repo.create(tx1).await.unwrap();
1006
1007 let mut tx2 = create_test_transaction("tx-2");
1008 tx2.status = TransactionStatus::Failed;
1009 if let NetworkTransactionData::Evm(ref mut data) = tx2.network_data {
1010 data.nonce = Some(2);
1011 }
1012 repo.create(tx2).await.unwrap();
1013
1014 let result = repo.get_nonce_occupancy("relayer-1", 1, 4).await.unwrap();
1015
1016 assert_eq!(result.len(), 3);
1017 assert_eq!(result[0], (1, Some(TransactionStatus::Pending)));
1018 assert_eq!(result[1], (2, Some(TransactionStatus::Failed)));
1019 assert_eq!(result[2], (3, None));
1020 }
1021
1022 #[tokio::test]
1023 async fn test_get_nonce_occupancy_empty_range() {
1024 let repo = InMemoryTransactionRepository::new();
1025
1026 let result = repo.get_nonce_occupancy("relayer-1", 5, 5).await.unwrap();
1028 assert!(result.is_empty());
1029
1030 let result = repo.get_nonce_occupancy("relayer-1", 10, 5).await.unwrap();
1031 assert!(result.is_empty());
1032 }
1033
1034 #[tokio::test]
1035 async fn test_get_nonce_occupancy_wrong_relayer() {
1036 let repo = InMemoryTransactionRepository::new();
1037
1038 let tx1 = create_test_transaction("tx-1"); repo.create(tx1).await.unwrap();
1040
1041 let result = repo.get_nonce_occupancy("relayer-999", 1, 2).await.unwrap();
1043 assert_eq!(result, vec![(1, None)]);
1044 }
1045
1046 #[tokio::test]
1047 async fn test_update_network_data() {
1048 let repo = InMemoryTransactionRepository::new();
1049 let tx = create_test_transaction("test-1");
1050
1051 repo.create(tx.clone()).await.unwrap();
1052
1053 let updated_network_data = NetworkTransactionData::Evm(EvmTransactionData {
1055 gas_price: Some(2000000000),
1056 gas_limit: Some(30000),
1057 nonce: Some(2),
1058 value: U256::from_str("2000000000000000000").unwrap(),
1059 data: Some("0xUpdated".to_string()),
1060 from: "0xSender".to_string(),
1061 to: Some("0xRecipient".to_string()),
1062 chain_id: 1,
1063 signature: None,
1064 hash: Some("0xUpdated".to_string()),
1065 raw: None,
1066 speed: None,
1067 max_fee_per_gas: None,
1068 max_priority_fee_per_gas: None,
1069 });
1070
1071 let updated = repo
1072 .update_network_data("test-1".to_string(), updated_network_data)
1073 .await
1074 .unwrap();
1075
1076 if let NetworkTransactionData::Evm(data) = &updated.network_data {
1078 assert_eq!(data.gas_price, Some(2000000000));
1079 assert_eq!(data.gas_limit, Some(30000));
1080 assert_eq!(data.nonce, Some(2));
1081 assert_eq!(data.hash, Some("0xUpdated".to_string()));
1082 assert_eq!(data.data, Some("0xUpdated".to_string()));
1083 } else {
1084 panic!("Expected EVM network data");
1085 }
1086 }
1087
1088 #[tokio::test]
1089 async fn test_set_sent_at() {
1090 let repo = InMemoryTransactionRepository::new();
1091 let tx = create_test_transaction("test-1");
1092
1093 repo.create(tx).await.unwrap();
1094
1095 let new_sent_at = "2025-02-01T10:00:00.000000+00:00".to_string();
1097
1098 let updated = repo
1099 .set_sent_at("test-1".to_string(), new_sent_at.clone())
1100 .await
1101 .unwrap();
1102
1103 assert_eq!(updated.sent_at, Some(new_sent_at.clone()));
1105
1106 let stored = repo.get_by_id("test-1".to_string()).await.unwrap();
1108 assert_eq!(stored.sent_at, Some(new_sent_at.clone()));
1109 }
1110
1111 #[tokio::test]
1112 async fn test_set_confirmed_at() {
1113 let repo = InMemoryTransactionRepository::new();
1114 let tx = create_test_transaction("test-1");
1115
1116 repo.create(tx).await.unwrap();
1117
1118 let new_confirmed_at = "2025-02-01T11:30:45.123456+00:00".to_string();
1120
1121 let updated = repo
1122 .set_confirmed_at("test-1".to_string(), new_confirmed_at.clone())
1123 .await
1124 .unwrap();
1125
1126 assert_eq!(updated.confirmed_at, Some(new_confirmed_at.clone()));
1128
1129 let stored = repo.get_by_id("test-1".to_string()).await.unwrap();
1131 assert_eq!(stored.confirmed_at, Some(new_confirmed_at.clone()));
1132 }
1133
1134 #[tokio::test]
1135 async fn test_find_by_relayer_id() {
1136 let repo = InMemoryTransactionRepository::new();
1137 let tx1 = create_test_transaction("test-1");
1138 let tx2 = create_test_transaction("test-2");
1139
1140 let mut tx3 = create_test_transaction("test-3");
1142 tx3.relayer_id = "relayer-2".to_string();
1143
1144 repo.create(tx1).await.unwrap();
1145 repo.create(tx2).await.unwrap();
1146 repo.create(tx3).await.unwrap();
1147
1148 let query = PaginationQuery {
1150 page: 1,
1151 per_page: 10,
1152 };
1153 let result = repo
1154 .find_by_relayer_id("relayer-1", query.clone())
1155 .await
1156 .unwrap();
1157 assert_eq!(result.total, 2);
1158 assert_eq!(result.items.len(), 2);
1159 assert!(result.items.iter().all(|tx| tx.relayer_id == "relayer-1"));
1160
1161 let result = repo
1163 .find_by_relayer_id("relayer-2", query.clone())
1164 .await
1165 .unwrap();
1166 assert_eq!(result.total, 1);
1167 assert_eq!(result.items.len(), 1);
1168 assert!(result.items.iter().all(|tx| tx.relayer_id == "relayer-2"));
1169
1170 let result = repo
1172 .find_by_relayer_id("non-existent", query.clone())
1173 .await
1174 .unwrap();
1175 assert_eq!(result.total, 0);
1176 assert_eq!(result.items.len(), 0);
1177 }
1178
1179 #[tokio::test]
1180 async fn test_find_by_relayer_id_sorted_by_created_at_newest_first() {
1181 let repo = InMemoryTransactionRepository::new();
1182
1183 let mut tx1 = create_test_transaction("test-1");
1185 tx1.created_at = "2025-01-27T10:00:00.000000+00:00".to_string(); let mut tx2 = create_test_transaction("test-2");
1188 tx2.created_at = "2025-01-27T12:00:00.000000+00:00".to_string(); let mut tx3 = create_test_transaction("test-3");
1191 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 {
1199 page: 1,
1200 per_page: 10,
1201 };
1202 let result = repo.find_by_relayer_id("relayer-1", query).await.unwrap();
1203
1204 assert_eq!(result.total, 3);
1205 assert_eq!(result.items.len(), 3);
1206
1207 assert_eq!(
1209 result.items[0].id, "test-3",
1210 "First item should be newest (test-3)"
1211 );
1212 assert_eq!(
1213 result.items[0].created_at,
1214 "2025-01-27T14:00:00.000000+00:00"
1215 );
1216
1217 assert_eq!(
1218 result.items[1].id, "test-2",
1219 "Second item should be middle (test-2)"
1220 );
1221 assert_eq!(
1222 result.items[1].created_at,
1223 "2025-01-27T12:00:00.000000+00:00"
1224 );
1225
1226 assert_eq!(
1227 result.items[2].id, "test-1",
1228 "Third item should be oldest (test-1)"
1229 );
1230 assert_eq!(
1231 result.items[2].created_at,
1232 "2025-01-27T10:00:00.000000+00:00"
1233 );
1234 }
1235
1236 #[tokio::test]
1237 async fn test_find_by_status() {
1238 let repo = InMemoryTransactionRepository::new();
1239 let tx1 = create_test_transaction_pending_state("tx1");
1240 let mut tx2 = create_test_transaction_pending_state("tx2");
1241 tx2.status = TransactionStatus::Submitted;
1242 let mut tx3 = create_test_transaction_pending_state("tx3");
1243 tx3.relayer_id = "relayer-2".to_string();
1244 tx3.status = TransactionStatus::Pending;
1245
1246 repo.create(tx1.clone()).await.unwrap();
1247 repo.create(tx2.clone()).await.unwrap();
1248 repo.create(tx3.clone()).await.unwrap();
1249
1250 let pending_txs = repo
1252 .find_by_status("relayer-1", &[TransactionStatus::Pending])
1253 .await
1254 .unwrap();
1255 assert_eq!(pending_txs.len(), 1);
1256 assert_eq!(pending_txs[0].id, "tx1");
1257
1258 let submitted_txs = repo
1259 .find_by_status("relayer-1", &[TransactionStatus::Submitted])
1260 .await
1261 .unwrap();
1262 assert_eq!(submitted_txs.len(), 1);
1263 assert_eq!(submitted_txs[0].id, "tx2");
1264
1265 let multiple_status_txs = repo
1267 .find_by_status(
1268 "relayer-1",
1269 &[TransactionStatus::Pending, TransactionStatus::Submitted],
1270 )
1271 .await
1272 .unwrap();
1273 assert_eq!(multiple_status_txs.len(), 2);
1274
1275 let relayer2_pending = repo
1277 .find_by_status("relayer-2", &[TransactionStatus::Pending])
1278 .await
1279 .unwrap();
1280 assert_eq!(relayer2_pending.len(), 1);
1281 assert_eq!(relayer2_pending[0].id, "tx3");
1282
1283 let no_txs = repo
1285 .find_by_status("non-existent", &[TransactionStatus::Pending])
1286 .await
1287 .unwrap();
1288 assert_eq!(no_txs.len(), 0);
1289 }
1290
1291 #[tokio::test]
1292 async fn test_find_by_status_sorted_by_created_at() {
1293 let repo = InMemoryTransactionRepository::new();
1294
1295 let create_tx_with_timestamp = |id: &str, timestamp: &str| -> TransactionRepoModel {
1297 let mut tx = create_test_transaction_pending_state(id);
1298 tx.created_at = timestamp.to_string();
1299 tx.status = TransactionStatus::Pending;
1300 tx
1301 };
1302
1303 let tx3 = create_tx_with_timestamp("tx3", "2025-01-27T17:00:00.000000+00:00"); let tx1 = create_tx_with_timestamp("tx1", "2025-01-27T15:00:00.000000+00:00"); let tx2 = create_tx_with_timestamp("tx2", "2025-01-27T16:00:00.000000+00:00"); repo.create(tx3.clone()).await.unwrap();
1310 repo.create(tx1.clone()).await.unwrap();
1311 repo.create(tx2.clone()).await.unwrap();
1312
1313 let result = repo
1315 .find_by_status("relayer-1", &[TransactionStatus::Pending])
1316 .await
1317 .unwrap();
1318
1319 assert_eq!(result.len(), 3);
1321 assert_eq!(result[0].id, "tx3"); assert_eq!(result[1].id, "tx2"); assert_eq!(result[2].id, "tx1"); assert_eq!(result[0].created_at, "2025-01-27T17:00:00.000000+00:00");
1327 assert_eq!(result[1].created_at, "2025-01-27T16:00:00.000000+00:00");
1328 assert_eq!(result[2].created_at, "2025-01-27T15:00:00.000000+00:00");
1329 }
1330
1331 #[tokio::test]
1332 async fn test_find_by_status_paginated() {
1333 let repo = InMemoryTransactionRepository::new();
1334
1335 let create_tx_with_timestamp =
1337 |id: &str, timestamp: &str, status: TransactionStatus| -> TransactionRepoModel {
1338 let mut tx = create_test_transaction_pending_state(id);
1339 tx.created_at = timestamp.to_string();
1340 tx.status = status;
1341 tx
1342 };
1343
1344 for i in 1..=5 {
1346 let tx = create_tx_with_timestamp(
1347 &format!("tx{i}"),
1348 &format!("2025-01-27T{:02}:00:00.000000+00:00", 10 + i),
1349 TransactionStatus::Pending,
1350 );
1351 repo.create(tx).await.unwrap();
1352 }
1353
1354 for i in 6..=7 {
1356 let tx = create_tx_with_timestamp(
1357 &format!("tx{i}"),
1358 &format!("2025-01-27T{:02}:00:00.000000+00:00", 10 + i),
1359 TransactionStatus::Confirmed,
1360 );
1361 repo.create(tx).await.unwrap();
1362 }
1363
1364 let query = PaginationQuery {
1366 page: 1,
1367 per_page: 2,
1368 };
1369 let result = repo
1370 .find_by_status_paginated("relayer-1", &[TransactionStatus::Pending], query, false)
1371 .await
1372 .unwrap();
1373
1374 assert_eq!(result.total, 5);
1375 assert_eq!(result.items.len(), 2);
1376 assert_eq!(result.page, 1);
1377 assert_eq!(result.per_page, 2);
1378 assert_eq!(result.items[0].id, "tx5");
1380 assert_eq!(result.items[1].id, "tx4");
1381
1382 let query = PaginationQuery {
1384 page: 2,
1385 per_page: 2,
1386 };
1387 let result = repo
1388 .find_by_status_paginated("relayer-1", &[TransactionStatus::Pending], query, false)
1389 .await
1390 .unwrap();
1391
1392 assert_eq!(result.total, 5);
1393 assert_eq!(result.items.len(), 2);
1394 assert_eq!(result.page, 2);
1395 assert_eq!(result.items[0].id, "tx3");
1397 assert_eq!(result.items[1].id, "tx2");
1398
1399 let query = PaginationQuery {
1401 page: 3,
1402 per_page: 2,
1403 };
1404 let result = repo
1405 .find_by_status_paginated("relayer-1", &[TransactionStatus::Pending], query, false)
1406 .await
1407 .unwrap();
1408
1409 assert_eq!(result.total, 5);
1410 assert_eq!(result.items.len(), 1);
1411 assert_eq!(result.page, 3);
1412 assert_eq!(result.items[0].id, "tx1");
1413
1414 let query = PaginationQuery {
1416 page: 10,
1417 per_page: 2,
1418 };
1419 let result = repo
1420 .find_by_status_paginated("relayer-1", &[TransactionStatus::Pending], query, false)
1421 .await
1422 .unwrap();
1423
1424 assert_eq!(result.total, 5);
1425 assert_eq!(result.items.len(), 0);
1426
1427 let query = PaginationQuery {
1429 page: 1,
1430 per_page: 10,
1431 };
1432 let result = repo
1433 .find_by_status_paginated(
1434 "relayer-1",
1435 &[TransactionStatus::Pending, TransactionStatus::Confirmed],
1436 query,
1437 false,
1438 )
1439 .await
1440 .unwrap();
1441
1442 assert_eq!(result.total, 7);
1443 assert_eq!(result.items.len(), 7);
1444
1445 let query = PaginationQuery {
1447 page: 1,
1448 per_page: 10,
1449 };
1450 let result = repo
1451 .find_by_status_paginated("relayer-1", &[TransactionStatus::Failed], query, false)
1452 .await
1453 .unwrap();
1454
1455 assert_eq!(result.total, 0);
1456 assert_eq!(result.items.len(), 0);
1457 }
1458
1459 #[tokio::test]
1460 async fn test_find_by_status_paginated_filtered_excludes_canceled() {
1461 let repo = InMemoryTransactionRepository::new();
1462
1463 let make = |id: &str, ts: &str, is_canceled: Option<bool>| -> TransactionRepoModel {
1464 let mut tx = create_test_transaction_pending_state(id);
1465 tx.created_at = ts.to_string();
1466 tx.status = TransactionStatus::Submitted;
1467 tx.is_canceled = is_canceled;
1468 tx
1469 };
1470
1471 repo.create(make("tx1", "2025-01-27T11:00:00.000000+00:00", None))
1473 .await
1474 .unwrap();
1475 repo.create(make("tx2", "2025-01-27T12:00:00.000000+00:00", Some(true)))
1476 .await
1477 .unwrap();
1478 repo.create(make("tx3", "2025-01-27T13:00:00.000000+00:00", Some(false)))
1479 .await
1480 .unwrap();
1481 repo.create(make("tx4", "2025-01-27T14:00:00.000000+00:00", None))
1482 .await
1483 .unwrap();
1484
1485 let statuses = [TransactionStatus::Submitted];
1486 let query = PaginationQuery {
1487 page: 1,
1488 per_page: 10,
1489 };
1490
1491 let all = repo
1493 .find_by_status_paginated_filtered("relayer-1", &statuses, query.clone(), false, false)
1494 .await
1495 .unwrap();
1496 assert_eq!(all.total, 4);
1497 assert_eq!(all.items.len(), 4);
1498
1499 let filtered = repo
1502 .find_by_status_paginated_filtered("relayer-1", &statuses, query, false, true)
1503 .await
1504 .unwrap();
1505 assert_eq!(filtered.total, 3, "total must reflect the filtered set");
1506 assert_eq!(filtered.items.len(), 3);
1507 let ids: Vec<&str> = filtered.items.iter().map(|t| t.id.as_str()).collect();
1508 assert!(
1509 !ids.contains(&"tx2"),
1510 "cancelled tx must be excluded, got {ids:?}"
1511 );
1512
1513 let p1 = repo
1516 .find_by_status_paginated_filtered(
1517 "relayer-1",
1518 &statuses,
1519 PaginationQuery {
1520 page: 1,
1521 per_page: 2,
1522 },
1523 false,
1524 true,
1525 )
1526 .await
1527 .unwrap();
1528 assert_eq!(p1.total, 3);
1529 assert_eq!(p1.items.len(), 2);
1530 let p2 = repo
1531 .find_by_status_paginated_filtered(
1532 "relayer-1",
1533 &statuses,
1534 PaginationQuery {
1535 page: 2,
1536 per_page: 2,
1537 },
1538 false,
1539 true,
1540 )
1541 .await
1542 .unwrap();
1543 assert_eq!(p2.total, 3);
1544 assert_eq!(p2.items.len(), 1);
1545 }
1546
1547 #[tokio::test]
1548 async fn test_find_by_status_paginated_oldest_first() {
1549 let repo = InMemoryTransactionRepository::new();
1550
1551 let create_tx_with_timestamp =
1553 |id: &str, timestamp: &str, status: TransactionStatus| -> TransactionRepoModel {
1554 let mut tx = create_test_transaction_pending_state(id);
1555 tx.created_at = timestamp.to_string();
1556 tx.status = status;
1557 tx
1558 };
1559
1560 for i in 1..=5 {
1562 let tx = create_tx_with_timestamp(
1563 &format!("tx{i}"),
1564 &format!("2025-01-27T{:02}:00:00.000000+00:00", 10 + i),
1565 TransactionStatus::Pending,
1566 );
1567 repo.create(tx).await.unwrap();
1568 }
1569
1570 let query = PaginationQuery {
1572 page: 1,
1573 per_page: 3,
1574 };
1575 let result = repo
1576 .find_by_status_paginated("relayer-1", &[TransactionStatus::Pending], query, true)
1577 .await
1578 .unwrap();
1579
1580 assert_eq!(result.total, 5);
1581 assert_eq!(result.items.len(), 3);
1582 assert_eq!(
1584 result.items[0].id, "tx1",
1585 "First item should be oldest (tx1)"
1586 );
1587 assert_eq!(result.items[1].id, "tx2", "Second item should be tx2");
1588 assert_eq!(result.items[2].id, "tx3", "Third item should be tx3");
1589
1590 let query = PaginationQuery {
1592 page: 2,
1593 per_page: 3,
1594 };
1595 let result = repo
1596 .find_by_status_paginated("relayer-1", &[TransactionStatus::Pending], query, true)
1597 .await
1598 .unwrap();
1599
1600 assert_eq!(result.total, 5);
1601 assert_eq!(result.items.len(), 2);
1602 assert_eq!(result.items[0].id, "tx4");
1604 assert_eq!(result.items[1].id, "tx5");
1605 }
1606
1607 #[tokio::test]
1608 async fn test_find_by_status_paginated_oldest_first_single_item() {
1609 let repo = InMemoryTransactionRepository::new();
1610
1611 let timestamps = [
1613 ("tx-oldest", "2025-01-27T08:00:00.000000+00:00"),
1614 ("tx-middle", "2025-01-27T10:00:00.000000+00:00"),
1615 ("tx-newest", "2025-01-27T12:00:00.000000+00:00"),
1616 ];
1617
1618 for (id, timestamp) in timestamps {
1619 let mut tx = create_test_transaction_pending_state(id);
1620 tx.created_at = timestamp.to_string();
1621 tx.status = TransactionStatus::Pending;
1622 repo.create(tx).await.unwrap();
1623 }
1624
1625 let query = PaginationQuery {
1627 page: 1,
1628 per_page: 1,
1629 };
1630 let result = repo
1631 .find_by_status_paginated(
1632 "relayer-1",
1633 &[TransactionStatus::Pending],
1634 query.clone(),
1635 true,
1636 )
1637 .await
1638 .unwrap();
1639
1640 assert_eq!(result.total, 3);
1641 assert_eq!(result.items.len(), 1);
1642 assert_eq!(
1643 result.items[0].id, "tx-oldest",
1644 "With oldest_first and per_page=1, should return the oldest transaction"
1645 );
1646
1647 let result = repo
1649 .find_by_status_paginated("relayer-1", &[TransactionStatus::Pending], query, false)
1650 .await
1651 .unwrap();
1652
1653 assert_eq!(result.items.len(), 1);
1654 assert_eq!(
1655 result.items[0].id, "tx-newest",
1656 "With oldest_first=false and per_page=1, should return the newest transaction"
1657 );
1658 }
1659
1660 #[tokio::test]
1661 async fn test_find_by_status_paginated_multi_status_oldest_first() {
1662 let repo = InMemoryTransactionRepository::new();
1663
1664 let transactions = [
1666 (
1667 "tx-pending-old",
1668 "2025-01-27T08:00:00.000000+00:00",
1669 TransactionStatus::Pending,
1670 ),
1671 (
1672 "tx-sent-mid",
1673 "2025-01-27T10:00:00.000000+00:00",
1674 TransactionStatus::Sent,
1675 ),
1676 (
1677 "tx-pending-new",
1678 "2025-01-27T12:00:00.000000+00:00",
1679 TransactionStatus::Pending,
1680 ),
1681 (
1682 "tx-sent-old",
1683 "2025-01-27T07:00:00.000000+00:00",
1684 TransactionStatus::Sent,
1685 ),
1686 ];
1687
1688 for (id, timestamp, status) in transactions {
1689 let mut tx = create_test_transaction_pending_state(id);
1690 tx.created_at = timestamp.to_string();
1691 tx.status = status;
1692 repo.create(tx).await.unwrap();
1693 }
1694
1695 let query = PaginationQuery {
1697 page: 1,
1698 per_page: 10,
1699 };
1700 let result = repo
1701 .find_by_status_paginated(
1702 "relayer-1",
1703 &[TransactionStatus::Pending, TransactionStatus::Sent],
1704 query,
1705 true,
1706 )
1707 .await
1708 .unwrap();
1709
1710 assert_eq!(result.total, 4);
1711 assert_eq!(result.items.len(), 4);
1712 assert_eq!(result.items[0].id, "tx-sent-old", "Oldest should be first");
1714 assert_eq!(result.items[1].id, "tx-pending-old");
1715 assert_eq!(result.items[2].id, "tx-sent-mid");
1716 assert_eq!(
1717 result.items[3].id, "tx-pending-new",
1718 "Newest should be last"
1719 );
1720 }
1721
1722 #[tokio::test]
1723 async fn test_has_entries() {
1724 let repo = InMemoryTransactionRepository::new();
1725 assert!(!repo.has_entries().await.unwrap());
1726
1727 let tx = create_test_transaction("test");
1728 repo.create(tx.clone()).await.unwrap();
1729
1730 assert!(repo.has_entries().await.unwrap());
1731 }
1732
1733 #[tokio::test]
1734 async fn test_drop_all_entries() {
1735 let repo = InMemoryTransactionRepository::new();
1736 let tx = create_test_transaction("test");
1737 repo.create(tx.clone()).await.unwrap();
1738
1739 assert!(repo.has_entries().await.unwrap());
1740
1741 repo.drop_all_entries().await.unwrap();
1742 assert!(!repo.has_entries().await.unwrap());
1743 }
1744
1745 #[tokio::test]
1748 async fn test_update_status_sets_delete_at_for_final_statuses() {
1749 let _lock = ENV_MUTEX.lock().await;
1750
1751 use chrono::{DateTime, Duration, Utc};
1752 use std::env;
1753
1754 env::set_var("TRANSACTION_EXPIRATION_HOURS", "6");
1756
1757 let repo = InMemoryTransactionRepository::new();
1758
1759 let final_statuses = [
1760 TransactionStatus::Canceled,
1761 TransactionStatus::Confirmed,
1762 TransactionStatus::Failed,
1763 TransactionStatus::Expired,
1764 ];
1765
1766 for (i, status) in final_statuses.iter().enumerate() {
1767 let tx_id = format!("test-final-{i}");
1768 let tx = create_test_transaction_pending_state(&tx_id);
1769
1770 assert!(tx.delete_at.is_none());
1772
1773 repo.create(tx).await.unwrap();
1774
1775 let before_update = Utc::now();
1776
1777 let updated = repo
1779 .update_status(tx_id.clone(), status.clone())
1780 .await
1781 .unwrap();
1782
1783 assert!(
1785 updated.delete_at.is_some(),
1786 "delete_at should be set for status: {status:?}"
1787 );
1788
1789 let delete_at_str = updated.delete_at.unwrap();
1791 let delete_at = DateTime::parse_from_rfc3339(&delete_at_str)
1792 .expect("delete_at should be valid RFC3339")
1793 .with_timezone(&Utc);
1794
1795 let duration_from_before = delete_at.signed_duration_since(before_update);
1796 let expected_duration = Duration::hours(6);
1797 let tolerance = Duration::minutes(5);
1798
1799 assert!(
1800 duration_from_before >= expected_duration - tolerance &&
1801 duration_from_before <= expected_duration + tolerance,
1802 "delete_at should be approximately 6 hours from now for status: {status:?}. Duration: {duration_from_before:?}"
1803 );
1804 }
1805
1806 env::remove_var("TRANSACTION_EXPIRATION_HOURS");
1808 }
1809
1810 #[tokio::test]
1811 async fn test_update_status_does_not_set_delete_at_for_non_final_statuses() {
1812 let _lock = ENV_MUTEX.lock().await;
1813
1814 use std::env;
1815
1816 env::set_var("TRANSACTION_EXPIRATION_HOURS", "4");
1817
1818 let repo = InMemoryTransactionRepository::new();
1819
1820 let non_final_statuses = [
1821 TransactionStatus::Pending,
1822 TransactionStatus::Sent,
1823 TransactionStatus::Submitted,
1824 TransactionStatus::Mined,
1825 ];
1826
1827 for (i, status) in non_final_statuses.iter().enumerate() {
1828 let tx_id = format!("test-non-final-{i}");
1829 let tx = create_test_transaction_pending_state(&tx_id);
1830
1831 repo.create(tx).await.unwrap();
1832
1833 let updated = repo
1835 .update_status(tx_id.clone(), status.clone())
1836 .await
1837 .unwrap();
1838
1839 assert!(
1841 updated.delete_at.is_none(),
1842 "delete_at should NOT be set for status: {status:?}"
1843 );
1844 }
1845
1846 env::remove_var("TRANSACTION_EXPIRATION_HOURS");
1848 }
1849
1850 #[tokio::test]
1851 async fn test_partial_update_sets_delete_at_for_final_statuses() {
1852 let _lock = ENV_MUTEX.lock().await;
1853
1854 use chrono::{DateTime, Duration, Utc};
1855 use std::env;
1856
1857 env::set_var("TRANSACTION_EXPIRATION_HOURS", "8");
1858
1859 let repo = InMemoryTransactionRepository::new();
1860 let tx = create_test_transaction_pending_state("test-partial-final");
1861
1862 repo.create(tx).await.unwrap();
1863
1864 let before_update = Utc::now();
1865
1866 let update = TransactionUpdateRequest {
1868 status: Some(TransactionStatus::Confirmed),
1869 status_reason: Some("Transaction completed".to_string()),
1870 confirmed_at: Some("2023-01-01T12:05:00Z".to_string()),
1871 ..Default::default()
1872 };
1873
1874 let updated = repo
1875 .partial_update("test-partial-final".to_string(), update)
1876 .await
1877 .unwrap();
1878
1879 assert!(
1881 updated.delete_at.is_some(),
1882 "delete_at should be set when updating to Confirmed status"
1883 );
1884
1885 let delete_at_str = updated.delete_at.unwrap();
1887 let delete_at = DateTime::parse_from_rfc3339(&delete_at_str)
1888 .expect("delete_at should be valid RFC3339")
1889 .with_timezone(&Utc);
1890
1891 let duration_from_before = delete_at.signed_duration_since(before_update);
1892 let expected_duration = Duration::hours(8);
1893 let tolerance = Duration::minutes(5);
1894
1895 assert!(
1896 duration_from_before >= expected_duration - tolerance
1897 && duration_from_before <= expected_duration + tolerance,
1898 "delete_at should be approximately 8 hours from now. Duration: {duration_from_before:?}"
1899 );
1900
1901 assert_eq!(updated.status, TransactionStatus::Confirmed);
1903 assert_eq!(
1904 updated.status_reason,
1905 Some("Transaction completed".to_string())
1906 );
1907 assert_eq!(
1908 updated.confirmed_at,
1909 Some("2023-01-01T12:05:00Z".to_string())
1910 );
1911
1912 env::remove_var("TRANSACTION_EXPIRATION_HOURS");
1914 }
1915
1916 #[tokio::test]
1917 async fn test_update_status_preserves_existing_delete_at() {
1918 let _lock = ENV_MUTEX.lock().await;
1919
1920 use std::env;
1921
1922 env::set_var("TRANSACTION_EXPIRATION_HOURS", "2");
1923
1924 let repo = InMemoryTransactionRepository::new();
1925 let mut tx = create_test_transaction_pending_state("test-preserve-delete-at");
1926
1927 let existing_delete_at = "2025-01-01T12:00:00Z".to_string();
1929 tx.delete_at = Some(existing_delete_at.clone());
1930
1931 repo.create(tx).await.unwrap();
1932
1933 let updated = repo
1935 .update_status(
1936 "test-preserve-delete-at".to_string(),
1937 TransactionStatus::Confirmed,
1938 )
1939 .await
1940 .unwrap();
1941
1942 assert_eq!(
1944 updated.delete_at,
1945 Some(existing_delete_at),
1946 "Existing delete_at should be preserved when updating to final status"
1947 );
1948
1949 env::remove_var("TRANSACTION_EXPIRATION_HOURS");
1951 }
1952
1953 #[tokio::test]
1954 async fn test_partial_update_without_status_change_preserves_delete_at() {
1955 let _lock = ENV_MUTEX.lock().await;
1956
1957 use std::env;
1958
1959 env::set_var("TRANSACTION_EXPIRATION_HOURS", "3");
1960
1961 let repo = InMemoryTransactionRepository::new();
1962 let tx = create_test_transaction_pending_state("test-preserve-no-status");
1963
1964 repo.create(tx).await.unwrap();
1965
1966 let updated1 = repo
1968 .update_status(
1969 "test-preserve-no-status".to_string(),
1970 TransactionStatus::Confirmed,
1971 )
1972 .await
1973 .unwrap();
1974
1975 assert!(updated1.delete_at.is_some());
1976 let original_delete_at = updated1.delete_at.clone();
1977
1978 let update = TransactionUpdateRequest {
1980 status: None, status_reason: Some("Updated reason".to_string()),
1982 confirmed_at: Some("2023-01-01T12:10:00Z".to_string()),
1983 ..Default::default()
1984 };
1985
1986 let updated2 = repo
1987 .partial_update("test-preserve-no-status".to_string(), update)
1988 .await
1989 .unwrap();
1990
1991 assert_eq!(
1993 updated2.delete_at, original_delete_at,
1994 "delete_at should be preserved when status is not updated"
1995 );
1996
1997 assert_eq!(updated2.status, TransactionStatus::Confirmed); assert_eq!(updated2.status_reason, Some("Updated reason".to_string()));
2000 assert_eq!(
2001 updated2.confirmed_at,
2002 Some("2023-01-01T12:10:00Z".to_string())
2003 );
2004
2005 env::remove_var("TRANSACTION_EXPIRATION_HOURS");
2007 }
2008
2009 #[tokio::test]
2010 async fn test_update_status_multiple_updates_idempotent() {
2011 let _lock = ENV_MUTEX.lock().await;
2012
2013 use std::env;
2014
2015 env::set_var("TRANSACTION_EXPIRATION_HOURS", "12");
2016
2017 let repo = InMemoryTransactionRepository::new();
2018 let tx = create_test_transaction_pending_state("test-idempotent");
2019
2020 repo.create(tx).await.unwrap();
2021
2022 let updated1 = repo
2024 .update_status("test-idempotent".to_string(), TransactionStatus::Confirmed)
2025 .await
2026 .unwrap();
2027
2028 assert!(updated1.delete_at.is_some());
2029 let first_delete_at = updated1.delete_at.clone();
2030
2031 let updated2 = repo
2033 .update_status("test-idempotent".to_string(), TransactionStatus::Failed)
2034 .await
2035 .unwrap();
2036
2037 assert_eq!(
2039 updated2.delete_at, first_delete_at,
2040 "delete_at should not change on subsequent final status updates"
2041 );
2042
2043 assert_eq!(updated2.status, TransactionStatus::Failed);
2045
2046 env::remove_var("TRANSACTION_EXPIRATION_HOURS");
2048 }
2049
2050 #[tokio::test]
2053 async fn test_delete_by_ids_empty_list() {
2054 let repo = InMemoryTransactionRepository::new();
2055
2056 let tx = create_test_transaction("test-1");
2058 repo.create(tx).await.unwrap();
2059
2060 let result = repo.delete_by_ids(vec![]).await.unwrap();
2062
2063 assert_eq!(result.deleted_count, 0);
2064 assert!(result.failed.is_empty());
2065
2066 assert!(repo.get_by_id("test-1".to_string()).await.is_ok());
2068 }
2069
2070 #[tokio::test]
2071 async fn test_delete_by_ids_single_transaction() {
2072 let repo = InMemoryTransactionRepository::new();
2073
2074 let tx = create_test_transaction("test-1");
2075 repo.create(tx).await.unwrap();
2076
2077 let result = repo
2078 .delete_by_ids(vec!["test-1".to_string()])
2079 .await
2080 .unwrap();
2081
2082 assert_eq!(result.deleted_count, 1);
2083 assert!(result.failed.is_empty());
2084
2085 assert!(repo.get_by_id("test-1".to_string()).await.is_err());
2087 }
2088
2089 #[tokio::test]
2090 async fn test_delete_by_ids_multiple_transactions() {
2091 let repo = InMemoryTransactionRepository::new();
2092
2093 for i in 1..=5 {
2095 let tx = create_test_transaction(&format!("test-{i}"));
2096 repo.create(tx).await.unwrap();
2097 }
2098
2099 assert_eq!(repo.count().await.unwrap(), 5);
2100
2101 let ids_to_delete = vec![
2103 "test-1".to_string(),
2104 "test-3".to_string(),
2105 "test-5".to_string(),
2106 ];
2107 let result = repo.delete_by_ids(ids_to_delete).await.unwrap();
2108
2109 assert_eq!(result.deleted_count, 3);
2110 assert!(result.failed.is_empty());
2111
2112 assert!(repo.get_by_id("test-1".to_string()).await.is_err());
2114 assert!(repo.get_by_id("test-2".to_string()).await.is_ok()); assert!(repo.get_by_id("test-3".to_string()).await.is_err());
2116 assert!(repo.get_by_id("test-4".to_string()).await.is_ok()); assert!(repo.get_by_id("test-5".to_string()).await.is_err());
2118
2119 assert_eq!(repo.count().await.unwrap(), 2);
2120 }
2121
2122 #[tokio::test]
2123 async fn test_delete_by_ids_nonexistent_transactions() {
2124 let repo = InMemoryTransactionRepository::new();
2125
2126 let ids_to_delete = vec!["nonexistent-1".to_string(), "nonexistent-2".to_string()];
2128 let result = repo.delete_by_ids(ids_to_delete).await.unwrap();
2129
2130 assert_eq!(result.deleted_count, 0);
2131 assert_eq!(result.failed.len(), 2);
2132
2133 assert!(result.failed.iter().any(|(id, _)| id == "nonexistent-1"));
2135 assert!(result.failed.iter().any(|(id, _)| id == "nonexistent-2"));
2136 }
2137
2138 #[tokio::test]
2139 async fn test_delete_by_ids_mixed_existing_and_nonexistent() {
2140 let repo = InMemoryTransactionRepository::new();
2141
2142 for i in 1..=3 {
2144 let tx = create_test_transaction(&format!("test-{i}"));
2145 repo.create(tx).await.unwrap();
2146 }
2147
2148 let ids_to_delete = vec![
2150 "test-1".to_string(), "nonexistent-1".to_string(), "test-2".to_string(), "nonexistent-2".to_string(), ];
2155 let result = repo.delete_by_ids(ids_to_delete).await.unwrap();
2156
2157 assert_eq!(result.deleted_count, 2);
2158 assert_eq!(result.failed.len(), 2);
2159
2160 assert!(repo.get_by_id("test-1".to_string()).await.is_err());
2162 assert!(repo.get_by_id("test-2".to_string()).await.is_err());
2163
2164 assert!(repo.get_by_id("test-3".to_string()).await.is_ok());
2166
2167 let failed_ids: Vec<&String> = result.failed.iter().map(|(id, _)| id).collect();
2169 assert!(failed_ids.contains(&&"nonexistent-1".to_string()));
2170 assert!(failed_ids.contains(&&"nonexistent-2".to_string()));
2171 }
2172
2173 #[tokio::test]
2174 async fn test_delete_by_ids_all_transactions() {
2175 let repo = InMemoryTransactionRepository::new();
2176
2177 for i in 1..=10 {
2179 let tx = create_test_transaction(&format!("test-{i}"));
2180 repo.create(tx).await.unwrap();
2181 }
2182
2183 assert_eq!(repo.count().await.unwrap(), 10);
2184
2185 let ids_to_delete: Vec<String> = (1..=10).map(|i| format!("test-{i}")).collect();
2187 let result = repo.delete_by_ids(ids_to_delete).await.unwrap();
2188
2189 assert_eq!(result.deleted_count, 10);
2190 assert!(result.failed.is_empty());
2191 assert_eq!(repo.count().await.unwrap(), 0);
2192 assert!(!repo.has_entries().await.unwrap());
2193 }
2194
2195 #[tokio::test]
2196 async fn test_delete_by_ids_duplicate_ids() {
2197 let repo = InMemoryTransactionRepository::new();
2198
2199 let tx = create_test_transaction("test-1");
2200 repo.create(tx).await.unwrap();
2201
2202 let ids_to_delete = vec![
2204 "test-1".to_string(),
2205 "test-1".to_string(), "test-1".to_string(), ];
2208 let result = repo.delete_by_ids(ids_to_delete).await.unwrap();
2209
2210 assert_eq!(result.deleted_count, 1);
2212 assert_eq!(result.failed.len(), 2);
2213
2214 assert!(repo.get_by_id("test-1".to_string()).await.is_err());
2216 }
2217
2218 #[tokio::test]
2219 async fn test_delete_by_ids_preserves_other_relayer_transactions() {
2220 let repo = InMemoryTransactionRepository::new();
2221
2222 let mut tx1 = create_test_transaction("tx-relayer-1");
2224 tx1.relayer_id = "relayer-1".to_string();
2225
2226 let mut tx2 = create_test_transaction("tx-relayer-2");
2227 tx2.relayer_id = "relayer-2".to_string();
2228
2229 repo.create(tx1).await.unwrap();
2230 repo.create(tx2).await.unwrap();
2231
2232 let result = repo
2234 .delete_by_ids(vec!["tx-relayer-1".to_string()])
2235 .await
2236 .unwrap();
2237
2238 assert_eq!(result.deleted_count, 1);
2239
2240 let remaining = repo.get_by_id("tx-relayer-2".to_string()).await.unwrap();
2242 assert_eq!(remaining.relayer_id, "relayer-2");
2243 }
2244
2245 #[tokio::test]
2248 async fn test_increment_status_check_failures_no_prior_metadata() {
2249 let repo = InMemoryTransactionRepository::new();
2250 let tx = create_test_transaction_pending_state("tx-inc-1");
2251 repo.create(tx).await.unwrap();
2252
2253 let updated = repo
2254 .increment_status_check_failures("tx-inc-1".to_string())
2255 .await
2256 .unwrap();
2257
2258 let meta = updated.metadata.expect("metadata should be set");
2259 assert_eq!(meta.consecutive_failures, 1);
2260 assert_eq!(meta.total_failures, 1);
2261 assert_eq!(meta.insufficient_fee_retries, 0);
2262 }
2263
2264 #[tokio::test]
2265 async fn test_increment_status_check_failures_accumulates() {
2266 let repo = InMemoryTransactionRepository::new();
2267 let tx = create_test_transaction_pending_state("tx-inc-2");
2268 repo.create(tx).await.unwrap();
2269
2270 repo.increment_status_check_failures("tx-inc-2".to_string())
2271 .await
2272 .unwrap();
2273 repo.increment_status_check_failures("tx-inc-2".to_string())
2274 .await
2275 .unwrap();
2276 let updated = repo
2277 .increment_status_check_failures("tx-inc-2".to_string())
2278 .await
2279 .unwrap();
2280
2281 let meta = updated.metadata.unwrap();
2282 assert_eq!(meta.consecutive_failures, 3);
2283 assert_eq!(meta.total_failures, 3);
2284 }
2285
2286 #[tokio::test]
2287 async fn test_increment_status_check_failures_noop_on_final_state() {
2288 let repo = InMemoryTransactionRepository::new();
2289 let mut tx = create_test_transaction_pending_state("tx-inc-final");
2290 tx.status = TransactionStatus::Confirmed;
2291 repo.create(tx).await.unwrap();
2292
2293 let result = repo
2294 .increment_status_check_failures("tx-inc-final".to_string())
2295 .await
2296 .unwrap();
2297
2298 assert!(result.metadata.is_none());
2300 assert_eq!(result.status, TransactionStatus::Confirmed);
2301 }
2302
2303 #[tokio::test]
2304 async fn test_increment_status_check_failures_not_found() {
2305 let repo = InMemoryTransactionRepository::new();
2306 let result = repo
2307 .increment_status_check_failures("nonexistent".to_string())
2308 .await;
2309
2310 assert!(matches!(result, Err(RepositoryError::NotFound(_))));
2311 }
2312
2313 #[tokio::test]
2316 async fn test_reset_consecutive_failures() {
2317 let repo = InMemoryTransactionRepository::new();
2318 let tx = create_test_transaction_pending_state("tx-reset-1");
2319 repo.create(tx).await.unwrap();
2320
2321 repo.increment_status_check_failures("tx-reset-1".to_string())
2323 .await
2324 .unwrap();
2325 repo.increment_status_check_failures("tx-reset-1".to_string())
2326 .await
2327 .unwrap();
2328
2329 let updated = repo
2330 .reset_status_check_consecutive_failures("tx-reset-1".to_string())
2331 .await
2332 .unwrap();
2333
2334 let meta = updated.metadata.unwrap();
2335 assert_eq!(meta.consecutive_failures, 0);
2336 assert_eq!(meta.total_failures, 2);
2338 }
2339
2340 #[tokio::test]
2341 async fn test_reset_consecutive_failures_noop_on_final_state() {
2342 let repo = InMemoryTransactionRepository::new();
2343 let mut tx = create_test_transaction_pending_state("tx-reset-final");
2344 tx.status = TransactionStatus::Failed;
2345 tx.metadata = Some(crate::models::TransactionMetadata {
2346 consecutive_failures: 5,
2347 total_failures: 10,
2348 insufficient_fee_retries: 0,
2349 try_again_later_retries: 0,
2350 nonce_too_high_retries: 0,
2351 });
2352 repo.create(tx).await.unwrap();
2353
2354 let result = repo
2355 .reset_status_check_consecutive_failures("tx-reset-final".to_string())
2356 .await
2357 .unwrap();
2358
2359 let meta = result.metadata.unwrap();
2361 assert_eq!(meta.consecutive_failures, 5);
2362 }
2363
2364 #[tokio::test]
2365 async fn test_reset_consecutive_failures_not_found() {
2366 let repo = InMemoryTransactionRepository::new();
2367 let result = repo
2368 .reset_status_check_consecutive_failures("nonexistent".to_string())
2369 .await;
2370
2371 assert!(matches!(result, Err(RepositoryError::NotFound(_))));
2372 }
2373
2374 #[tokio::test]
2377 async fn test_record_insufficient_fee_retry() {
2378 let repo = InMemoryTransactionRepository::new();
2379 let mut tx = create_test_transaction_pending_state("tx-fee-1");
2380 tx.status = TransactionStatus::Sent;
2381 tx.sent_at = None;
2382 repo.create(tx).await.unwrap();
2383
2384 let updated = repo
2385 .record_stellar_insufficient_fee_retry(
2386 "tx-fee-1".to_string(),
2387 "2025-03-18T10:00:00Z".to_string(),
2388 )
2389 .await
2390 .unwrap();
2391
2392 assert_eq!(updated.sent_at.as_deref(), Some("2025-03-18T10:00:00Z"));
2393 let meta = updated.metadata.unwrap();
2394 assert_eq!(meta.insufficient_fee_retries, 1);
2395 assert_eq!(meta.consecutive_failures, 0);
2396 assert_eq!(meta.total_failures, 0);
2397 }
2398
2399 #[tokio::test]
2400 async fn test_record_insufficient_fee_retry_accumulates() {
2401 let repo = InMemoryTransactionRepository::new();
2402 let mut tx = create_test_transaction_pending_state("tx-fee-2");
2403 tx.status = TransactionStatus::Sent;
2404 repo.create(tx).await.unwrap();
2405
2406 repo.record_stellar_insufficient_fee_retry(
2407 "tx-fee-2".to_string(),
2408 "2025-03-18T10:00:00Z".to_string(),
2409 )
2410 .await
2411 .unwrap();
2412
2413 let updated = repo
2414 .record_stellar_insufficient_fee_retry(
2415 "tx-fee-2".to_string(),
2416 "2025-03-18T10:01:00Z".to_string(),
2417 )
2418 .await
2419 .unwrap();
2420
2421 assert_eq!(updated.sent_at.as_deref(), Some("2025-03-18T10:01:00Z"));
2422 let meta = updated.metadata.unwrap();
2423 assert_eq!(meta.insufficient_fee_retries, 2);
2424 }
2425
2426 #[tokio::test]
2427 async fn test_record_insufficient_fee_retry_noop_on_final_state() {
2428 let repo = InMemoryTransactionRepository::new();
2429 let mut tx = create_test_transaction_pending_state("tx-fee-final");
2430 tx.status = TransactionStatus::Confirmed;
2431 tx.sent_at = Some("old-time".to_string());
2432 repo.create(tx).await.unwrap();
2433
2434 let result = repo
2435 .record_stellar_insufficient_fee_retry(
2436 "tx-fee-final".to_string(),
2437 "new-time".to_string(),
2438 )
2439 .await
2440 .unwrap();
2441
2442 assert_eq!(result.sent_at.as_deref(), Some("old-time"));
2444 assert!(result.metadata.is_none());
2445 }
2446
2447 #[tokio::test]
2448 async fn test_record_insufficient_fee_retry_not_found() {
2449 let repo = InMemoryTransactionRepository::new();
2450 let result = repo
2451 .record_stellar_insufficient_fee_retry(
2452 "nonexistent".to_string(),
2453 "2025-03-18T10:00:00Z".to_string(),
2454 )
2455 .await;
2456
2457 assert!(matches!(result, Err(RepositoryError::NotFound(_))));
2458 }
2459
2460 #[tokio::test]
2463 async fn test_record_try_again_later_retry() {
2464 let repo = InMemoryTransactionRepository::new();
2465 let mut tx = create_test_transaction_pending_state("tx-tal-1");
2466 tx.status = TransactionStatus::Sent;
2467 tx.sent_at = None;
2468 repo.create(tx).await.unwrap();
2469
2470 let updated = repo
2471 .record_stellar_try_again_later_retry(
2472 "tx-tal-1".to_string(),
2473 "2025-03-18T10:00:00Z".to_string(),
2474 )
2475 .await
2476 .unwrap();
2477
2478 assert_eq!(updated.sent_at.as_deref(), Some("2025-03-18T10:00:00Z"));
2479 let meta = updated.metadata.unwrap();
2480 assert_eq!(meta.try_again_later_retries, 1);
2481 assert_eq!(meta.consecutive_failures, 0);
2482 assert_eq!(meta.total_failures, 0);
2483 }
2484
2485 #[tokio::test]
2486 async fn test_record_try_again_later_retry_accumulates() {
2487 let repo = InMemoryTransactionRepository::new();
2488 let mut tx = create_test_transaction_pending_state("tx-tal-2");
2489 tx.status = TransactionStatus::Sent;
2490 repo.create(tx).await.unwrap();
2491
2492 repo.record_stellar_try_again_later_retry(
2493 "tx-tal-2".to_string(),
2494 "2025-03-18T10:00:00Z".to_string(),
2495 )
2496 .await
2497 .unwrap();
2498
2499 let updated = repo
2500 .record_stellar_try_again_later_retry(
2501 "tx-tal-2".to_string(),
2502 "2025-03-18T10:01:00Z".to_string(),
2503 )
2504 .await
2505 .unwrap();
2506
2507 assert_eq!(updated.sent_at.as_deref(), Some("2025-03-18T10:01:00Z"));
2508 let meta = updated.metadata.unwrap();
2509 assert_eq!(meta.try_again_later_retries, 2);
2510 }
2511
2512 #[tokio::test]
2513 async fn test_record_try_again_later_retry_noop_on_final_state() {
2514 let repo = InMemoryTransactionRepository::new();
2515 let mut tx = create_test_transaction_pending_state("tx-tal-final");
2516 tx.status = TransactionStatus::Confirmed;
2517 tx.sent_at = Some("old-time".to_string());
2518 repo.create(tx).await.unwrap();
2519
2520 let result = repo
2521 .record_stellar_try_again_later_retry(
2522 "tx-tal-final".to_string(),
2523 "new-time".to_string(),
2524 )
2525 .await
2526 .unwrap();
2527
2528 assert_eq!(result.sent_at.as_deref(), Some("old-time"));
2530 assert!(result.metadata.is_none());
2531 }
2532
2533 #[tokio::test]
2534 async fn test_record_try_again_later_retry_not_found() {
2535 let repo = InMemoryTransactionRepository::new();
2536 let result = repo
2537 .record_stellar_try_again_later_retry(
2538 "nonexistent".to_string(),
2539 "2025-03-18T10:00:00Z".to_string(),
2540 )
2541 .await;
2542
2543 assert!(matches!(result, Err(RepositoryError::NotFound(_))));
2544 }
2545}