Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
44 changes: 9 additions & 35 deletions dash-spv/src/network/manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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;
Expand All @@ -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;
}
}
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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;
}
Expand All @@ -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
Expand Down Expand Up @@ -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;
}
});
}
Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -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(())
}
Expand Down
3 changes: 2 additions & 1 deletion dash-spv/src/network/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)]
Expand All @@ -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;

Expand Down
203 changes: 46 additions & 157 deletions dash-spv/src/network/reputation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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<RwLock<HashMap<SocketAddr, PeerReputation>>>,

/// Recent reputation events for monitoring
recent_events: Arc<RwLock<Vec<ReputationEvent>>>,

/// Maximum number of events to keep
max_events: usize,
}

impl Default for PeerReputationManager {
Expand All @@ -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();

Expand Down Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -404,11 +325,6 @@ impl PeerReputationManager {
reputations.clone()
}

/// Get recent reputation events
pub async fn get_recent_events(&self) -> Vec<ReputationEvent> {
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;
Expand All @@ -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;
Expand Down
Loading
Loading