diff --git a/dash-spv/src/network/manager.rs b/dash-spv/src/network/manager.rs index 130eb715c..680eedf4d 100644 --- a/dash-spv/src/network/manager.rs +++ b/dash-spv/src/network/manager.rs @@ -16,9 +16,7 @@ use crate::network::addrv2::AddrV2Handler; use crate::network::constants::*; use crate::network::discovery::DnsDiscovery; use crate::network::pool::PeerPool; -use crate::network::reputation::{ - misbehavior_scores, positive_scores, PeerReputationManager, ReputationAware, -}; +use crate::network::reputation::{ChangeReason, PeerReputationManager, ReputationAware}; use crate::network::{ HandshakeManager, Message, MessageDispatcher, MessageType, NetworkEvent, NetworkManager, NetworkRequest, Peer, RequestSender, @@ -360,11 +358,7 @@ impl PeerNetworkManager { pool.remove_peer(&addr).await; // Update reputation for handshake failure reputation_manager - .update_reputation( - addr, - misbehavior_scores::INVALID_MESSAGE, - "Handshake failed", - ) + .update_reputation(addr, ChangeReason::HandshakeFailed) .await; // For handshake failures, try again later tokio::time::sleep(RECONNECT_DELAY).await; @@ -377,11 +371,7 @@ impl PeerNetworkManager { pool.remove_peer(&addr).await; // Minor reputation penalty for connection failure reputation_manager - .update_reputation( - addr, - misbehavior_scores::TIMEOUT / 2, - "Connection failed", - ) + .update_reputation(addr, ChangeReason::ConnectionFailed) .await; } } @@ -638,8 +628,7 @@ impl PeerNetworkManager { reputation_manager .update_reputation( addr, - misbehavior_scores::INVALID_MESSAGE, - "Headers2 decompression failed", + ChangeReason::Headers2DecompressionFailed, ) .await; continue; // Don't forward corrupted message @@ -698,11 +687,7 @@ impl PeerNetworkManager { tracing::debug!("Timeout reading from {}, continuing...", addr); // Minor reputation penalty for timeout reputation_manager - .update_reputation( - addr, - misbehavior_scores::TIMEOUT, - "Read timeout", - ) + .update_reputation(addr, ChangeReason::ReadTimeout) .await; continue; } @@ -722,8 +707,7 @@ impl PeerNetworkManager { reputation_manager .update_reputation( addr, - misbehavior_scores::INVALID_TRANSACTION, - "Invalid transaction type in block", + ChangeReason::InvalidTransactionInBlock, ) .await; } else if error_msg @@ -784,9 +768,7 @@ impl PeerNetworkManager { let conn_duration = Duration::from_secs(60 * loop_iteration); // Rough estimate if conn_duration > Duration::from_secs(3600) { // 1 hour - reputation_manager - .update_reputation(addr, positive_scores::LONG_UPTIME, "Long connection uptime") - .await; + reputation_manager.update_reputation(addr, ChangeReason::LongUptime).await; } }); } @@ -1013,9 +995,7 @@ impl PeerNetworkManager { if let Err(e) = peer_guard.send_ping().await { tracing::error!("Failed to ping {}: {}", addr, e); // Update reputation for ping failure - self.reputation_manager - .update_reputation(addr, misbehavior_scores::TIMEOUT, "Ping failed") - .await; + self.reputation_manager.update_reputation(addr, ChangeReason::PingFailed).await; } } let has_expired = peer_guard.remove_expired_pings(); @@ -1316,13 +1296,7 @@ impl PeerNetworkManager { self.disconnect_peer(addr, reason).await?; // Update reputation to trigger ban - self.reputation_manager - .update_reputation( - *addr, - misbehavior_scores::INVALID_HEADER * 2, // Severe penalty - reason, - ) - .await; + self.reputation_manager.update_reputation(*addr, ChangeReason::ManuallyBanned).await; Ok(()) } diff --git a/dash-spv/src/network/mod.rs b/dash-spv/src/network/mod.rs index a428ddde1..017af47d9 100644 --- a/dash-spv/src/network/mod.rs +++ b/dash-spv/src/network/mod.rs @@ -9,7 +9,7 @@ pub mod manager; mod message_dispatcher; pub mod peer; pub mod pool; -pub mod reputation; +mod reputation; mod message_type; #[cfg(test)] @@ -35,6 +35,7 @@ pub use manager::PeerNetworkManager; pub use message_dispatcher::{Message, MessageDispatcher}; pub use message_type::MessageType; pub use peer::Peer; +pub(crate) use reputation::PeerReputation; use std::net::SocketAddr; use tokio::sync::mpsc::UnboundedReceiver; diff --git a/dash-spv/src/network/reputation.rs b/dash-spv/src/network/reputation.rs index 5c65ce32b..f90656584 100644 --- a/dash-spv/src/network/reputation.rs +++ b/dash-spv/src/network/reputation.rs @@ -14,58 +14,51 @@ use std::sync::Arc; use std::time::{Duration, Instant}; use tokio::sync::RwLock; -/// Misbehavior score thresholds for different violations -pub mod misbehavior_scores { - /// Invalid message format or protocol violation - pub const INVALID_MESSAGE: i32 = 10; - - /// Invalid block header - pub const INVALID_HEADER: i32 = 50; - - /// Invalid compact filter - pub const INVALID_FILTER: i32 = 25; - - /// Timeout or slow response - pub const TIMEOUT: i32 = 5; - - /// Sending unsolicited data - pub const UNSOLICITED_DATA: i32 = 15; - - /// Invalid transaction - pub const INVALID_TRANSACTION: i32 = 20; - - /// Invalid masternode list diff - pub const INVALID_MASTERNODE_DIFF: i32 = 30; - - /// Invalid ChainLock - pub const INVALID_CHAINLOCK: i32 = 40; - - /// Invalid InstantLock - pub const INVALID_INSTANTLOCK: i32 = 35; - - /// Duplicate message - pub const DUPLICATE_MESSAGE: i32 = 5; - - /// Connection flood attempt - pub const CONNECTION_FLOOD: i32 = 20; +/// Reason for a peer reputation change. Each reason owns its score delta +/// (positive = penalty, negative = reward) and a human-readable label. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum ChangeReason { + HandshakeFailed, + ConnectionFailed, + Headers2DecompressionFailed, + ReadTimeout, + PingFailed, + InvalidTransactionInBlock, + ManuallyBanned, + LongUptime, } -/// Positive behavior scores -pub mod positive_scores { - /// Successfully provided valid headers - pub const VALID_HEADERS: i32 = -5; - - /// Successfully provided valid filters - pub const VALID_FILTERS: i32 = -3; - - /// Successfully provided valid block - pub const VALID_BLOCK: i32 = -10; - - /// Fast response time - pub const FAST_RESPONSE: i32 = -2; +impl ChangeReason { + /// Score delta for this reason: positive for misbehavior (penalty), + /// negative for good behavior (reward). + pub fn score(&self) -> i32 { + match self { + ChangeReason::HandshakeFailed => 10, + ChangeReason::ConnectionFailed => 2, + ChangeReason::Headers2DecompressionFailed => 10, + ChangeReason::ReadTimeout => 5, + ChangeReason::PingFailed => 5, + ChangeReason::InvalidTransactionInBlock => 20, + ChangeReason::ManuallyBanned => 100, + ChangeReason::LongUptime => -5, + } + } +} - /// Long uptime connection - pub const LONG_UPTIME: i32 = -5; +impl std::fmt::Display for ChangeReason { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + let label = match self { + ChangeReason::HandshakeFailed => "Handshake failed", + ChangeReason::ConnectionFailed => "Connection failed", + ChangeReason::Headers2DecompressionFailed => "Headers2 decompression failed", + ChangeReason::ReadTimeout => "Read timeout", + ChangeReason::PingFailed => "Ping failed", + ChangeReason::InvalidTransactionInBlock => "Invalid transaction type in block", + ChangeReason::ManuallyBanned => "Manually banned", + ChangeReason::LongUptime => "Long connection uptime", + }; + f.write_str(label) + } } /// Ban duration for misbehaving peers @@ -223,25 +216,10 @@ impl PeerReputation { } } -/// Reputation change event -#[derive(Debug, Clone)] -pub struct ReputationEvent { - pub peer: SocketAddr, - pub change: i32, - pub reason: String, - pub timestamp: Instant, -} - /// Peer reputation manager pub struct PeerReputationManager { /// Reputation data for each peer reputations: Arc>>, - - /// Recent reputation events for monitoring - recent_events: Arc>>, - - /// Maximum number of events to keep - max_events: usize, } impl Default for PeerReputationManager { @@ -255,18 +233,13 @@ impl PeerReputationManager { pub fn new() -> Self { Self { reputations: Arc::new(RwLock::new(HashMap::new())), - recent_events: Arc::new(RwLock::new(Vec::new())), - max_events: 1000, } } - /// Update peer reputation - pub async fn update_reputation( - &self, - peer: SocketAddr, - score_change: i32, - reason: &str, - ) -> bool { + /// Update peer reputation by the score delta of `reason`. + pub async fn update_reputation(&self, peer: SocketAddr, reason: ChangeReason) -> bool { + let score_change = reason.score(); + let mut reputations = self.reputations.write().await; let reputation = reputations.entry(peer).or_default(); @@ -311,32 +284,9 @@ impl PeerReputationManager { ); } - // Record event - let event = ReputationEvent { - peer, - change: score_change, - reason: reason.to_string(), - timestamp: Instant::now(), - }; - - drop(reputations); // Release lock before recording event - self.record_event(event).await; - should_ban } - /// Record a reputation event - async fn record_event(&self, event: ReputationEvent) { - let mut events = self.recent_events.write().await; - events.push(event); - - // Keep only recent events - if events.len() > self.max_events { - let drain_count = events.len() - self.max_events; - events.drain(0..drain_count); - } - } - /// Check if a peer is banned pub async fn is_banned(&self, peer: &SocketAddr) -> bool { let mut reputations = self.reputations.write().await; @@ -348,35 +298,6 @@ impl PeerReputationManager { } } - /// Get peer reputation score - pub async fn get_score(&self, peer: &SocketAddr) -> i32 { - let mut reputations = self.reputations.write().await; - if let Some(reputation) = reputations.get_mut(peer) { - reputation.apply_decay(); - reputation.score - } else { - 0 - } - } - - /// Temporarily ban a peer for a specified duration, regardless of score. - /// This can be used for critical protocol violations (e.g., invalid ChainLocks). - pub async fn temporary_ban_peer(&self, peer: SocketAddr, duration: Duration, reason: &str) { - let mut reputations = self.reputations.write().await; - let reputation = reputations.entry(peer).or_default(); - - reputation.banned_until = Some(Instant::now() + duration); - reputation.ban_count += 1; - - tracing::warn!( - "Peer {} temporarily banned for {:?} (ban #{}, reason: {})", - peer, - duration, - reputation.ban_count, - reason - ); - } - /// Record a connection attempt pub async fn record_connection_attempt(&self, peer: SocketAddr) { let mut reputations = self.reputations.write().await; @@ -404,11 +325,6 @@ impl PeerReputationManager { reputations.clone() } - /// Get recent reputation events - pub async fn get_recent_events(&self) -> Vec { - self.recent_events.read().await.clone() - } - /// Clear banned status for a peer (admin function) pub async fn unban_peer(&self, peer: &SocketAddr) { let mut reputations = self.reputations.write().await; @@ -419,33 +335,6 @@ impl PeerReputationManager { } } - /// Reset reputation for a peer - pub async fn reset_reputation(&self, peer: &SocketAddr) { - let mut reputations = self.reputations.write().await; - reputations.remove(peer); - tracing::info!("Reset reputation for peer {}", peer); - } - - /// Get peers sorted by reputation (best first) - pub async fn get_peers_by_reputation(&self) -> Vec<(SocketAddr, i32)> { - let mut reputations = self.reputations.write().await; - - // Apply decay and collect scores - let mut peer_scores: Vec<(SocketAddr, i32)> = reputations - .iter_mut() - .map(|(addr, rep)| { - rep.apply_decay(); - (*addr, rep.score) - }) - .filter(|(_, score)| *score < MAX_MISBEHAVIOR_SCORE) // Exclude banned peers - .collect(); - - // Sort by score (lower is better) - peer_scores.sort_by_key(|(_, score)| *score); - - peer_scores - } - /// Save reputation data to persistent storage pub async fn save_to_storage(&self, storage: &impl PeerStorage) -> std::io::Result<()> { let reputations = self.reputations.read().await; diff --git a/dash-spv/src/network/reputation_tests.rs b/dash-spv/src/network/reputation_tests.rs index b8b140057..68b74e13b 100644 --- a/dash-spv/src/network/reputation_tests.rs +++ b/dash-spv/src/network/reputation_tests.rs @@ -7,23 +7,22 @@ mod tests { use super::super::*; use std::net::SocketAddr; + async fn score(manager: &PeerReputationManager, peer: &SocketAddr) -> i32 { + manager.get_all_reputations().await.get(peer).map_or(0, |rep| rep.score) + } + #[tokio::test] async fn test_basic_reputation_operations() { let manager = PeerReputationManager::new(); let peer: SocketAddr = "127.0.0.1:8333".parse().unwrap(); - // Initial score should be 0 - assert_eq!(manager.get_score(&peer).await, 0); + assert_eq!(score(&manager, &peer).await, 0); - // Test misbehavior - manager - .update_reputation(peer, misbehavior_scores::INVALID_MESSAGE, "Test invalid message") - .await; - assert_eq!(manager.get_score(&peer).await, 10); + manager.update_reputation(peer, ChangeReason::HandshakeFailed).await; + assert_eq!(score(&manager, &peer).await, 10); - // Test positive behavior - manager.update_reputation(peer, positive_scores::VALID_HEADERS, "Test valid headers").await; - assert_eq!(manager.get_score(&peer).await, 5); + manager.update_reputation(peer, ChangeReason::LongUptime).await; + assert_eq!(score(&manager, &peer).await, 5); } #[tokio::test] @@ -31,17 +30,9 @@ mod tests { let manager = PeerReputationManager::new(); let peer: SocketAddr = "192.168.1.1:8333".parse().unwrap(); - // Accumulate misbehavior + // Banned on the 10th violation (10 * 10 = 100). for i in 0..10 { - let banned = manager - .update_reputation( - peer, - misbehavior_scores::INVALID_MESSAGE, - &format!("Violation {}", i), - ) - .await; - - // Should be banned on the 10th violation (total score = 100) + let banned = manager.update_reputation(peer, ChangeReason::HandshakeFailed).await; if i == 9 { assert!(banned); } else { @@ -58,11 +49,10 @@ mod tests { let peer1: SocketAddr = "10.0.0.1:8333".parse().unwrap(); let peer2: SocketAddr = "10.0.0.2:8333".parse().unwrap(); - // Set reputations - manager.update_reputation(peer1, -10, "Good peer").await; - manager.update_reputation(peer2, 50, "Bad peer").await; + manager.update_reputation(peer1, ChangeReason::LongUptime).await; + manager.update_reputation(peer1, ChangeReason::LongUptime).await; + manager.update_reputation(peer2, ChangeReason::InvalidTransactionInBlock).await; - // Save and load let temp_dir = tempfile::TempDir::new().unwrap(); let peer_storage = PersistentPeerStorage::open(temp_dir.path()) .await @@ -72,9 +62,8 @@ mod tests { let new_manager = PeerReputationManager::new(); new_manager.load_from_storage(&peer_storage).await.unwrap(); - // Verify scores were preserved - assert_eq!(new_manager.get_score(&peer1).await, -10); - assert_eq!(new_manager.get_score(&peer2).await, 50); + assert_eq!(score(&new_manager, &peer1).await, -10); + assert_eq!(score(&new_manager, &peer2).await, 20); } #[tokio::test] @@ -85,15 +74,17 @@ mod tests { let neutral_peer = AddrV2Message::dummy(0, "2.2.2.2".parse().unwrap(), 8333); let bad_peer = AddrV2Message::dummy(0, "3.3.3.3".parse().unwrap(), 8333); - // Set different reputations - manager.update_reputation(good_peer.socket_addr().unwrap(), -20, "Very good").await; - manager.update_reputation(bad_peer.socket_addr().unwrap(), 80, "Very bad").await; - // neutral_peer has default score of 0 + manager.update_reputation(good_peer.socket_addr().unwrap(), ChangeReason::LongUptime).await; + manager + .update_reputation( + bad_peer.socket_addr().unwrap(), + ChangeReason::InvalidTransactionInBlock, + ) + .await; let all_peers = vec![good_peer.clone(), neutral_peer.clone(), bad_peer.clone()]; let selected = manager.select_best_peers(all_peers, 2).await; - // Should select good_peer first, then neutral_peer assert_eq!(selected.len(), 2); assert_eq!(selected[0], good_peer.socket_addr().unwrap()); assert_eq!(selected[1], neutral_peer.socket_addr().unwrap()); diff --git a/dash-spv/src/storage/peers.rs b/dash-spv/src/storage/peers.rs index 69e075ab9..360e83650 100644 --- a/dash-spv/src/storage/peers.rs +++ b/dash-spv/src/storage/peers.rs @@ -10,7 +10,7 @@ use dashcore::{ use crate::{ error::StorageResult, - network::reputation::PeerReputation, + network::PeerReputation, storage::{io::atomic_write, PersistentStorage}, StorageError, }; @@ -97,12 +97,12 @@ impl PeerStorage for PersistentPeerStorage { return Ok(Vec::new()); }; - let mut peers = Vec::new(); - let peers = tokio::task::spawn_blocking(move || { let file = File::open(&peers_file)?; let mut reader = BufReader::new(file); + let mut peers = Vec::new(); + loop { match AddrV2Message::consensus_decode(&mut reader) { Ok(peer) => peers.push(peer), @@ -116,7 +116,6 @@ impl PeerStorage for PersistentPeerStorage { } } } - Ok(peers) }) .await