moirai-runtime 0.7.0

High-performance hybrid concurrency library for Rust - weaving the threads of fate
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
//! Financial Transaction Processing - Real-World Concurrency Example
//!
//! This example demonstrates:
//! - Race condition handling in financial transfers
//! - Transaction isolation and consistency
//! - Deadlock prevention in multi-account operations
//! - Audit trail maintenance under high concurrency
//! - Error handling and recovery in financial systems

#![expect(
    clippy::unwrap_used,
    reason = "test scope: failed precondition = test failure"
)]
#![expect(
    dead_code,
    reason = "This example defines audit and transaction domain variants beyond the short executable scenario"
)]

use moirai::{Moirai, Priority};
use std::collections::HashMap;
use std::fmt;
use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex, RwLock};
use std::time::{Instant, SystemTime, UNIX_EPOCH};

/// Represents an account in the financial system
#[derive(Debug, Clone)]
struct Account {
    id: u64,
    balance: Arc<AtomicU64>, // Using atomic for lock-free balance operations
    version: Arc<AtomicU64>, // Optimistic locking version
}

impl Account {
    fn new(id: u64, initial_balance: u64) -> Self {
        Self {
            id,
            balance: Arc::new(AtomicU64::new(initial_balance)),
            version: Arc::new(AtomicU64::new(1)),
        }
    }

    fn balance(&self) -> u64 {
        self.balance.load(Ordering::Acquire)
    }

    fn version(&self) -> u64 {
        self.version.load(Ordering::Acquire)
    }
}

/// Different types of financial transactions
#[derive(Debug, Clone, PartialEq)]
enum TransactionType {
    Transfer,
    Deposit,
    Withdrawal,
    Fee,
}

/// Represents a financial transaction
#[derive(Debug, Clone)]
struct Transaction {
    id: u64,
    from_account: Option<u64>,
    to_account: Option<u64>,
    amount: u64,
    transaction_type: TransactionType,
    timestamp: u64,
    status: TransactionStatus,
}

#[derive(Debug, Clone, PartialEq)]
enum TransactionStatus {
    Pending,
    Completed,
    Failed,
    Cancelled,
}

impl fmt::Display for TransactionStatus {
    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
        match self {
            TransactionStatus::Pending => write!(f, "PENDING"),
            TransactionStatus::Completed => write!(f, "COMPLETED"),
            TransactionStatus::Failed => write!(f, "FAILED"),
            TransactionStatus::Cancelled => write!(f, "CANCELLED"),
        }
    }
}

/// Audit trail for financial transactions
#[derive(Debug)]
struct AuditTrail {
    entries: Arc<Mutex<Vec<AuditEntry>>>,
    entry_count: Arc<AtomicUsize>,
}

#[derive(Debug, Clone)]
struct AuditEntry {
    transaction_id: u64,
    account_id: u64,
    balance_before: u64,
    balance_after: u64,
    timestamp: u64,
    operation: String,
}

impl AuditTrail {
    fn new() -> Self {
        Self {
            entries: Arc::new(Mutex::new(Vec::new())),
            entry_count: Arc::new(AtomicUsize::new(0)),
        }
    }

    fn log_entry(&self, entry: AuditEntry) -> Result<(), String> {
        let mut entries = self
            .entries
            .lock()
            .map_err(|_| "Failed to acquire audit lock")?;
        entries.push(entry);
        self.entry_count.fetch_add(1, Ordering::Relaxed);
        Ok(())
    }

    fn entry_count(&self) -> usize {
        self.entry_count.load(Ordering::Relaxed)
    }

    fn get_entries_for_account(&self, account_id: u64) -> Result<Vec<AuditEntry>, String> {
        let entries = self
            .entries
            .lock()
            .map_err(|_| "Failed to acquire audit lock")?;
        Ok(entries
            .iter()
            .filter(|entry| entry.account_id == account_id)
            .cloned()
            .collect())
    }
}

/// Financial transaction processing engine
struct TransactionEngine {
    accounts: Arc<RwLock<HashMap<u64, Account>>>,
    audit_trail: Arc<AuditTrail>,
    transaction_counter: Arc<AtomicU64>,
    successful_transactions: Arc<AtomicUsize>,
    failed_transactions: Arc<AtomicUsize>,
    runtime: Moirai,
}

impl TransactionEngine {
    fn new() -> Result<Self, String> {
        let runtime = Moirai::new().map_err(|_| "Failed to create Moirai runtime")?;

        Ok(Self {
            accounts: Arc::new(RwLock::new(HashMap::new())),
            audit_trail: Arc::new(AuditTrail::new()),
            transaction_counter: Arc::new(AtomicU64::new(1)),
            successful_transactions: Arc::new(AtomicUsize::new(0)),
            failed_transactions: Arc::new(AtomicUsize::new(0)),
            runtime,
        })
    }

    fn create_account(&self, account_id: u64, initial_balance: u64) -> Result<(), String> {
        let mut accounts = self
            .accounts
            .write()
            .map_err(|_| "Failed to acquire accounts write lock")?;

        if accounts.contains_key(&account_id) {
            return Err(format!("Account {} already exists", account_id));
        }

        accounts.insert(account_id, Account::new(account_id, initial_balance));
        Ok(())
    }

    fn get_account_balance(&self, account_id: u64) -> Result<u64, String> {
        let accounts = self
            .accounts
            .read()
            .map_err(|_| "Failed to acquire accounts read lock")?;

        accounts
            .get(&account_id)
            .map(|account| account.balance())
            .ok_or_else(|| format!("Account {} not found", account_id))
    }

    /// Process a transfer with comprehensive error handling and audit trailing
    fn process_transfer(
        &self,
        from_account_id: u64,
        to_account_id: u64,
        amount: u64,
        priority: Priority,
    ) -> Result<u64, String> {
        if from_account_id == to_account_id {
            return Err("Cannot transfer to the same account".to_string());
        }

        if amount == 0 {
            return Err("Transfer amount must be greater than zero".to_string());
        }

        let transaction_id = self.transaction_counter.fetch_add(1, Ordering::Relaxed);
        let timestamp = SystemTime::now()
            .duration_since(UNIX_EPOCH)
            .unwrap()
            .as_secs();

        let accounts = self.accounts.clone();
        let audit_trail = self.audit_trail.clone();
        let successful_counter = self.successful_transactions.clone();
        let failed_counter = self.failed_transactions.clone();

        // Use async execution for the transaction to demonstrate real-world async processing
        let handle = self.runtime.spawn_fn_with_priority(
            move || {
                // Always acquire locks in consistent order (by account ID) to prevent deadlocks
                let (_first_id, _second_id) = if from_account_id < to_account_id {
                    (from_account_id, to_account_id)
                } else {
                    (to_account_id, from_account_id)
                };

                // Get account references
                let accounts_read = accounts.read().map_err(|_| {
                    format!(
                        "Failed to acquire accounts lock for transaction {}",
                        transaction_id
                    )
                })?;

                let from_account = accounts_read
                    .get(&from_account_id)
                    .ok_or_else(|| format!("Source account {} not found", from_account_id))?
                    .clone();

                let to_account = accounts_read
                    .get(&to_account_id)
                    .ok_or_else(|| format!("Destination account {} not found", to_account_id))?
                    .clone();

                // Release read lock before processing
                drop(accounts_read);

                // Check balance with optimistic locking
                let from_balance_before = from_account.balance();
                let _from_version_before = from_account.version();

                if from_balance_before < amount {
                    failed_counter.fetch_add(1, Ordering::Relaxed);
                    return Err(format!(
                        "Insufficient funds in account {}: {} < {}",
                        from_account_id, from_balance_before, amount
                    ));
                }

                // Perform atomic balance updates
                let to_balance_before = to_account.balance();

                // Debit from source account
                let from_balance_after = from_account.balance.fetch_sub(amount, Ordering::AcqRel);
                if from_balance_after < amount {
                    // Race condition detected - restore balance and fail
                    from_account.balance.fetch_add(amount, Ordering::AcqRel);
                    failed_counter.fetch_add(1, Ordering::Relaxed);
                    return Err(format!(
                        "Race condition detected in account {}",
                        from_account_id
                    ));
                }
                let from_balance_after = from_balance_after - amount;

                // Credit to destination account
                let to_balance_after =
                    to_account.balance.fetch_add(amount, Ordering::AcqRel) + amount;

                // Update versions for optimistic locking
                from_account.version.fetch_add(1, Ordering::AcqRel);
                to_account.version.fetch_add(1, Ordering::AcqRel);

                // Log audit entries
                let from_audit = AuditEntry {
                    transaction_id,
                    account_id: from_account_id,
                    balance_before: from_balance_before,
                    balance_after: from_balance_after,
                    timestamp,
                    operation: format!("DEBIT {}", amount),
                };

                let to_audit = AuditEntry {
                    transaction_id,
                    account_id: to_account_id,
                    balance_before: to_balance_before,
                    balance_after: to_balance_after,
                    timestamp,
                    operation: format!("CREDIT {}", amount),
                };

                audit_trail
                    .log_entry(from_audit)
                    .map_err(|e| format!("Failed to log debit audit: {}", e))?;
                audit_trail
                    .log_entry(to_audit)
                    .map_err(|e| format!("Failed to log credit audit: {}", e))?;

                successful_counter.fetch_add(1, Ordering::Relaxed);
                Ok(transaction_id)
            },
            priority,
        );

        // Return transaction ID immediately for async processing
        match handle.join() {
            Some(Ok(result)) => result,
            Some(Err(_)) | None => {
                self.failed_transactions.fetch_add(1, Ordering::Relaxed);
                Err("Transaction execution failed".to_string())
            }
        }
    }

    /// Process multiple transactions concurrently to test edge cases
    fn process_batch_transfers(
        &self,
        transfers: Vec<(u64, u64, u64)>,
    ) -> Result<Vec<Result<u64, String>>, String> {
        let mut handles = Vec::new();

        for (from, to, amount) in transfers {
            let engine = self;
            let priority = if amount > 1000 {
                Priority::High
            } else {
                Priority::Normal
            };

            // Process each transfer with appropriate priority
            match engine.process_transfer(from, to, amount, priority) {
                Ok(tx_id) => handles.push(Ok(tx_id)),
                Err(e) => handles.push(Err(e)),
            }
        }

        Ok(handles)
    }

    fn get_statistics(&self) -> (usize, usize, usize) {
        (
            self.successful_transactions.load(Ordering::Relaxed),
            self.failed_transactions.load(Ordering::Relaxed),
            self.audit_trail.entry_count(),
        )
    }
}

fn main() -> Result<(), Box<dyn std::error::Error>> {
    println!("Financial Transaction Processing - Real-World Concurrency");
    println!("=========================================================");

    let engine = TransactionEngine::new()?;

    // Create test accounts
    println!("\n1. Setting up test accounts...");
    engine.create_account(1001, 10000)?; // Alice
    engine.create_account(1002, 5000)?; // Bob
    engine.create_account(1003, 7500)?; // Carol
    engine.create_account(1004, 2000)?; // Dave
    engine.create_account(1005, 15000)?; // Eve

    println!("  Created 5 accounts with initial balances");

    // Edge Case 1: High-frequency concurrent transfers
    println!("\n2. High-frequency concurrent transfers...");
    let start_time = Instant::now();

    let concurrent_transfers = vec![
        (1001, 1002, 100), // Alice -> Bob
        (1002, 1003, 50),  // Bob -> Carol
        (1003, 1004, 75),  // Carol -> Dave
        (1004, 1005, 25),  // Dave -> Eve
        (1005, 1001, 200), // Eve -> Alice
        (1001, 1003, 150), // Alice -> Carol
        (1002, 1004, 80),  // Bob -> Dave
        (1003, 1005, 120), // Carol -> Eve
        (1004, 1001, 90),  // Dave -> Alice
        (1005, 1002, 110), // Eve -> Bob
    ];

    let results = engine.process_batch_transfers(concurrent_transfers)?;
    let processing_time = start_time.elapsed();

    println!(
        "  Processed {} transfers in {:?}",
        results.len(),
        processing_time
    );

    let successful_count = results.iter().filter(|r| r.is_ok()).count();
    let failed_count = results.len() - successful_count;
    println!(
        "  Results: {} successful, {} failed",
        successful_count, failed_count
    );

    // Edge Case 2: Insufficient funds scenario
    println!("\n3. Testing insufficient funds edge case...");
    match engine.process_transfer(1004, 1005, 5000, Priority::High) {
        Ok(_) => println!("  ERROR: Transfer should have failed!"),
        Err(e) => println!("  Expected failure: {}", e),
    }

    // Edge Case 3: Same account transfer
    println!("\n4. Testing same account transfer edge case...");
    match engine.process_transfer(1001, 1001, 100, Priority::Normal) {
        Ok(_) => println!("  ERROR: Same account transfer should have failed!"),
        Err(e) => println!("  Expected failure: {}", e),
    }

    // Edge Case 4: Zero amount transfer
    println!("\n5. Testing zero amount transfer edge case...");
    match engine.process_transfer(1001, 1002, 0, Priority::Normal) {
        Ok(_) => println!("  ERROR: Zero amount transfer should have failed!"),
        Err(e) => println!("  Expected failure: {}", e),
    }

    // Edge Case 5: Race condition simulation
    println!("\n6. Simulating race conditions with rapid transfers...");
    let race_start = Instant::now();

    // Create multiple rapid transfers from the same account
    let rapid_transfers = vec![
        (1001, 1002, 1000),
        (1001, 1003, 1000),
        (1001, 1004, 1000),
        (1001, 1005, 1000),
        (1001, 1002, 1000),
        (1001, 1003, 1000),
        (1001, 1004, 1000),
        (1001, 1005, 1000),
    ];

    let race_results = engine.process_batch_transfers(rapid_transfers)?;
    let race_time = race_start.elapsed();

    let race_successful = race_results.iter().filter(|r| r.is_ok()).count();
    let race_failed = race_results.len() - race_successful;

    println!(
        "  Rapid transfers: {} successful, {} failed in {:?}",
        race_successful, race_failed, race_time
    );

    // Display final account balances
    println!("\n7. Final account balances:");
    for account_id in 1001..=1005 {
        match engine.get_account_balance(account_id) {
            Ok(balance) => println!("  Account {}: ${}", account_id, balance),
            Err(e) => println!("  Account {}: Error - {}", account_id, e),
        }
    }

    // Display audit statistics
    let (successful, failed, audit_entries) = engine.get_statistics();
    println!("\n8. Transaction Statistics:");
    println!("  Successful transactions: {}", successful);
    println!("  Failed transactions: {}", failed);
    println!("  Audit entries: {}", audit_entries);
    println!(
        "  Success rate: {:.2}%",
        (successful as f64 / (successful + failed) as f64) * 100.0
    );

    // Edge Case 6: Audit trail verification
    println!("\n9. Audit trail verification for Account 1001:");
    match engine.audit_trail.get_entries_for_account(1001) {
        Ok(entries) => {
            println!("  Found {} audit entries:", entries.len());
            for (i, entry) in entries.iter().take(5).enumerate() {
                println!(
                    "    {}. TX{}: {} (${} -> ${})",
                    i + 1,
                    entry.transaction_id,
                    entry.operation,
                    entry.balance_before,
                    entry.balance_after
                );
            }
            if entries.len() > 5 {
                println!("    ... and {} more entries", entries.len() - 5);
            }
        }
        Err(e) => println!("  Failed to get audit entries: {}", e),
    }

    println!("\nFinancial transaction processing completed successfully!");
    println!("All edge cases handled appropriately with proper error recovery.");

    Ok(())
}