@@ -353,6 +353,12 @@ impl EventMonitor {
353353 // Auto-emit audit log for security-relevant events
354354 self . emit_audit_log ( event) . await ;
355355
356+ // Update user analytics for trade lifecycle events
357+ self . update_user_analytics ( event) . await ;
358+
359+ // Auto-emit audit log for security-relevant events
360+ self . emit_audit_log ( event) . await ;
361+
356362 let mut queue = self . job_queue . lock ( ) . await ;
357363 if let Err ( e) = queue. enqueue ( event_job) . await {
358364 error ! ( "Failed to enqueue event job for event {}: {}" , event. id, e) ;
@@ -536,6 +542,116 @@ impl EventMonitor {
536542 }
537543 }
538544
545+ async fn update_user_analytics ( & self , event : & Event ) {
546+ let d = & event. data ;
547+ let upsert = |address : & str , seller : bool , buyer : bool , amount : i64 , completed : bool , disputed : bool , cancelled : bool | {
548+ let pool = self . database . pool ( ) . clone ( ) ;
549+ let address = address. to_string ( ) ;
550+ tokio:: spawn ( async move {
551+ let _ = sqlx:: query (
552+ r#"
553+ INSERT INTO user_analytics
554+ (address, total_trades, trades_as_seller, trades_as_buyer,
555+ total_volume, completed_trades, disputed_trades, cancelled_trades, updated_at)
556+ VALUES ($1, $2, $3, $4, $5, $6, $7, $8, NOW())
557+ ON CONFLICT (address) DO UPDATE SET
558+ total_trades = user_analytics.total_trades + EXCLUDED.total_trades,
559+ trades_as_seller = user_analytics.trades_as_seller + EXCLUDED.trades_as_seller,
560+ trades_as_buyer = user_analytics.trades_as_buyer + EXCLUDED.trades_as_buyer,
561+ total_volume = user_analytics.total_volume + EXCLUDED.total_volume,
562+ completed_trades = user_analytics.completed_trades + EXCLUDED.completed_trades,
563+ disputed_trades = user_analytics.disputed_trades + EXCLUDED.disputed_trades,
564+ cancelled_trades = user_analytics.cancelled_trades + EXCLUDED.cancelled_trades,
565+ updated_at = NOW()
566+ "# ,
567+ )
568+ . bind ( & address)
569+ . bind ( if seller || buyer { 1i32 } else { 0 } )
570+ . bind ( if seller { 1i32 } else { 0 } )
571+ . bind ( if buyer { 1i32 } else { 0 } )
572+ . bind ( amount)
573+ . bind ( if completed { 1i32 } else { 0 } )
574+ . bind ( if disputed { 1i32 } else { 0 } )
575+ . bind ( if cancelled { 1i32 } else { 0 } )
576+ . execute ( & pool)
577+ . await ;
578+ } ) ;
579+ } ;
580+
581+ match event. event_type . as_str ( ) {
582+ "trade_created" => {
583+ let seller = d. get ( "seller" ) . and_then ( |v| v. as_str ( ) ) . unwrap_or_default ( ) ;
584+ let buyer = d. get ( "buyer" ) . and_then ( |v| v. as_str ( ) ) . unwrap_or_default ( ) ;
585+ let amount = d. get ( "amount" ) . and_then ( |v| v. as_i64 ( ) ) . unwrap_or ( 0 ) ;
586+ upsert ( seller, true , false , amount, false , false , false ) ;
587+ upsert ( buyer, false , true , amount, false , false , false ) ;
588+ }
589+ "trade_confirmed" => {
590+ let seller = d. get ( "seller" ) . and_then ( |v| v. as_str ( ) ) . unwrap_or_default ( ) ;
591+ let buyer = d. get ( "buyer" ) . and_then ( |v| v. as_str ( ) ) . unwrap_or_default ( ) ;
592+ upsert ( seller, false , false , 0 , true , false , false ) ;
593+ upsert ( buyer, false , false , 0 , true , false , false ) ;
594+ }
595+ "trade_disputed" => {
596+ let seller = d. get ( "seller" ) . and_then ( |v| v. as_str ( ) ) . unwrap_or_default ( ) ;
597+ let buyer = d. get ( "buyer" ) . and_then ( |v| v. as_str ( ) ) . unwrap_or_default ( ) ;
598+ upsert ( seller, false , false , 0 , false , true , false ) ;
599+ upsert ( buyer, false , false , 0 , false , true , false ) ;
600+ }
601+ "trade_cancelled" => {
602+ let seller = d. get ( "seller" ) . and_then ( |v| v. as_str ( ) ) . unwrap_or_default ( ) ;
603+ upsert ( seller, false , false , 0 , false , false , true ) ;
604+ }
605+ _ => { }
606+ }
607+ }
608+
609+ async fn emit_audit_log ( & self , event : & Event ) {
610+ use crate :: models:: { AuditCategory , AuditOutcome , AuditSeverity , NewAuditLog } ;
611+
612+ let ( category, action, severity) = match event. event_type . as_str ( ) {
613+ "trade_created" => ( AuditCategory :: Trade , "trade.created" , AuditSeverity :: Info ) ,
614+ "trade_funded" => ( AuditCategory :: Trade , "trade.funded" , AuditSeverity :: Info ) ,
615+ "trade_confirmed" => ( AuditCategory :: Trade , "trade.confirmed" , AuditSeverity :: Info ) ,
616+ "trade_cancelled" => ( AuditCategory :: Trade , "trade.cancelled" , AuditSeverity :: Warn ) ,
617+ "dispute_raised" => ( AuditCategory :: Trade , "trade.disputed" , AuditSeverity :: Warn ) ,
618+ "dispute_resolved" => ( AuditCategory :: Trade , "trade.resolved" , AuditSeverity :: Info ) ,
619+ "arb_reg" => ( AuditCategory :: Admin , "arbitrator.registered" , AuditSeverity :: Info ) ,
620+ "arb_rem" => ( AuditCategory :: Admin , "arbitrator.removed" , AuditSeverity :: Warn ) ,
621+ "fee_updated" => ( AuditCategory :: Admin , "fee.updated" , AuditSeverity :: Warn ) ,
622+ "paused" => ( AuditCategory :: Security , "contract.paused" , AuditSeverity :: Error ) ,
623+ "unpaused" => ( AuditCategory :: Security , "contract.unpaused" , AuditSeverity :: Warn ) ,
624+ "emrg_wd" => ( AuditCategory :: Security , "emergency.withdraw" , AuditSeverity :: Critical ) ,
625+ _ => return , // skip non-auditable events
626+ } ;
627+
628+ let actor = event. data . get ( "seller" )
629+ . or_else ( || event. data . get ( "admin" ) )
630+ . or_else ( || event. data . get ( "arbitrator" ) )
631+ . and_then ( |v| v. as_str ( ) )
632+ . unwrap_or ( "contract" )
633+ . to_string ( ) ;
634+
635+ let trade_id = event. data . get ( "trade_id" ) . and_then ( |v| v. as_str ( ) ) . map ( String :: from) ;
636+
637+ let entry = NewAuditLog {
638+ actor,
639+ category,
640+ action : action. to_string ( ) ,
641+ resource_type : Some ( "trade" . to_string ( ) ) ,
642+ resource_id : trade_id,
643+ outcome : AuditOutcome :: Success ,
644+ ledger : Some ( event. ledger ) ,
645+ tx_hash : Some ( event. transaction_hash . clone ( ) ) ,
646+ metadata : event. data . clone ( ) ,
647+ severity,
648+ } ;
649+
650+ if let Err ( e) = self . database . insert_audit_log ( & entry) . await {
651+ error ! ( "Failed to emit audit log for {}: {}" , event. event_type, e) ;
652+ }
653+ }
654+
539655 async fn get_latest_ledger ( & self ) -> Result < i64 , AppError > {
540656 let url = format ! ( "{}/ledgers?order=desc&limit=1" , self . config. horizon_url) ;
541657 let response: HorizonResponse < Ledger > = self . client . get ( & url) . send ( ) . await ?. json ( ) . await ?;
0 commit comments