use crate::error::ValidationError;
use crate::peer_score::RejectReason;
use crate::MessageId;
use tet_libp2p_core::PeerId;
use log::debug;
use rand::seq::SliceRandom;
use rand::thread_rng;
use std::collections::HashMap;
use wasm_timer::Instant;
#[derive(Default)]
pub(crate) struct GossipPromises {
promises: HashMap<MessageId, HashMap<PeerId, Instant>>,
}
impl GossipPromises {
pub fn add_promise(&mut self, peer: PeerId, messages: &[MessageId], expires: Instant) {
let mut rng = thread_rng();
if let Some(message_id) = messages.choose(&mut rng) {
self.promises
.entry(message_id.clone())
.or_insert_with(HashMap::new)
.entry(peer)
.or_insert(expires);
}
}
pub fn message_delivered(&mut self, message_id: &MessageId) {
self.promises.remove(message_id);
}
pub fn reject_message(&mut self, message_id: &MessageId, reason: &RejectReason) {
match reason {
RejectReason::ValidationError(ValidationError::InvalidSignature) => (),
RejectReason::SelfOrigin => (),
_ => {
self.promises.remove(message_id);
}
};
}
pub fn get_broken_promises(&mut self) -> HashMap<PeerId, usize> {
let now = Instant::now();
let mut result = HashMap::new();
self.promises.retain(|msg, peers| {
peers.retain(|peer_id, expires| {
if *expires < now {
let count = result.entry(peer_id.clone()).or_insert(0);
*count += 1;
debug!(
"The peer {} broke the promise to deliver message {} in time!",
peer_id, msg
);
false
} else {
true
}
});
!peers.is_empty()
});
result
}
}