use chrono::{DateTime, Utc};
use rust_decimal::Decimal;
use rust_decimal_macros::dec;
use serde::{Deserialize, Serialize};
use std::collections::{HashMap, HashSet};
use uuid::Uuid;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub enum Region {
US,
EU,
UK,
APAC,
ME,
LATAM,
Other,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub enum RestrictionLevel {
None,
Partial,
Full,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum SuspiciousActivity {
UnusualVolume,
RapidFunding,
WashTrading,
Layering,
Spoofing,
FrontRunning,
UnusualWithdrawal,
FailedTransactions,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SuspiciousActivityAlert {
pub id: Uuid,
pub user_id: Uuid,
pub activity_type: SuspiciousActivity,
pub description: String,
pub risk_score: u8,
pub transaction_ids: Vec<Uuid>,
pub detected_at: DateTime<Utc>,
pub reviewed: bool,
pub notes: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct RegionalCompliance {
pub region: Region,
pub trading_allowed: bool,
pub daily_limit: Option<Decimal>,
pub monthly_limit: Option<Decimal>,
pub kyc_required: bool,
pub min_age: u8,
pub restricted_tokens: HashSet<Uuid>,
}
impl RegionalCompliance {
pub fn us_standard() -> Self {
Self {
region: Region::US,
trading_allowed: true,
daily_limit: Some(dec!(50000)),
monthly_limit: Some(dec!(500000)),
kyc_required: true,
min_age: 18,
restricted_tokens: HashSet::new(),
}
}
pub fn eu_standard() -> Self {
Self {
region: Region::EU,
trading_allowed: true,
daily_limit: Some(dec!(40000)),
monthly_limit: Some(dec!(400000)),
kyc_required: true,
min_age: 18,
restricted_tokens: HashSet::new(),
}
}
pub fn is_token_restricted(&self, token_id: Uuid) -> bool {
self.restricted_tokens.contains(&token_id)
}
pub fn exceeds_limits(&self, daily_volume: Decimal, monthly_volume: Decimal) -> bool {
if let Some(daily_limit) = self.daily_limit {
if daily_volume > daily_limit {
return true;
}
}
if let Some(monthly_limit) = self.monthly_limit {
if monthly_volume > monthly_limit {
return true;
}
}
false
}
}
pub struct TransactionMonitor {
daily_volumes: HashMap<Uuid, Decimal>,
monthly_volumes: HashMap<Uuid, Decimal>,
alerts: Vec<SuspiciousActivityAlert>,
regional_rules: HashMap<Region, RegionalCompliance>,
user_regions: HashMap<Uuid, Region>,
}
impl TransactionMonitor {
pub fn new() -> Self {
let mut regional_rules = HashMap::new();
regional_rules.insert(Region::US, RegionalCompliance::us_standard());
regional_rules.insert(Region::EU, RegionalCompliance::eu_standard());
Self {
daily_volumes: HashMap::new(),
monthly_volumes: HashMap::new(),
alerts: Vec::new(),
regional_rules,
user_regions: HashMap::new(),
}
}
pub fn set_user_region(&mut self, user_id: Uuid, region: Region) {
self.user_regions.insert(user_id, region);
}
pub fn add_regional_rule(&mut self, compliance: RegionalCompliance) {
self.regional_rules.insert(compliance.region, compliance);
}
pub fn check_transaction(
&self,
user_id: Uuid,
token_id: Uuid,
amount: Decimal,
) -> Result<(), String> {
let region = self
.user_regions
.get(&user_id)
.copied()
.unwrap_or(Region::Other);
let compliance = self
.regional_rules
.get(®ion)
.ok_or_else(|| format!("No compliance rules for region {:?}", region))?;
if !compliance.trading_allowed {
return Err("Trading not allowed in your region".to_string());
}
if compliance.is_token_restricted(token_id) {
return Err("Token is restricted in your region".to_string());
}
let daily_volume = self
.daily_volumes
.get(&user_id)
.copied()
.unwrap_or(Decimal::ZERO);
let monthly_volume = self
.monthly_volumes
.get(&user_id)
.copied()
.unwrap_or(Decimal::ZERO);
if compliance.exceeds_limits(daily_volume + amount, monthly_volume + amount) {
return Err("Transaction would exceed regional limits".to_string());
}
Ok(())
}
pub fn record_transaction(&mut self, user_id: Uuid, amount: Decimal) {
*self.daily_volumes.entry(user_id).or_insert(Decimal::ZERO) += amount;
*self.monthly_volumes.entry(user_id).or_insert(Decimal::ZERO) += amount;
}
pub fn detect_suspicious_activity(
&mut self,
user_id: Uuid,
transaction_ids: Vec<Uuid>,
description: String,
activity_type: SuspiciousActivity,
risk_score: u8,
) {
let alert = SuspiciousActivityAlert {
id: Uuid::new_v4(),
user_id,
activity_type,
description,
risk_score,
transaction_ids,
detected_at: Utc::now(),
reviewed: false,
notes: None,
};
self.alerts.push(alert);
}
pub fn get_unreviewed_alerts(&self) -> Vec<&SuspiciousActivityAlert> {
self.alerts.iter().filter(|a| !a.reviewed).collect()
}
pub fn get_high_risk_alerts(&self, threshold: u8) -> Vec<&SuspiciousActivityAlert> {
self.alerts
.iter()
.filter(|a| a.risk_score >= threshold)
.collect()
}
pub fn review_alert(&mut self, alert_id: Uuid, notes: String) -> Result<(), String> {
if let Some(alert) = self.alerts.iter_mut().find(|a| a.id == alert_id) {
alert.reviewed = true;
alert.notes = Some(notes);
Ok(())
} else {
Err("Alert not found".to_string())
}
}
pub fn reset_daily_volumes(&mut self) {
self.daily_volumes.clear();
}
pub fn reset_monthly_volumes(&mut self) {
self.monthly_volumes.clear();
}
pub fn generate_report(&self) -> ComplianceReport {
let total_users = self.user_regions.len();
let total_alerts = self.alerts.len();
let unreviewed_alerts = self.get_unreviewed_alerts().len();
let high_risk_alerts = self.get_high_risk_alerts(75).len();
let mut users_by_region: HashMap<Region, usize> = HashMap::new();
for region in self.user_regions.values() {
*users_by_region.entry(*region).or_insert(0) += 1;
}
ComplianceReport {
generated_at: Utc::now(),
total_users,
total_alerts,
unreviewed_alerts,
high_risk_alerts,
users_by_region,
}
}
}
impl Default for TransactionMonitor {
fn default() -> Self {
Self::new()
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ComplianceReport {
pub generated_at: DateTime<Utc>,
pub total_users: usize,
pub total_alerts: usize,
pub unreviewed_alerts: usize,
pub high_risk_alerts: usize,
pub users_by_region: HashMap<Region, usize>,
}
pub struct RealTimeSurveillance {
rules: Vec<SurveillanceRule>,
violations: Vec<Violation>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SurveillanceRule {
pub id: Uuid,
pub name: String,
pub rule_type: RuleType,
pub enabled: bool,
pub threshold: Option<f64>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub enum RuleType {
VolumeSpike,
PriceManipulation,
InsiderTrading,
MarketManipulation,
Layering,
Spoofing,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Violation {
pub id: Uuid,
pub rule_id: Uuid,
pub user_id: Uuid,
pub detected_at: DateTime<Utc>,
pub severity: ViolationSeverity,
pub description: String,
pub evidence: Vec<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub enum ViolationSeverity {
Low,
Medium,
High,
Critical,
}
impl RealTimeSurveillance {
pub fn new() -> Self {
Self {
rules: Vec::new(),
violations: Vec::new(),
}
}
pub fn add_rule(&mut self, rule: SurveillanceRule) {
self.rules.push(rule);
}
pub fn record_violation(&mut self, violation: Violation) {
self.violations.push(violation);
}
pub fn get_violations_by_severity(&self, severity: ViolationSeverity) -> Vec<&Violation> {
self.violations
.iter()
.filter(|v| v.severity == severity)
.collect()
}
}
impl Default for RealTimeSurveillance {
fn default() -> Self {
Self::new()
}
}
pub struct AutomatedSAR {
pending_sars: Vec<SuspiciousActivityReport>,
submitted_sars: Vec<SuspiciousActivityReport>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SuspiciousActivityReport {
pub id: Uuid,
pub user_id: Uuid,
pub activity_type: SuspiciousActivity,
pub description: String,
pub start_date: DateTime<Utc>,
pub end_date: DateTime<Utc>,
pub total_amount: Decimal,
pub transaction_count: usize,
pub risk_indicators: Vec<String>,
pub status: SARStatus,
pub created_at: DateTime<Utc>,
pub submitted_at: Option<DateTime<Utc>>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub enum SARStatus {
Draft,
PendingReview,
Approved,
Submitted,
Rejected,
}
impl AutomatedSAR {
pub fn new() -> Self {
Self {
pending_sars: Vec::new(),
submitted_sars: Vec::new(),
}
}
#[allow(clippy::too_many_arguments)]
pub fn create_sar(
&mut self,
user_id: Uuid,
activity_type: SuspiciousActivity,
description: String,
start_date: DateTime<Utc>,
end_date: DateTime<Utc>,
total_amount: Decimal,
transaction_count: usize,
) -> Uuid {
let sar = SuspiciousActivityReport {
id: Uuid::new_v4(),
user_id,
activity_type,
description,
start_date,
end_date,
total_amount,
transaction_count,
risk_indicators: Vec::new(),
status: SARStatus::Draft,
created_at: Utc::now(),
submitted_at: None,
};
let id = sar.id;
self.pending_sars.push(sar);
id
}
pub fn submit_sar(&mut self, sar_id: Uuid) -> Result<(), String> {
if let Some(pos) = self.pending_sars.iter().position(|s| s.id == sar_id) {
let mut sar = self.pending_sars.remove(pos);
sar.status = SARStatus::Submitted;
sar.submitted_at = Some(Utc::now());
self.submitted_sars.push(sar);
Ok(())
} else {
Err("SAR not found".to_string())
}
}
pub fn get_pending_sars(&self) -> &[SuspiciousActivityReport] {
&self.pending_sars
}
}
impl Default for AutomatedSAR {
fn default() -> Self {
Self::new()
}
}
pub struct MarketManipulationDetector {
threshold: f64,
price_history: Vec<(DateTime<Utc>, Decimal)>,
}
impl MarketManipulationDetector {
pub fn new(threshold: f64) -> Self {
Self {
threshold,
price_history: Vec::new(),
}
}
pub fn add_price_point(&mut self, timestamp: DateTime<Utc>, price: Decimal) {
self.price_history.push((timestamp, price));
}
pub fn detect_pump_and_dump(&self, window_size: usize) -> Option<ManipulationAlert> {
if self.price_history.len() < window_size {
return None;
}
let recent = &self.price_history[self.price_history.len() - window_size..];
let prices: Vec<f64> = recent
.iter()
.map(|(_, p)| p.to_string().parse::<f64>().unwrap_or(0.0))
.collect();
let max_price = prices.iter().cloned().fold(f64::NEG_INFINITY, f64::max);
let min_price = prices.iter().cloned().fold(f64::INFINITY, f64::min);
let price_range = (max_price - min_price) / min_price;
if price_range > self.threshold {
Some(ManipulationAlert {
alert_type: ManipulationType::PumpAndDump,
detected_at: Utc::now(),
confidence: (price_range / self.threshold).min(1.0),
description: format!(
"Price increased by {:.2}% then dropped, indicating potential pump and dump",
price_range * 100.0
),
})
} else {
None
}
}
pub fn detect_wash_trading(
&self,
user_trades: &[(Uuid, Uuid, Decimal)], ) -> Option<ManipulationAlert> {
let mut same_party_trades = 0;
for (buyer, seller, _) in user_trades {
if buyer == seller {
same_party_trades += 1;
}
}
let wash_trade_ratio = same_party_trades as f64 / user_trades.len() as f64;
if wash_trade_ratio > 0.3 {
Some(ManipulationAlert {
alert_type: ManipulationType::WashTrading,
detected_at: Utc::now(),
confidence: wash_trade_ratio,
description: format!(
"{:.1}% of trades are between same parties",
wash_trade_ratio * 100.0
),
})
} else {
None
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub enum ManipulationType {
PumpAndDump,
WashTrading,
Layering,
Spoofing,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ManipulationAlert {
pub alert_type: ManipulationType,
pub detected_at: DateTime<Utc>,
pub confidence: f64,
pub description: String,
}
pub struct InsiderTradingDetector {
suspicious_trades: Vec<SuspiciousTrade>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SuspiciousTrade {
pub trade_id: Uuid,
pub user_id: Uuid,
pub token_id: Uuid,
pub amount: Decimal,
pub trade_timestamp: DateTime<Utc>,
pub event_timestamp: Option<DateTime<Utc>>,
pub event_type: Option<String>,
pub suspicion_score: f64,
}
impl InsiderTradingDetector {
pub fn new() -> Self {
Self {
suspicious_trades: Vec::new(),
}
}
#[allow(clippy::too_many_arguments)]
pub fn detect_pre_event_trading(
&mut self,
trade_id: Uuid,
user_id: Uuid,
token_id: Uuid,
amount: Decimal,
trade_timestamp: DateTime<Utc>,
event_timestamp: DateTime<Utc>,
event_type: String,
) {
let time_diff = (event_timestamp - trade_timestamp).num_hours();
if time_diff > 0 && time_diff <= 72 {
let suspicion_score = 1.0 - (time_diff as f64 / 72.0);
self.suspicious_trades.push(SuspiciousTrade {
trade_id,
user_id,
token_id,
amount,
trade_timestamp,
event_timestamp: Some(event_timestamp),
event_type: Some(event_type),
suspicion_score,
});
}
}
pub fn get_high_confidence_trades(&self, threshold: f64) -> Vec<&SuspiciousTrade> {
self.suspicious_trades
.iter()
.filter(|t| t.suspicion_score >= threshold)
.collect()
}
}
impl Default for InsiderTradingDetector {
fn default() -> Self {
Self::new()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_regional_compliance() {
let us_compliance = RegionalCompliance::us_standard();
assert!(us_compliance.trading_allowed);
assert!(us_compliance.kyc_required);
assert_eq!(us_compliance.min_age, 18);
}
#[test]
fn test_token_restriction() {
let mut compliance = RegionalCompliance::us_standard();
let token_id = Uuid::new_v4();
assert!(!compliance.is_token_restricted(token_id));
compliance.restricted_tokens.insert(token_id);
assert!(compliance.is_token_restricted(token_id));
}
#[test]
fn test_exceeds_limits() {
let compliance = RegionalCompliance::us_standard();
assert!(!compliance.exceeds_limits(dec!(1000), dec!(10000)));
assert!(compliance.exceeds_limits(dec!(60000), dec!(10000))); assert!(compliance.exceeds_limits(dec!(1000), dec!(600000))); }
#[test]
fn test_transaction_monitoring() {
let mut monitor = TransactionMonitor::new();
let user_id = Uuid::new_v4();
let token_id = Uuid::new_v4();
monitor.set_user_region(user_id, Region::US);
let result = monitor.check_transaction(user_id, token_id, dec!(1000));
assert!(result.is_ok());
monitor.record_transaction(user_id, dec!(1000));
let daily_volume = monitor.daily_volumes.get(&user_id).copied().unwrap();
assert_eq!(daily_volume, dec!(1000));
}
#[test]
fn test_suspicious_activity_detection() {
let mut monitor = TransactionMonitor::new();
let user_id = Uuid::new_v4();
monitor.detect_suspicious_activity(
user_id,
vec![Uuid::new_v4()],
"Unusual volume spike".to_string(),
SuspiciousActivity::UnusualVolume,
85,
);
let alerts = monitor.get_unreviewed_alerts();
assert_eq!(alerts.len(), 1);
assert_eq!(alerts[0].risk_score, 85);
}
#[test]
fn test_alert_review() {
let mut monitor = TransactionMonitor::new();
let user_id = Uuid::new_v4();
monitor.detect_suspicious_activity(
user_id,
vec![Uuid::new_v4()],
"Test alert".to_string(),
SuspiciousActivity::WashTrading,
75,
);
let alert_id = monitor.alerts[0].id;
let result = monitor.review_alert(alert_id, "Reviewed and cleared".to_string());
assert!(result.is_ok());
assert!(monitor.alerts[0].reviewed);
}
#[test]
fn test_high_risk_alerts() {
let mut monitor = TransactionMonitor::new();
let user_id = Uuid::new_v4();
monitor.detect_suspicious_activity(
user_id,
vec![Uuid::new_v4()],
"Low risk".to_string(),
SuspiciousActivity::FailedTransactions,
50,
);
monitor.detect_suspicious_activity(
user_id,
vec![Uuid::new_v4()],
"High risk".to_string(),
SuspiciousActivity::FrontRunning,
90,
);
let high_risk = monitor.get_high_risk_alerts(75);
assert_eq!(high_risk.len(), 1);
assert_eq!(high_risk[0].risk_score, 90);
}
#[test]
fn test_compliance_report() {
let mut monitor = TransactionMonitor::new();
monitor.set_user_region(Uuid::new_v4(), Region::US);
monitor.set_user_region(Uuid::new_v4(), Region::EU);
monitor.set_user_region(Uuid::new_v4(), Region::US);
monitor.detect_suspicious_activity(
Uuid::new_v4(),
vec![],
"Test".to_string(),
SuspiciousActivity::UnusualVolume,
50,
);
let report = monitor.generate_report();
assert_eq!(report.total_users, 3);
assert_eq!(report.total_alerts, 1);
assert_eq!(*report.users_by_region.get(&Region::US).unwrap(), 2);
assert_eq!(*report.users_by_region.get(&Region::EU).unwrap(), 1);
}
#[test]
fn test_volume_reset() {
let mut monitor = TransactionMonitor::new();
let user_id = Uuid::new_v4();
monitor.record_transaction(user_id, dec!(1000));
assert_eq!(
monitor.daily_volumes.get(&user_id).copied().unwrap(),
dec!(1000)
);
monitor.reset_daily_volumes();
assert_eq!(
monitor
.daily_volumes
.get(&user_id)
.copied()
.unwrap_or(Decimal::ZERO),
Decimal::ZERO
);
}
#[test]
fn test_real_time_surveillance() {
let mut surveillance = RealTimeSurveillance::new();
let rule = SurveillanceRule {
id: Uuid::new_v4(),
name: "Volume Spike Detection".to_string(),
rule_type: RuleType::VolumeSpike,
enabled: true,
threshold: Some(2.0),
};
surveillance.add_rule(rule);
let violation = Violation {
id: Uuid::new_v4(),
rule_id: Uuid::new_v4(),
user_id: Uuid::new_v4(),
detected_at: Utc::now(),
severity: ViolationSeverity::High,
description: "Unusual volume detected".to_string(),
evidence: vec!["Trade ID: 123".to_string()],
};
surveillance.record_violation(violation);
let high_severity = surveillance.get_violations_by_severity(ViolationSeverity::High);
assert_eq!(high_severity.len(), 1);
}
#[test]
fn test_automated_sar() {
let mut sar_system = AutomatedSAR::new();
let sar_id = sar_system.create_sar(
Uuid::new_v4(),
SuspiciousActivity::UnusualVolume,
"High volume spike detected".to_string(),
Utc::now(),
Utc::now(),
dec!(100000),
50,
);
assert_eq!(sar_system.get_pending_sars().len(), 1);
let result = sar_system.submit_sar(sar_id);
assert!(result.is_ok());
assert_eq!(sar_system.get_pending_sars().len(), 0);
}
#[test]
fn test_market_manipulation_detection() {
let mut detector = MarketManipulationDetector::new(0.5);
detector.add_price_point(Utc::now(), dec!(100));
detector.add_price_point(Utc::now(), dec!(150));
detector.add_price_point(Utc::now(), dec!(200));
detector.add_price_point(Utc::now(), dec!(120));
let alert = detector.detect_pump_and_dump(4);
assert!(alert.is_some());
if let Some(alert) = alert {
assert_eq!(alert.alert_type, ManipulationType::PumpAndDump);
assert!(alert.confidence > 0.0);
}
}
#[test]
fn test_wash_trading_detection() {
let detector = MarketManipulationDetector::new(0.5);
let user_id = Uuid::new_v4();
let trades = vec![
(user_id, user_id, dec!(100)), (user_id, Uuid::new_v4(), dec!(200)),
(user_id, user_id, dec!(150)), ];
let alert = detector.detect_wash_trading(&trades);
assert!(alert.is_some());
if let Some(alert) = alert {
assert_eq!(alert.alert_type, ManipulationType::WashTrading);
}
}
#[test]
fn test_insider_trading_detection() {
let mut detector = InsiderTradingDetector::new();
let trade_time = Utc::now();
let event_time = trade_time + chrono::Duration::hours(24);
detector.detect_pre_event_trading(
Uuid::new_v4(),
Uuid::new_v4(),
Uuid::new_v4(),
dec!(10000),
trade_time,
event_time,
"Major announcement".to_string(),
);
let suspicious = detector.get_high_confidence_trades(0.5);
assert_eq!(suspicious.len(), 1);
assert!(suspicious[0].suspicion_score > 0.0);
}
}