Skip to content
This repository was archived by the owner on Nov 15, 2023. It is now read-only.
Closed
Show file tree
Hide file tree
Changes from 1 commit
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
188 changes: 94 additions & 94 deletions core/network/src/consensus_gossip.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@ use rand::{self, seq::SliceRandom};
use lru_cache::LruCache;
use network_libp2p::{Severity, NodeIndex};
use runtime_primitives::traits::{Block as BlockT, Hash, HashFor};
pub use crate::message::generic::{Message, ConsensusMessage};
pub use crate::message::generic::{Message, ConsensusMessage, ConsensusStatusMessage};
use crate::protocol::Context;
use crate::config::Roles;
use crate::ConsensusEngineId;
Expand All @@ -40,26 +40,16 @@ struct PeerConsensus<H> {
is_authority: bool,
}

#[derive(Clone, Copy)]
enum Status {
Live,
Future,
}

struct MessageEntry<B: BlockT> {
message_hash: B::Hash,
topic: B::Hash,
message: ConsensusMessage,
timestamp: Instant,
status: Status,
}

/// Message validation result.
pub enum ValidationResult<H> {
/// Message is valid with this topic.
Valid(H),
/// Message is future with this topic.
Future(H),
/// Invalid message.
Invalid,
/// Obsolete message.
Expand All @@ -68,13 +58,30 @@ pub enum ValidationResult<H> {

/// Validates consensus messages.
pub trait Validator<H> {
/// New peer is connected.
fn new_peer(&self, _who: NodeIndex, _roles: Roles) {
}

/// New connection is dropped.
fn peer_disconnected(&self, _who: NodeIndex) {
}

/// Validate consensus message.
fn validate(&self, data: &[u8]) -> ValidationResult<H>;

/// Process status message.
fn on_status(&self, _who: NodeIndex, _data: &[u8]) {
}

/// Filter out departing messages.
fn should_send_to(&self, _who: NodeIndex, _topic: &H, _data: &[u8]) -> bool {

@rphmeier rphmeier Mar 18, 2019

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

a bit unfortunate that we've got to re-decode all messages each time we try to send them out.

In GRANDPA we'll likely do it by topic, so it'll be zero-cost...

true
}

/// Produce a closure for validating messages on a given topic.
fn message_expired<'a>(&'a self) -> Box<FnMut(H, &[u8]) -> bool + 'a> {
Box::new(move |_topic, data| match self.validate(data) {
ValidationResult::Valid(_) | ValidationResult::Future(_) => false,
ValidationResult::Valid(_) => false,
ValidationResult::Invalid | ValidationResult::Expired => true,
})
}
Expand Down Expand Up @@ -112,36 +119,22 @@ impl<B: BlockT> ConsensusGossip<B> {
}

/// Handle new connected peer.
pub fn new_peer(&mut self, protocol: &mut Context<B>, who: NodeIndex, roles: Roles) {
if roles.intersects(Roles::AUTHORITY) {
trace!(target:"gossip", "Registering {:?} {}", roles, who);
let now = Instant::now();
// Send out all known messages to authorities.
let mut known_messages = HashSet::new();
for entry in self.messages.iter() {
if entry.timestamp + MESSAGE_LIFETIME < now { continue }
if let Status::Future = entry.status { continue }

known_messages.insert(entry.message_hash);
protocol.send_message(who, Message::Consensus(entry.message.clone()));
}
self.peers.insert(who, PeerConsensus {
known_messages,
is_authority: true,
});
}
else if roles.intersects(Roles::FULL) {
self.peers.insert(who, PeerConsensus {
known_messages: HashSet::new(),
is_authority: false,
});
pub fn new_peer(&mut self, _protocol: &mut Context<B>, who: NodeIndex, roles: Roles) {
trace!(target:"gossip", "Registering {:?} {}", roles, who);
self.peers.insert(who, PeerConsensus {
known_messages: HashSet::new(),
is_authority: roles.intersects(Roles::AUTHORITY),
Comment thread
arkpar marked this conversation as resolved.
Outdated
});
for v in self.validators.values() {
v.new_peer(who, roles);
}
}

fn propagate<F>(
&mut self,
protocol: &mut Context<B>,
message_hash: B::Hash,
topic: B::Hash,
get_message: F,
)
where F: Fn() -> ConsensusMessage,
Expand All @@ -164,44 +157,47 @@ impl<B: BlockT> ConsensusGossip<B> {
};

for (id, ref mut peer) in self.peers.iter_mut() {
if peer.known_messages.contains(&message_hash) {
continue
}
let message = get_message();
if let Some(validator) = self.validators.get(&message.engine_id) {
if !validator.should_send_to(*id, &topic, &message.data) {
continue
}
}
peer.known_messages.insert(message_hash.clone());
if peer.is_authority {
if peer.known_messages.insert(message_hash.clone()) {
let message = get_message();
trace!(target:"gossip", "Propagating to authority {}: {:?}", id, message);
protocol.send_message(*id, Message::Consensus(message));
}
} else if non_authorities.contains(&id) {
if peer.known_messages.insert(message_hash.clone()) {
let message = get_message();
trace!(target:"gossip", "Propagating to {}: {:?}", id, message);
protocol.send_message(*id, Message::Consensus(message));
}
}
protocol.send_message(*id, Message::Consensus(message));
}
}

fn register_message<F>(
&mut self,
message_hash: B::Hash,
topic: B::Hash,
status: Status,
get_message: F,
)
where F: Fn() -> ConsensusMessage
{
if self.known_messages.insert(message_hash, ()).is_none() {
self.messages.push(MessageEntry {
topic,
message_hash,
message: get_message(),
timestamp: Instant::now(),
status,
});
}
}

/// Call when a peer has been disconnected to stop tracking gossip status.
pub fn peer_disconnected(&mut self, _protocol: &mut Context<B>, who: NodeIndex) {
for v in self.validators.values() {
v.peer_disconnected(who);
}
self.peers.remove(&who);
}

Expand Down Expand Up @@ -254,39 +250,11 @@ impl<B: BlockT> ConsensusGossip<B> {
-> mpsc::UnboundedReceiver<Vec<u8>>
{
let (tx, rx) = mpsc::unbounded();

let validator = match self.validators.get(&engine_id) {
None => {
self.live_message_sinks.entry((engine_id, topic)).or_default().push(tx);
return rx;
}
Some(v) => v,
};

for entry in self.messages.iter_mut()
.filter(|e| e.topic == topic && e.message.engine_id == engine_id)
{
let live = match entry.status {
Status::Live => true,
Status::Future => match validator.validate(&entry.message.data) {
ValidationResult::Valid(_) => {
entry.status = Status::Live;
true
}
_ => {
// don't send messages considered to be future still.
// if messages are considered expired they'll be cleaned up when we
// collect garbage.
false
}
}
};

if live {
entry.status = Status::Live;
tx.unbounded_send(entry.message.data.clone())
.expect("receiver known to be live; qed");
}
tx.unbounded_send(entry.message.data.clone())
.expect("receiver known to be live; qed");
}

self.live_message_sinks.entry((engine_id, topic)).or_default().push(tx);
Expand Down Expand Up @@ -316,11 +284,10 @@ impl<B: BlockT> ConsensusGossip<B> {

let engine_id = message.engine_id;
//validate the message
let (topic, status) = match self.validators.get(&engine_id)
let topic = match self.validators.get(&engine_id)
.map(|v| v.validate(&message.data))
{
Some(ValidationResult::Valid(topic)) => (topic, Status::Live),
Some(ValidationResult::Future(topic)) => (topic, Status::Future),
Some(ValidationResult::Valid(topic)) => topic,
Some(ValidationResult::Invalid) => {
trace!(target:"gossip", "Invalid message from {}", who);
protocol.report_peer(
Expand Down Expand Up @@ -357,14 +324,52 @@ impl<B: BlockT> ConsensusGossip<B> {
entry.remove_entry();
}
}
self.multicast_inner(protocol, message_hash, topic, status, || message.clone());
Comment thread
arkpar marked this conversation as resolved.
Outdated
self.multicast_inner(protocol, message_hash, topic, || message.clone());
Some((topic, message))
} else {
trace!(target:"gossip", "Ignored statement from unregistered peer {}", who);
None
}
}


/// Handle an incoming ConsensusStatusMessage. Forwards it to the consensus engine handler.
pub fn on_status(
&mut self,
protocol: &mut Context<B>,
who: NodeIndex,
message: ConsensusStatusMessage,
) {
if let Some(_) = self.peers.get_mut(&who) {
let engine_id = message.engine_id;
//validate the message
if let Some(validator) = self.validators.get(&engine_id) {
validator.on_status(who, &message.data)
Comment thread
arkpar marked this conversation as resolved.
Outdated
} else {
protocol.report_peer(
who,
Severity::Useless(format!("Sent status with unknown consensus engine id")),
);
trace!(target:"gossip", "Unknown status message engine id {:?} from {}", engine_id, who);
}
} else {
trace!(target:"gossip", "Ignored status from unregistered peer {}", who);
}
}

/// Send out consensus status message to a particular peer.
pub fn send_consensus_status(
&mut self,
protocol: &mut Context<B>,
peer: NodeIndex,
message: ConsensusStatusMessage,
) {
if let Some(_) = self.peers.get(&peer) {
trace!(target:"gossip", "Sending consensus status to {}", peer);
protocol.send_message(peer, Message::ConsensusStatus(message));
}
}

/// Multicast a message to all peers.
pub fn multicast(
&mut self,
Expand All @@ -373,23 +378,20 @@ impl<B: BlockT> ConsensusGossip<B> {
message: ConsensusMessage,
) {
let message_hash = HashFor::<B>::hash(&message.data);
self.multicast_inner(protocol, message_hash, topic, Status::Live, || message.clone());
self.multicast_inner(protocol, message_hash, topic, || message.clone());
}

fn multicast_inner<F>(
&mut self,
protocol: &mut Context<B>,
message_hash: B::Hash,
topic: B::Hash,
status: Status,
get_message: F,
)
where F: Fn() -> ConsensusMessage
{
self.register_message(message_hash, topic, status, &get_message);
if let Status::Live = status {
self.propagate(protocol, message_hash, get_message);
}
self.register_message(message_hash, topic, &get_message);
self.propagate(protocol, message_hash, topic, get_message);
}

/// Note new consensus session.
Expand All @@ -413,10 +415,8 @@ mod tests {
if $consensus.known_messages.insert($hash, ()).is_none() {
$consensus.messages.push(MessageEntry {
topic: $topic,
message_hash: $hash,
message: ConsensusMessage { data: $m, engine_id: [0, 0, 0, 0]},
timestamp: $now,
status: Status::Live,
});
}
}
Expand Down Expand Up @@ -493,7 +493,7 @@ mod tests {
let message_hash = HashFor::<Block>::hash(&message.data);
let topic = HashFor::<Block>::hash(&[1,2,3]);

consensus.register_message(message_hash, topic, Status::Live, || message.clone());
consensus.register_message(message_hash, topic, || message.clone());
let stream = consensus.messages_for([0, 0, 0, 0], topic);

assert_eq!(stream.wait().next(), Some(Ok(message.data)));
Expand All @@ -507,8 +507,8 @@ mod tests {
let msg_a = ConsensusMessage { data: vec![1, 2, 3], engine_id: [0, 0, 0, 0] };
let msg_b = ConsensusMessage { data: vec![4, 5, 6], engine_id: [0, 0, 0, 0] };

consensus.register_message(HashFor::<Block>::hash(&msg_a.data), topic, Status::Live, || msg_a.clone());
consensus.register_message(HashFor::<Block>::hash(&msg_b.data), topic, Status::Live, || msg_b.clone());
consensus.register_message(HashFor::<Block>::hash(&msg_a.data), topic, || msg_a.clone());
consensus.register_message(HashFor::<Block>::hash(&msg_b.data), topic, || msg_b.clone());

assert_eq!(consensus.messages.len(), 2);
}
Expand All @@ -523,7 +523,7 @@ mod tests {
let message_hash = HashFor::<Block>::hash(&message.data);
let topic = HashFor::<Block>::hash(&[1,2,3]);

consensus.register_message(message_hash, topic, Status::Live, || message.clone());
consensus.register_message(message_hash, topic, || message.clone());

let stream1 = consensus.messages_for([0, 0, 0, 0], topic);
let stream2 = consensus.messages_for([0, 0, 0, 0], topic);
Expand All @@ -541,8 +541,8 @@ mod tests {
let msg_a = ConsensusMessage { data: vec![1, 2, 3], engine_id: [0, 0, 0, 0] };
let msg_b = ConsensusMessage { data: vec![4, 5, 6], engine_id: [0, 0, 0, 1] };

consensus.register_message(HashFor::<Block>::hash(&msg_a.data), topic, Status::Live, || msg_a.clone());
consensus.register_message(HashFor::<Block>::hash(&msg_b.data), topic, Status::Live, || msg_b.clone());
consensus.register_message(HashFor::<Block>::hash(&msg_a.data), topic, || msg_a.clone());
consensus.register_message(HashFor::<Block>::hash(&msg_b.data), topic, || msg_b.clone());

let mut stream = consensus.messages_for([0, 0, 0, 0], topic).wait();

Expand Down
12 changes: 12 additions & 0 deletions core/network/src/message.rs
Original file line number Diff line number Diff line change
Expand Up @@ -185,6 +185,8 @@ pub mod generic {
Transactions(Transactions<Extrinsic>),
/// Consensus protocol message.
Consensus(ConsensusMessage),
/// Consensus status message.
ConsensusStatus(ConsensusStatusMessage),
/// Remote method call request.
RemoteCallRequest(RemoteCallRequest<Hash>),
/// Remote method call response.
Expand Down Expand Up @@ -225,6 +227,7 @@ pub mod generic {
Message::BlockAnnounce(_) => CustomMessageId::OneWay,
Message::Transactions(_) => CustomMessageId::OneWay,
Message::Consensus(_) => CustomMessageId::OneWay,
Message::ConsensusStatus(_) => CustomMessageId::OneWay,
Message::RemoteCallRequest(ref req) => CustomMessageId::Request(req.id),
Message::RemoteCallResponse(ref resp) => CustomMessageId::Response(resp.id),
Message::RemoteReadRequest(ref req) => CustomMessageId::Request(req.id),
Expand Down Expand Up @@ -257,6 +260,15 @@ pub mod generic {
pub chain_status: Vec<u8>,
}

/// Consensus status broadcasted periodically and on request.
#[derive(Debug, PartialEq, Eq, Clone, Encode, Decode)]
pub struct ConsensusStatusMessage {
/// Engine id
pub engine_id: ConsensusEngineId,
/// Engine-specific status message.
pub data: Vec<u8>,
}

/// Request block data from a peer.
#[derive(Debug, PartialEq, Eq, Clone, Encode, Decode)]
pub struct BlockRequest<Hash, Number> {
Expand Down
Loading