Skip to main content

subspace_networking/
constructor.rs

1pub(crate) mod temporary_bans;
2mod transport;
3
4use crate::behavior::persistent_parameters::{KnownPeersRegistry, StubNetworkingParametersManager};
5use crate::behavior::{Behavior, BehaviorConfig};
6use crate::constructor::temporary_bans::TemporaryBans;
7use crate::constructor::transport::build_transport;
8use crate::node::Node;
9use crate::node_runner::{NodeRunner, NodeRunnerConfig};
10use crate::protocols::autonat_wrapper::Config as AutonatWrapperConfig;
11use crate::protocols::request_response::request_response_factory::RequestHandler;
12use crate::protocols::reserved_peers::Config as ReservedPeersConfig;
13use crate::shared::Shared;
14use crate::utils::rate_limiter::RateLimiter;
15use crate::utils::{SubspaceMetrics, strip_peer_id};
16use backon::ExponentialBuilder;
17use futures::channel::mpsc;
18use libp2p::autonat::Config as AutonatConfig;
19use libp2p::connection_limits::ConnectionLimits;
20use libp2p::gossipsub::{
21    Config as GossipsubConfig, ConfigBuilder as GossipsubConfigBuilder,
22    Message as GossipsubMessage, MessageId, ValidationMode,
23};
24use libp2p::identify::Config as IdentifyConfig;
25use libp2p::kad::store::RecordStore;
26use libp2p::kad::{
27    BucketInserts, Config as KademliaConfig, Mode, ProviderRecord, Record, RecordKey, StoreInserts,
28    store,
29};
30use libp2p::metrics::Metrics;
31use libp2p::multiaddr::Protocol;
32use libp2p::yamux::Config as YamuxConfig;
33use libp2p::{Multiaddr, PeerId, StreamProtocol, SwarmBuilder, TransportError, identity};
34use parking_lot::Mutex;
35use prometheus_client::registry::Registry;
36use std::borrow::Cow;
37use std::iter::Empty;
38use std::sync::Arc;
39use std::time::Duration;
40use std::{fmt, io, iter};
41use subspace_core_primitives::hashes;
42use subspace_core_primitives::pieces::Piece;
43use thiserror::Error;
44use tracing::{debug, info};
45
46const DEFAULT_NETWORK_PROTOCOL_VERSION: &str = "dev";
47const KADEMLIA_PROTOCOL: &str = "/subspace/kad/0.1.0";
48const GOSSIPSUB_PROTOCOL_PREFIX: &str = "subspace/gossipsub";
49
50/// Defines max_negotiating_inbound_streams constant for the swarm.
51/// It must be set for large plots.
52const SWARM_MAX_NEGOTIATING_INBOUND_STREAMS: usize = 100000;
53/// How long will connection be allowed to be open without any usage
54const IDLE_CONNECTION_TIMEOUT: Duration = Duration::from_secs(3);
55/// The default maximum established incoming connection number for the swarm.
56const SWARM_MAX_ESTABLISHED_INCOMING_CONNECTIONS: u32 = 100;
57/// The default maximum established incoming connection number for the swarm.
58const SWARM_MAX_ESTABLISHED_OUTGOING_CONNECTIONS: u32 = 100;
59/// The default maximum pending incoming connection number for the swarm.
60const SWARM_MAX_PENDING_INCOMING_CONNECTIONS: u32 = 80;
61/// The default maximum pending incoming connection number for the swarm.
62const SWARM_MAX_PENDING_OUTGOING_CONNECTIONS: u32 = 80;
63const KADEMLIA_QUERY_TIMEOUT: Duration = Duration::from_secs(40);
64const SWARM_MAX_ESTABLISHED_CONNECTIONS_PER_PEER: u32 = 3;
65const MAX_CONCURRENT_STREAMS_PER_CONNECTION: usize = 10;
66// TODO: Consider moving this constant to configuration or removing `Toggle` wrapper when we find a
67//  use-case for gossipsub protocol.
68const ENABLE_GOSSIP_PROTOCOL: bool = false;
69
70const TEMPORARY_BANS_CACHE_SIZE: u32 = 10_000;
71const TEMPORARY_BANS_DEFAULT_BACKOFF_INITIAL_INTERVAL: Duration = Duration::from_secs(5);
72const TEMPORARY_BANS_DEFAULT_BACKOFF_MULTIPLIER: f64 = 1.5;
73const TEMPORARY_BANS_DEFAULT_MAX_INTERVAL: Duration = Duration::from_secs(30 * 60);
74
75/// We pause between reserved peers dialing otherwise we could do multiple dials to offline peers
76/// wasting resources and producing a ton of log records.
77const DIALING_INTERVAL_IN_SECS: Duration = Duration::from_secs(1);
78
79/// Max confidence for autonat protocol. Could affect Kademlia mode change.
80pub(crate) const AUTONAT_MAX_CONFIDENCE: usize = 3;
81/// We set a very long pause before autonat initialization (Duration::Max panics).
82const AUTONAT_SERVER_PROBE_DELAY: Duration = Duration::from_secs(3600 * 24 * 365);
83
84/// Defines Kademlia mode
85#[derive(Clone, Debug)]
86pub enum KademliaMode {
87    /// The Kademlia mode is static for the duration of the application.
88    Static(Mode),
89    /// Kademlia mode will be changed using Autonat protocol when max confidence reached.
90    Dynamic,
91}
92
93impl KademliaMode {
94    /// Returns true if the mode is Dynamic.
95    pub fn is_dynamic(&self) -> bool {
96        matches!(self, Self::Dynamic)
97    }
98
99    /// Returns true if the mode is Static.
100    pub fn is_static(&self) -> bool {
101        matches!(self, Self::Static(..))
102    }
103}
104
105pub(crate) struct DummyRecordStore;
106
107impl RecordStore for DummyRecordStore {
108    type RecordsIter<'a>
109        = Empty<Cow<'a, Record>>
110    where
111        Self: 'a;
112    type ProvidedIter<'a>
113        = Empty<Cow<'a, ProviderRecord>>
114    where
115        Self: 'a;
116
117    #[inline]
118    fn get(&self, _key: &RecordKey) -> Option<Cow<'_, Record>> {
119        // Not supported
120        None
121    }
122
123    #[inline]
124    fn put(&mut self, _record: Record) -> store::Result<()> {
125        // Not supported
126        Ok(())
127    }
128
129    #[inline]
130    fn remove(&mut self, _key: &RecordKey) {
131        // Not supported
132    }
133
134    #[inline]
135    fn records(&self) -> Self::RecordsIter<'_> {
136        // Not supported
137        iter::empty()
138    }
139
140    #[inline]
141    fn add_provider(&mut self, _record: ProviderRecord) -> store::Result<()> {
142        // Not supported
143        Ok(())
144    }
145
146    #[inline]
147    fn providers(&self, _key: &RecordKey) -> Vec<ProviderRecord> {
148        // Not supported
149        Vec::new()
150    }
151
152    #[inline]
153    fn provided(&self) -> Self::ProvidedIter<'_> {
154        // Not supported
155        iter::empty()
156    }
157
158    #[inline]
159    fn remove_provider(&mut self, _key: &RecordKey, _provider: &PeerId) {
160        // Not supported
161    }
162}
163
164/// [`Node`] configuration.
165pub struct Config {
166    /// Identity keypair of a node used for authenticated connections.
167    pub keypair: identity::Keypair,
168    /// List of [`Multiaddr`] on which to listen for incoming connections.
169    pub listen_on: Vec<Multiaddr>,
170    /// Fallback to random port if specified (or default) port is already occupied.
171    pub listen_on_fallback_to_random_port: bool,
172    /// Adds a timeout to the setup and protocol upgrade process for all inbound and outbound
173    /// connections established through the transport.
174    pub timeout: Duration,
175    /// The configuration for the Identify behaviour.
176    pub identify: IdentifyConfig,
177    /// The configuration for the Kademlia behaviour.
178    pub kademlia: KademliaConfig,
179    /// The configuration for the Gossip behaviour.
180    pub gossipsub: Option<GossipsubConfig>,
181    /// Yamux multiplexing configuration.
182    pub yamux_config: YamuxConfig,
183    /// Should non-global addresses be added to the DHT?
184    pub allow_non_global_addresses_in_dht: bool,
185    /// How frequently should random queries be done using Kademlia DHT to populate routing table.
186    pub initial_random_query_interval: Duration,
187    /// A reference to the `NetworkingParametersRegistry` implementation.
188    pub known_peers_registry: Box<dyn KnownPeersRegistry>,
189    /// The configuration for the `RequestResponsesBehaviour` protocol.
190    pub request_response_protocols: Vec<Box<dyn RequestHandler>>,
191    /// Defines set of peers with a permanent connection (and reconnection if necessary).
192    pub reserved_peers: Vec<Multiaddr>,
193    /// Established incoming swarm connection limit.
194    pub max_established_incoming_connections: u32,
195    /// Established outgoing swarm connection limit.
196    pub max_established_outgoing_connections: u32,
197    /// Pending incoming swarm connection limit.
198    pub max_pending_incoming_connections: u32,
199    /// Pending outgoing swarm connection limit.
200    pub max_pending_outgoing_connections: u32,
201    /// How many temporarily banned unreachable peers to keep in memory.
202    pub temporary_bans_cache_size: u32,
203    /// Backoff policy for temporary banning of unreachable peers.
204    pub temporary_ban_backoff: ExponentialBuilder,
205    /// Optional libp2p prometheus metrics. None will disable metrics gathering.
206    pub libp2p_metrics: Option<Metrics>,
207    /// Internal prometheus metrics. None will disable metrics gathering.
208    pub metrics: Option<SubspaceMetrics>,
209    /// Defines protocol version for the network peers. Affects network partition.
210    pub protocol_version: String,
211    /// Addresses to bootstrap Kademlia network
212    pub bootstrap_addresses: Vec<Multiaddr>,
213    /// Kademlia mode. The default value is set to Static(Client). The peer won't add its address
214    /// to other peers` Kademlia routing table. Changing this behaviour implies that a peer can
215    /// provide pieces to others.
216    pub kademlia_mode: KademliaMode,
217    /// Known external addresses to the local peer. The addresses will be added on the swarm start
218    /// and enable peer to notify others about its reachable address.
219    pub external_addresses: Vec<Multiaddr>,
220}
221
222impl fmt::Debug for Config {
223    #[inline]
224    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
225        f.debug_struct("Config").finish()
226    }
227}
228
229/// This default can only be used for `dev` networks.
230/// Other networks should use `Config::new()` to apply the correct prefix to the protocol version.
231impl Default for Config {
232    #[inline]
233    fn default() -> Self {
234        let ed25519_keypair = identity::ed25519::Keypair::generate();
235        let keypair = identity::Keypair::from(ed25519_keypair);
236
237        Self::new(DEFAULT_NETWORK_PROTOCOL_VERSION.to_string(), keypair, None)
238    }
239}
240
241impl Config {
242    /// Creates a new [`Config`].
243    /// Applies a subspace-specific version prefix to the `protocol_version`.
244    pub fn new(
245        protocol_version: String,
246        keypair: identity::Keypair,
247        prometheus_registry: Option<&mut Registry>,
248    ) -> Self {
249        let (libp2p_metrics, metrics) = prometheus_registry
250            .map(|registry| {
251                (
252                    Some(Metrics::new(registry)),
253                    Some(SubspaceMetrics::new(registry)),
254                )
255            })
256            .unwrap_or((None, None));
257
258        let mut kademlia = KademliaConfig::new(
259            StreamProtocol::try_from_owned(KADEMLIA_PROTOCOL.to_owned())
260                .expect("Manual protocol name creation."),
261        );
262        kademlia
263            .set_query_timeout(KADEMLIA_QUERY_TIMEOUT)
264            .disjoint_query_paths(true)
265            .set_max_packet_size(2 * Piece::SIZE)
266            .set_kbucket_inserts(BucketInserts::Manual)
267            .set_record_filtering(StoreInserts::FilterBoth)
268            // We don't use records and providers publication.
269            .set_provider_record_ttl(None)
270            .set_provider_publication_interval(None)
271            .set_record_ttl(None)
272            .set_replication_interval(None);
273
274        // NOTE: Do not call deprecated setters like `set_max_num_streams()` on this config.
275        // They silently downgrade from yamux 0.13 to 0.12, which has a remote DoS vulnerability.
276        let yamux_config = YamuxConfig::default();
277
278        let gossipsub = ENABLE_GOSSIP_PROTOCOL.then(|| {
279            GossipsubConfigBuilder::default()
280                .protocol_id_prefix(GOSSIPSUB_PROTOCOL_PREFIX)
281                // TODO: Do we want message signing?
282                .validation_mode(ValidationMode::None)
283                // To content-address message, we can take the hash of message and use it as an ID.
284                .message_id_fn(|message: &GossipsubMessage| {
285                    MessageId::from(*hashes::blake3_hash(&message.data))
286                })
287                .max_transmit_size(2 * 1024 * 1024) // 2MB
288                .build()
289                .expect("Default config for gossipsub is always correct; qed")
290        });
291
292        let protocol_version = format!("/subspace/2/{protocol_version}");
293        let identify = IdentifyConfig::new(protocol_version.clone(), keypair.public());
294
295        let temporary_ban_backoff = ExponentialBuilder::default()
296            .with_factor(TEMPORARY_BANS_DEFAULT_BACKOFF_MULTIPLIER as f32)
297            .with_min_delay(TEMPORARY_BANS_DEFAULT_BACKOFF_INITIAL_INTERVAL)
298            .with_max_delay(TEMPORARY_BANS_DEFAULT_MAX_INTERVAL)
299            .without_max_times();
300
301        Self {
302            keypair,
303            listen_on: vec![],
304            listen_on_fallback_to_random_port: true,
305            timeout: Duration::from_secs(10),
306            identify,
307            kademlia,
308            gossipsub,
309            allow_non_global_addresses_in_dht: false,
310            initial_random_query_interval: Duration::from_secs(1),
311            known_peers_registry: StubNetworkingParametersManager.boxed(),
312            request_response_protocols: Vec::new(),
313            yamux_config,
314            reserved_peers: Vec::new(),
315            max_established_incoming_connections: SWARM_MAX_ESTABLISHED_INCOMING_CONNECTIONS,
316            max_established_outgoing_connections: SWARM_MAX_ESTABLISHED_OUTGOING_CONNECTIONS,
317            max_pending_incoming_connections: SWARM_MAX_PENDING_INCOMING_CONNECTIONS,
318            max_pending_outgoing_connections: SWARM_MAX_PENDING_OUTGOING_CONNECTIONS,
319            temporary_bans_cache_size: TEMPORARY_BANS_CACHE_SIZE,
320            temporary_ban_backoff,
321            libp2p_metrics,
322            metrics,
323            protocol_version,
324            bootstrap_addresses: Vec::new(),
325            kademlia_mode: KademliaMode::Static(Mode::Client),
326            external_addresses: Vec::new(),
327        }
328    }
329}
330
331/// Errors that might happen during network creation.
332#[derive(Debug, Error)]
333pub enum CreationError {
334    /// Circuit relay client error.
335    #[error("Expected relay server node.")]
336    RelayServerExpected,
337    /// I/O error.
338    #[error("I/O error: {0}")]
339    Io(#[from] io::Error),
340    /// Transport creation error.
341    #[error("Transport creation error: {0}")]
342    // TODO: Restore `#[from] TransportError` once https://github.com/libp2p/rust-libp2p/issues/4824
343    //  is resolved
344    TransportCreationError(Box<dyn std::error::Error + Send + Sync>),
345    /// Transport error when attempting to listen on multiaddr.
346    #[error("Transport error when attempting to listen on multiaddr: {0}")]
347    TransportError(#[from] TransportError<io::Error>),
348}
349
350/// Converts public key from keypair to PeerId.
351/// It serves as the shared PeerId generating algorithm.
352pub fn peer_id(keypair: &identity::Keypair) -> PeerId {
353    keypair.public().to_peer_id()
354}
355
356/// Create a new network node and node runner instances.
357pub fn construct(config: Config) -> Result<(Node, NodeRunner), CreationError> {
358    let Config {
359        keypair,
360        listen_on,
361        listen_on_fallback_to_random_port,
362        timeout,
363        identify,
364        kademlia,
365        gossipsub,
366        yamux_config,
367        allow_non_global_addresses_in_dht,
368        initial_random_query_interval,
369        known_peers_registry,
370        request_response_protocols,
371        reserved_peers,
372        max_established_incoming_connections,
373        max_established_outgoing_connections,
374        max_pending_incoming_connections,
375        max_pending_outgoing_connections,
376        temporary_bans_cache_size,
377        temporary_ban_backoff,
378        libp2p_metrics,
379        metrics,
380        protocol_version,
381        bootstrap_addresses,
382        kademlia_mode,
383        external_addresses,
384    } = config;
385    let local_peer_id = peer_id(&keypair);
386
387    info!(
388        %allow_non_global_addresses_in_dht,
389        peer_id = %local_peer_id,
390        %protocol_version,
391        "DSN instance configured."
392    );
393
394    let connection_limits = ConnectionLimits::default()
395        .with_max_established_per_peer(Some(SWARM_MAX_ESTABLISHED_CONNECTIONS_PER_PEER))
396        .with_max_pending_incoming(Some(max_pending_incoming_connections))
397        .with_max_pending_outgoing(Some(max_pending_outgoing_connections))
398        .with_max_established_incoming(Some(max_established_incoming_connections))
399        .with_max_established_outgoing(Some(max_established_outgoing_connections));
400
401    debug!(?connection_limits, "DSN connection limits set.");
402
403    let autonat_boot_delay = if kademlia_mode.is_static() || !external_addresses.is_empty() {
404        AUTONAT_SERVER_PROBE_DELAY
405    } else {
406        AutonatConfig::default().boot_delay
407    };
408
409    debug!(
410        ?autonat_boot_delay,
411        ?kademlia_mode,
412        ?external_addresses,
413        "Autonat boot delay set."
414    );
415
416    let mut behaviour = Behavior::new(BehaviorConfig {
417        peer_id: local_peer_id,
418        identify,
419        kademlia,
420        gossipsub,
421        request_response_protocols,
422        request_response_max_concurrent_streams: {
423            let max_num_connections = max_established_incoming_connections as usize
424                + max_established_outgoing_connections as usize;
425            max_num_connections * MAX_CONCURRENT_STREAMS_PER_CONNECTION
426        },
427        connection_limits,
428        reserved_peers: ReservedPeersConfig {
429            reserved_peers: reserved_peers.clone(),
430            dialing_interval: DIALING_INTERVAL_IN_SECS,
431        },
432        autonat: AutonatWrapperConfig {
433            inner_config: AutonatConfig {
434                use_connected: true,
435                only_global_ips: !config.allow_non_global_addresses_in_dht,
436                confidence_max: AUTONAT_MAX_CONFIDENCE,
437                boot_delay: autonat_boot_delay,
438                ..Default::default()
439            },
440            local_peer_id,
441            servers: bootstrap_addresses.clone(),
442        },
443    });
444
445    match (kademlia_mode, external_addresses.is_empty()) {
446        (KademliaMode::Static(mode), _) => {
447            behaviour.kademlia.set_mode(Some(mode));
448        }
449        (KademliaMode::Dynamic, false) => {
450            behaviour.kademlia.set_mode(Some(Mode::Server));
451        }
452        _ => {
453            // Autonat will figure it out
454        }
455    };
456
457    let temporary_bans = Arc::new(Mutex::new(TemporaryBans::new(
458        temporary_bans_cache_size,
459        temporary_ban_backoff,
460    )));
461
462    let mut swarm = SwarmBuilder::with_existing_identity(keypair)
463        .with_tokio()
464        .with_other_transport(|keypair| {
465            Ok(build_transport(
466                allow_non_global_addresses_in_dht,
467                keypair,
468                Arc::clone(&temporary_bans),
469                timeout,
470                yamux_config,
471            )?)
472        })
473        .map_err(|error| CreationError::TransportCreationError(error.into()))?
474        .with_behaviour(move |_keypair| Ok(behaviour))
475        .expect("Not fallible; qed")
476        .with_swarm_config(|config| {
477            config
478                .with_max_negotiating_inbound_streams(SWARM_MAX_NEGOTIATING_INBOUND_STREAMS)
479                .with_idle_connection_timeout(IDLE_CONNECTION_TIMEOUT)
480        })
481        .build();
482
483    let is_listening = !listen_on.is_empty();
484
485    // Setup listen_on addresses
486    for mut addr in listen_on {
487        if let Err(error) = swarm.listen_on(addr.clone()) {
488            if !listen_on_fallback_to_random_port {
489                return Err(error.into());
490            }
491
492            let addr_string = addr.to_string();
493            // Listen on random port if specified is already occupied
494            if let Some(Protocol::Tcp(_port)) = addr.pop() {
495                info!("Failed to listen on {addr_string} ({error}), falling back to random port");
496                addr.push(Protocol::Tcp(0));
497                swarm.listen_on(addr)?;
498            }
499        }
500    }
501
502    // Setup external addresses
503    for addr in external_addresses.iter().cloned() {
504        info!("DSN external address added: {addr}");
505        swarm.add_external_address(addr);
506    }
507
508    // Create final structs
509    let (command_sender, command_receiver) = mpsc::channel(1);
510
511    let rate_limiter = RateLimiter::new(
512        max_established_outgoing_connections,
513        max_pending_outgoing_connections,
514    );
515
516    let shared = Arc::new(Shared::new(local_peer_id, command_sender, rate_limiter));
517    let shared_weak = Arc::downgrade(&shared);
518
519    let node = Node::new(shared);
520    let node_runner = NodeRunner::new(NodeRunnerConfig {
521        allow_non_global_addresses_in_dht,
522        is_listening,
523        command_receiver,
524        swarm,
525        shared_weak,
526        next_random_query_interval: initial_random_query_interval,
527        known_peers_registry,
528        reserved_peers: strip_peer_id(reserved_peers).into_iter().collect(),
529        temporary_bans,
530        libp2p_metrics,
531        metrics,
532        protocol_version,
533        bootstrap_addresses,
534    });
535
536    Ok((node, node_runner))
537}