Skip to main content

consensus_core/network/
mod.rs

1// Copyright (c) Mysten Labs, Inc.
2// SPDX-License-Identifier: Apache-2.0
3
4//! This module defines the network interface, and provides network implementations for the
5//! consensus protocol.
6//!
7//! Having an abstract network interface allows
8//! - simplifying the semantics of sending data and serving requests over the network
9//! - hiding implementation specific types and semantics from the consensus protocol
10//! - allowing easy swapping of network implementations, for better performance or testing
11//!
12//! When modifying the client and server interfaces, the principle is to keep the interfaces
13//! low level, close to underlying implementations in semantics. For example, the client interface
14//! exposes sending messages to a specific peer, instead of broadcasting to all peers. Subscribing
15//! to a stream of blocks gets back the stream via response, instead of delivering the stream
16//! directly to the server. This keeps the logic agnostics to the underlying network outside of
17//! this module, so they can be reused easily across network implementations.
18
19use std::{
20    collections::BTreeSet,
21    fmt::{Display, Formatter},
22    net::SocketAddrV6,
23    pin::Pin,
24    sync::Arc,
25    time::Duration,
26};
27
28use async_trait::async_trait;
29use bytes::Bytes;
30use consensus_config::{AuthorityIndex, NetworkKeyPair, NetworkPublicKey};
31use consensus_types::block::{BlockRef, Round};
32use fastcrypto::encoding::{Encoding, Hex};
33use futures::Stream;
34use mysten_network::{Multiaddr, multiaddr::Protocol};
35use prost::Message as _;
36
37use crate::{
38    block::{ExtendedBlock, VerifiedBlock},
39    commit::{CommitRange, TrustedCommit},
40    context::Context,
41    error::{ConsensusError, ConsensusResult},
42};
43
44/// Identifies an observer node by its network public key.
45pub type NodeId = NetworkPublicKey;
46
47/// Identifies a peer in the network, which can be either a validator or an observer.
48/// The Observer variant is boxed to keep the enum small, since `NodeId` (32 bytes) is
49/// much larger than `AuthorityIndex` (4 bytes).
50#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
51pub enum PeerId {
52    /// A validator node identified by its authority index.
53    Validator(AuthorityIndex),
54    /// An observer node identified by its network public key.
55    Observer(Box<NodeId>),
56}
57
58impl Display for PeerId {
59    fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
60        match self {
61            PeerId::Validator(authority) => write!(f, "[{}]", authority),
62            PeerId::Observer(node_id) => {
63                let bytes = node_id.to_bytes();
64                let s = Hex::encode(bytes.get(0..4).ok_or(std::fmt::Error)?);
65                write!(f, "o#{}..", s)
66            }
67        }
68    }
69}
70
71impl PeerId {
72    /// Returns a human-readable name suitable for logging. For observers, prints
73    /// the first 8 hex digits of the public key.
74    pub(crate) fn hostname(&self, context: &Context) -> String {
75        match self {
76            PeerId::Validator(index) => context.committee.authority(*index).hostname.to_string(),
77            PeerId::Observer(node_id) => {
78                let bytes = node_id.to_bytes();
79                format!(
80                    "[Observer]{:02x}{:02x}{:02x}{:02x}",
81                    bytes[0], bytes[1], bytes[2], bytes[3]
82                )
83            }
84        }
85    }
86
87    /// Returns a short label suitable for use in metrics. Does not include
88    /// full public keys to avoid high-cardinality metric labels.
89    pub(crate) fn labelname(&self, context: &Context) -> String {
90        match self {
91            PeerId::Validator(index) => context.committee.authority(*index).hostname.to_string(),
92            PeerId::Observer(_) => "observer".to_string(),
93        }
94    }
95}
96
97// Tonic generated RPC stubs.
98mod tonic_gen {
99    include!(concat!(env!("OUT_DIR"), "/consensus.ConsensusService.rs"));
100    include!(concat!(env!("OUT_DIR"), "/consensus.ObserverService.rs"));
101}
102
103mod clients;
104pub(crate) mod metrics;
105mod metrics_layer;
106#[cfg(all(test, not(msim)))]
107mod network_tests;
108#[cfg(not(msim))]
109pub(crate) mod observer;
110#[cfg(msim)]
111pub mod observer;
112#[cfg(test)]
113pub(crate) mod test_network;
114#[cfg(not(msim))]
115pub(crate) mod tonic_network;
116#[cfg(msim)]
117pub mod tonic_network;
118mod tonic_tls;
119
120/// A stream of serialized filtered blocks returned over the network.
121pub(crate) type BlockStream = Pin<Box<dyn Stream<Item = ExtendedSerializedBlock> + Send>>;
122
123/// Validator network client for communicating with validator peers.
124///
125/// NOTE: the timeout parameters help saving resources at client and potentially server.
126/// But it is up to the server implementation if the timeout is honored.
127/// - To bound server resources, server should implement own timeout for incoming requests.
128#[async_trait]
129pub(crate) trait ValidatorNetworkClient: Send + Sync + Sized + 'static {
130    /// Subscribes to blocks from a peer after last_received round.
131    async fn subscribe_blocks(
132        &self,
133        peer: AuthorityIndex,
134        last_received: Round,
135        timeout: Duration,
136    ) -> ConsensusResult<BlockStream>;
137
138    // TODO: add a parameter for maximum total size of blocks returned.
139    /// Fetches serialized `SignedBlock`s from a peer. It also might return additional ancestor blocks
140    /// of the requested blocks according to the provided `fetch_after_rounds`. The `fetch_after_rounds`
141    /// length should be equal to the committee size. If `fetch_after_rounds` is empty then it will
142    /// be simply ignored.
143    async fn fetch_blocks(
144        &self,
145        peer: AuthorityIndex,
146        block_refs: Vec<BlockRef>,
147        fetch_after_rounds: Vec<Round>,
148        fetch_missing_ancestors: bool,
149        timeout: Duration,
150    ) -> ConsensusResult<Vec<Bytes>>;
151
152    /// Fetches serialized commits in the commit range from a peer.
153    /// Returns a tuple of both the serialized commits, and serialized blocks that contain
154    /// votes certifying the last commit.
155    async fn fetch_commits(
156        &self,
157        peer: AuthorityIndex,
158        commit_range: CommitRange,
159        timeout: Duration,
160    ) -> ConsensusResult<(Vec<Bytes>, Vec<Bytes>)>;
161
162    /// Fetches the latest block from `peer` for the requested `authorities`. The latest blocks
163    /// are returned in the serialised format of `SignedBlocks`. The method can return multiple
164    /// blocks per peer as its possible to have equivocations.
165    async fn fetch_latest_blocks(
166        &self,
167        peer: AuthorityIndex,
168        authorities: Vec<AuthorityIndex>,
169        timeout: Duration,
170    ) -> ConsensusResult<Vec<Bytes>>;
171
172    /// Gets the latest received & accepted rounds of all authorities from the peer.
173    async fn get_latest_rounds(
174        &self,
175        peer: AuthorityIndex,
176        timeout: Duration,
177    ) -> ConsensusResult<(Vec<Round>, Vec<Round>)>;
178
179    /// Sends a serialized SignedBlock to a peer.
180    #[cfg(test)]
181    async fn send_block(
182        &self,
183        peer: AuthorityIndex,
184        block: &VerifiedBlock,
185        timeout: Duration,
186    ) -> ConsensusResult<()>;
187}
188
189/// Validator network service for handling requests from validator peers.
190#[async_trait]
191pub(crate) trait ValidatorNetworkService: Send + Sync + 'static {
192    /// Handles the block sent from the peer via either unicast RPC or subscription stream.
193    /// Peer value can be trusted to be a valid authority index.
194    /// But serialized_block must be verified before its contents are trusted.
195    /// Excluded ancestors are also included as part of an effort to further propagate
196    /// blocks to peers despite the current exclusion.
197    async fn handle_send_block(
198        &self,
199        peer: AuthorityIndex,
200        block: ExtendedSerializedBlock,
201    ) -> ConsensusResult<()>;
202
203    /// Handles the subscription request from the peer.
204    /// A stream of newly proposed blocks is returned to the peer.
205    /// The stream continues until the end of epoch, peer unsubscribes, or a network error / crash
206    /// occurs.
207    async fn handle_subscribe_blocks(
208        &self,
209        peer: AuthorityIndex,
210        last_received: Round,
211    ) -> ConsensusResult<BlockStream>;
212
213    /// Handles the request to fetch blocks by references from the peer.
214    async fn handle_fetch_blocks(
215        &self,
216        peer: AuthorityIndex,
217        block_refs: Vec<BlockRef>,
218        fetch_after_rounds: Vec<Round>,
219        fetch_missing_ancestors: bool,
220    ) -> ConsensusResult<Vec<Bytes>>;
221
222    /// Handles the request to fetch commits by index range from the peer.
223    async fn handle_fetch_commits(
224        &self,
225        peer: AuthorityIndex,
226        commit_range: CommitRange,
227    ) -> ConsensusResult<(Vec<TrustedCommit>, Vec<VerifiedBlock>)>;
228
229    /// Handles the request to fetch the latest block for the provided `authorities`.
230    async fn handle_fetch_latest_blocks(
231        &self,
232        peer: AuthorityIndex,
233        authorities: Vec<AuthorityIndex>,
234    ) -> ConsensusResult<Vec<Bytes>>;
235
236    /// Handles the request to get the latest received & accepted rounds of all authorities.
237    async fn handle_get_latest_rounds(
238        &self,
239        peer: AuthorityIndex,
240    ) -> ConsensusResult<(Vec<Round>, Vec<Round>)>;
241}
242
243/// Handler for randomness round signatures exchanged between validators and
244/// observer nodes via the consensus block stream.
245pub trait RandomnessSignatureHandler: Send + Sync + 'static {
246    /// Called by the observer subscriber for each randomness round signature
247    /// received from the block stream.
248    fn handle_randomness_signature(&self, data: Bytes);
249
250    /// Returns a receiver for broadcast randomness signatures. Called by
251    /// the observer service to merge signatures into the outgoing block stream.
252    fn subscribe_randomness_signatures(&self) -> tokio::sync::broadcast::Receiver<Bytes>;
253}
254
255/// Filter sent by an observer when opening a block stream, restricting which authorities'
256/// blocks the server releases on that stream. Only validated committee members appear here:
257/// the transport layer converts and checks the raw wire representation before the filter
258/// reaches the service.
259#[derive(Clone, Debug, Default, PartialEq, Eq)]
260pub(crate) enum BlockStreamFilter {
261    /// Blocks from all authorities are streamed.
262    #[default]
263    All,
264    /// Only blocks authored by the given committee members are streamed. Must be
265    /// non-empty; the wire format cannot express an empty author set.
266    Authors(BTreeSet<AuthorityIndex>),
267}
268
269/// A single item in the observer block stream, carrying both blocks and auxiliary data.
270pub(crate) struct ObserverStreamItem {
271    pub(crate) blocks: Vec<Bytes>,
272    pub(crate) auxiliary_data: observer::AuxiliaryData,
273}
274
275/// Observer block stream type.
276pub(crate) type ObserverBlockStream =
277    Pin<Box<dyn Stream<Item = ObserverStreamItem> + Send + 'static>>;
278
279/// Observer network service for handling requests from observer nodes.
280/// Unlike ValidatorNetworkService which uses AuthorityIndex, this uses NodeId (NetworkPublicKey)
281/// to identify peers since observers are not part of the committee.
282#[async_trait]
283pub(crate) trait ObserverNetworkService: Send + Sync + 'static {
284    /// Handles a block received from a peer subscription. Used by ObserverSubscriber to process
285    /// blocks streamed from validators or other observers.
286    async fn handle_block(&self, peer: PeerId, block: Bytes) -> ConsensusResult<()>;
287
288    /// Handles the block streaming request from an observer peer.
289    /// Returns a stream of blocks with the highest commit index for each block.
290    /// Blocks with rounds higher than the highest_round_per_authority will be streamed,
291    /// restricted to the authors requested in the filter.
292    async fn handle_stream_blocks(
293        &self,
294        peer: NodeId,
295        highest_round_per_authority: Vec<Round>,
296        filter: BlockStreamFilter,
297    ) -> ConsensusResult<ObserverBlockStream>;
298
299    /// Handles the request to fetch blocks by references from an observer peer.
300    async fn handle_fetch_blocks(
301        &self,
302        peer: NodeId,
303        block_refs: Vec<BlockRef>,
304        fetch_after_rounds: Vec<Round>,
305        fetch_missing_ancestors: bool,
306    ) -> ConsensusResult<Vec<Bytes>>;
307
308    /// Handles the request to fetch commits by index range from an observer peer.
309    /// Returns serialized commits and certifier blocks.
310    async fn handle_fetch_commits(
311        &self,
312        peer: NodeId,
313        commit_range: CommitRange,
314    ) -> ConsensusResult<(Vec<TrustedCommit>, Vec<VerifiedBlock>)>;
315}
316
317/// Observer network client for communicating with validators' observer ports or other observers.
318/// Unlike ValidatorNetworkClient which uses AuthorityIndex, this uses PeerId to identify peers
319/// since the observer server can serve both validators and observer nodes.
320#[async_trait]
321pub(crate) trait ObserverNetworkClient: Send + Sync + Sized + 'static {
322    /// Initiates block streaming with a peer (validator or observer).
323    /// Returns a stream of blocks with the highest commit index.
324    /// Blocks with rounds higher than the highest_round_per_authority will be streamed,
325    /// restricted to the authors requested in the filter.
326    async fn stream_blocks(
327        &self,
328        peer: PeerId,
329        highest_round_per_authority: Vec<Round>,
330        filter: BlockStreamFilter,
331        timeout: Duration,
332    ) -> ConsensusResult<ObserverBlockStream>;
333
334    /// Fetches serialized blocks by references from a peer.
335    async fn fetch_blocks(
336        &self,
337        peer: PeerId,
338        block_refs: Vec<BlockRef>,
339        fetch_after_rounds: Vec<Round>,
340        fetch_missing_ancestors: bool,
341        timeout: Duration,
342    ) -> ConsensusResult<Vec<Bytes>>;
343
344    /// Fetches serialized commits in the commit range from a peer.
345    /// Returns a tuple of both the serialized commits, and serialized blocks that contain
346    /// votes certifying the last commit.
347    async fn fetch_commits(
348        &self,
349        peer: PeerId,
350        commit_range: CommitRange,
351        timeout: Duration,
352    ) -> ConsensusResult<(Vec<Bytes>, Vec<Bytes>)>;
353}
354
355/// An `AuthorityNode` holds a `NetworkManager` until shutdown.
356/// Dropping `NetworkManager` will shutdown the network service.
357pub(crate) trait NetworkManager: Send + Sync {
358    type ValidatorClient: ValidatorNetworkClient;
359    type ObserverClient: ObserverNetworkClient;
360
361    /// Creates a new network manager.
362    fn new(context: Arc<Context>, network_keypair: NetworkKeyPair) -> Self;
363
364    /// Returns the validator network client.
365    fn validator_client(&self) -> Arc<Self::ValidatorClient>;
366
367    /// Returns the observer network client.
368    fn observer_client(&self) -> Arc<Self::ObserverClient>;
369
370    /// Starts the validator network server with the provided service.
371    async fn start_validator_server<V>(&mut self, service: Arc<V>)
372    where
373        V: ValidatorNetworkService;
374
375    /// Starts the observer network server with the provided service.
376    async fn start_observer_server<O>(&mut self, service: Arc<O>)
377    where
378        O: ObserverNetworkService;
379
380    /// Stops the network service.
381    async fn stop(&mut self);
382
383    /// Updates the network address for a peer identified by their authority index.
384    /// If address is None, the override is cleared and the committee address will be used.
385    fn update_peer_address(&self, peer: AuthorityIndex, address: Option<Multiaddr>);
386}
387
388// Re-export the concrete client implementations.
389pub(crate) use clients::{CommitSyncerClient, SynchronizerClient};
390
391/// Serialized block with extended information from the proposing authority.
392#[derive(Clone, PartialEq, Eq, Debug)]
393pub(crate) struct ExtendedSerializedBlock {
394    pub(crate) block: SerializedBlockForm,
395    // Serialized BlockRefs that are excluded from the blocks ancestors.
396    pub(crate) excluded_ancestors: Vec<Vec<u8>>,
397}
398
399/// Wire framing for one block on live subscription streams, carried inside a bytes
400/// field so batch streams can reuse it.
401#[derive(Clone, PartialEq, prost::Message)]
402pub(crate) struct SerializedBlockEnvelope {
403    #[prost(oneof = "SerializedBlockForm", tags = "1, 2")]
404    pub(crate) block: Option<SerializedBlockForm>,
405}
406
407/// Which representation of a block a payload holds.
408#[derive(Clone, PartialEq, Eq, prost::Oneof)]
409pub(crate) enum SerializedBlockForm {
410    /// The block exactly as serialized and signed by its author.
411    #[prost(bytes = "bytes", tag = "1")]
412    Full(Bytes),
413    /// The same block with its ancestor digests stripped, to be rebuilt by the
414    /// receiver. Nothing emits this form yet; the encoder lands with the codec.
415    #[prost(bytes = "bytes", tag = "2")]
416    Slim(Bytes),
417}
418
419impl SerializedBlockEnvelope {
420    pub(crate) fn encode_form(form: SerializedBlockForm) -> Bytes {
421        Self { block: Some(form) }.encode_to_vec().into()
422    }
423
424    pub(crate) fn decode_form(bytes: &[u8]) -> ConsensusResult<SerializedBlockForm> {
425        Self::decode(bytes)
426            .map_err(ConsensusError::MalformedBlockEnvelope)?
427            .block
428            .ok_or_else(|| {
429                ConsensusError::MalformedBlockEnvelope(prost::DecodeError::new(
430                    "empty block envelope",
431                ))
432            })
433    }
434}
435
436impl From<ExtendedBlock> for ExtendedSerializedBlock {
437    fn from(extended_block: ExtendedBlock) -> Self {
438        Self {
439            block: SerializedBlockForm::Full(extended_block.block.serialized().clone()),
440            excluded_ancestors: extended_block
441                .excluded_ancestors
442                .iter()
443                .filter_map(|r| match bcs::to_bytes(r) {
444                    Ok(serialized) => Some(serialized),
445                    Err(e) => {
446                        tracing::debug!("Failed to serialize block ref {:?}: {e:?}", r);
447                        None
448                    }
449                })
450                .collect(),
451        }
452    }
453}
454
455/// Attempts to convert a multiaddr of the form `/[ip4,ip6,dns]/{}/udp/{port}` into
456/// a host:port string.
457pub(crate) fn to_host_port_str(addr: &Multiaddr) -> Result<String, String> {
458    let mut iter = addr.iter();
459
460    match (iter.next(), iter.next()) {
461        (Some(Protocol::Ip4(ipaddr)), Some(Protocol::Udp(port))) => {
462            Ok(format!("{}:{}", ipaddr, port))
463        }
464        (Some(Protocol::Ip6(ipaddr)), Some(Protocol::Udp(port))) => {
465            Ok(format!("{}", SocketAddrV6::new(ipaddr, port, 0, 0)))
466        }
467        (Some(Protocol::Dns(hostname)), Some(Protocol::Udp(port))) => {
468            Ok(format!("{}:{}", hostname, port))
469        }
470
471        _ => Err(format!("unsupported multiaddr: {addr}")),
472    }
473}