extern crate alloc;
#[cfg(feature = "portable-mdns")]
use alloc::collections::BTreeMap;
#[cfg(all(feature = "smoltcp", feature = "pubsub"))]
use alloc::collections::VecDeque;
#[cfg(feature = "portable-mdns")]
use alloc::string::ToString;
use alloc::{string::String, vec::Vec};
pub use minip2p_core::{Multiaddr, PeerAddr, PeerId, Protocol, TransportKind};
pub use minip2p_identity::Ed25519Keypair;
pub use minip2p_platform::{Deadline as PollDeadline, EntropySource, Now, SharedEntropy};
pub use minip2p_swarm::{
DriverError, IdentifyMessage, SwarmBuilder, SwarmError, SwarmEvent, SwarmRuntime,
};
pub use minip2p_transport::{ConnectionId, StreamId, Transport, TransportError};
#[cfg(feature = "portable-mdns")]
pub use minip2p_discovery::{
BeaconConfig, DiscoveryEvent, DiscoverySource, KnownPeer, PeerDiscoveryConfig,
};
#[cfg(feature = "portable-mdns")]
pub use minip2p_mdns::{MdnsConfig, MdnsConfigError, MdnsError, MdnsEvent, MdnsIo};
#[cfg(feature = "smoltcp")]
pub use minip2p_mdns::{SmoltcpMdnsConfig, SmoltcpMdnsIo};
#[cfg(all(feature = "portable-autonat", not(feature = "nat")))]
pub use minip2p_nat::{
ConnectId, NatConfig, NatEvent, Path, ReachabilityState, ReservationInfo, ReservationPolicy,
};
#[cfg(all(feature = "pubsub", not(feature = "std")))]
pub use minip2p_pubsub::{
FloodsubConfig, GossipsubConfig, PublishError, PubsubConfig, PubsubConfigError, PubsubEvent,
TopicError,
};
#[cfg(feature = "smoltcp")]
pub use minip2p_tcp::{SmoltcpConfig, SmoltcpStack, SmoltcpTcpProvider, smoltcp};
#[cfg(feature = "tcp")]
pub use minip2p_tcp::{TcpConfig, TcpProvider, TcpTransport};
#[cfg(feature = "portable-autonat")]
mod nat;
#[cfg(not(feature = "std"))]
pub struct Endpoint;
#[cfg(not(feature = "std"))]
impl Endpoint {
pub fn portable<E: EntropySource>(
identity: &Ed25519Keypair,
entropy: E,
) -> PortableEndpointBuilder<E> {
PortableEndpointBuilder::new(identity, entropy)
}
}
pub struct PortableEndpoint<T: Transport, E: EntropySource> {
runtime: SwarmRuntime<T, E>,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct PortableEndpointStats {
pub connected_peers: usize,
pub ready_peers: usize,
pub listen_addresses: Vec<Multiaddr>,
}
impl<T: Transport, E: EntropySource> PortableEndpoint<T, E> {
pub fn peer_id(&self) -> &PeerId {
self.runtime.local_peer_id()
}
pub fn runtime(&self) -> &SwarmRuntime<T, E> {
&self.runtime
}
pub fn runtime_mut(&mut self) -> &mut SwarmRuntime<T, E> {
&mut self.runtime
}
pub fn into_runtime(self) -> SwarmRuntime<T, E> {
self.runtime
}
pub fn listen(&mut self, address: &Multiaddr) -> Result<Multiaddr, DriverError> {
self.runtime.listen(address)
}
pub fn listen_all(&mut self) -> Result<Vec<PeerAddr>, DriverError> {
self.runtime.listen_on_bound_addrs()
}
pub fn dial(&mut self, address: &PeerAddr) -> Result<ConnectionId, DriverError> {
self.runtime.dial(address)
}
pub fn connected_peers(&self) -> Vec<PeerId> {
self.runtime.connected_peers()
}
pub fn is_peer_ready(&self, peer_id: &PeerId) -> bool {
self.runtime.is_peer_ready(peer_id)
}
pub fn peer_info(&self, peer_id: &PeerId) -> Option<&IdentifyMessage> {
self.runtime.peer_info(peer_id)
}
pub fn set_external_addresses(&mut self, addresses: Vec<Multiaddr>) {
self.runtime.set_external_addresses(addresses);
}
pub fn ping(&mut self, peer_id: &PeerId, now: Now) -> Result<(), DriverError> {
self.runtime.ping(peer_id, now.monotonic_ms)
}
pub fn disconnect(&mut self, peer_id: &PeerId, now: Now) -> Result<(), DriverError> {
self.runtime.disconnect(peer_id, now.monotonic_ms)
}
pub fn open_stream(
&mut self,
peer_id: &PeerId,
protocol_id: &str,
now: Now,
) -> Result<StreamId, DriverError> {
self.runtime
.open_stream(peer_id, protocol_id, now.monotonic_ms)
}
pub fn send_stream(
&mut self,
peer_id: &PeerId,
stream_id: StreamId,
data: Vec<u8>,
now: Now,
) -> Result<(), DriverError> {
self.runtime
.send_stream(peer_id, stream_id, data, now.monotonic_ms)
}
pub fn close_stream_write(
&mut self,
peer_id: &PeerId,
stream_id: StreamId,
now: Now,
) -> Result<(), DriverError> {
self.runtime
.close_stream_write(peer_id, stream_id, now.monotonic_ms)
}
pub fn reset_stream(
&mut self,
peer_id: &PeerId,
stream_id: StreamId,
now: Now,
) -> Result<(), DriverError> {
self.runtime
.reset_stream(peer_id, stream_id, now.monotonic_ms)
}
pub fn abandon_stream(
&mut self,
peer_id: &PeerId,
stream_id: StreamId,
now: Now,
) -> Result<(), DriverError> {
self.runtime
.abandon_stream(peer_id, stream_id, now.monotonic_ms)
}
pub fn stats(&self) -> PortableEndpointStats {
let connected = self.runtime.connected_peers();
PortableEndpointStats {
ready_peers: connected
.iter()
.filter(|peer_id| self.runtime.is_peer_ready(peer_id))
.count(),
connected_peers: connected.len(),
listen_addresses: self.runtime.transport().local_addresses(),
}
}
pub fn poll(&mut self, now: Now) -> Result<alloc::vec::Vec<SwarmEvent>, DriverError> {
self.runtime.poll(now)
}
pub fn next_deadline(&self, now: Now) -> Option<PollDeadline> {
self.runtime.next_deadline(now)
}
pub fn shutdown(mut self, now: Now) -> Result<Vec<SwarmEvent>, DriverError> {
let mut first_error = None;
for peer_id in self.runtime.connected_peers() {
if let Err(error) = self.runtime.disconnect(&peer_id, now.monotonic_ms)
&& first_error.is_none()
{
first_error = Some(error);
}
}
let events = self.runtime.poll(now);
match first_error {
Some(error) => Err(error),
None => events,
}
}
#[cfg(feature = "portable-mdns")]
pub fn with_mdns<I: MdnsIo>(
self,
io: I,
config: MdnsConfig,
seed: [u8; 32],
) -> Result<PortableMdnsEndpoint<T, E, I>, PortableMdnsConfigError> {
self.with_mdns_config(io, config, PeerDiscoveryConfig::default(), seed)
}
#[cfg(feature = "portable-mdns")]
pub fn with_mdns_config<I: MdnsIo>(
self,
io: I,
mdns_config: MdnsConfig,
discovery_config: PeerDiscoveryConfig,
seed: [u8; 32],
) -> Result<PortableMdnsEndpoint<T, E, I>, PortableMdnsConfigError> {
let agent =
minip2p_mdns::MdnsAgent::new(self.peer_id().clone(), mdns_config.clone(), seed)?;
let local_peer_id = self.peer_id().clone();
Ok(PortableMdnsEndpoint {
endpoint: self,
mdns: minip2p_mdns::MdnsDriver::new(agent, io, &mdns_config),
discovery: minip2p_discovery::PeerDiscoveryAgent::new(local_peer_id, discovery_config)?,
active_dials: BTreeMap::new(),
})
}
}
#[cfg(feature = "portable-mdns")]
#[derive(Debug)]
pub enum PortableMdnsConfigError {
Mdns(MdnsConfigError),
Discovery(minip2p_discovery::DiscoveryConfigError),
}
#[cfg(feature = "portable-mdns")]
impl From<MdnsConfigError> for PortableMdnsConfigError {
fn from(error: MdnsConfigError) -> Self {
Self::Mdns(error)
}
}
#[cfg(feature = "portable-mdns")]
impl From<minip2p_discovery::DiscoveryConfigError> for PortableMdnsConfigError {
fn from(error: minip2p_discovery::DiscoveryConfigError) -> Self {
Self::Discovery(error)
}
}
#[cfg(feature = "portable-mdns")]
impl core::fmt::Display for PortableMdnsConfigError {
fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
match self {
Self::Mdns(error) => error.fmt(formatter),
Self::Discovery(error) => error.fmt(formatter),
}
}
}
#[cfg(feature = "portable-mdns")]
#[derive(Clone, Debug)]
pub enum PortableMdnsEvent {
Endpoint(SwarmEvent),
Discovery(DiscoveryEvent),
}
#[cfg(feature = "portable-mdns")]
pub struct PortableMdnsEndpoint<T: Transport, E: EntropySource, I: MdnsIo> {
endpoint: PortableEndpoint<T, E>,
mdns: minip2p_mdns::MdnsDriver<I>,
discovery: minip2p_discovery::PeerDiscoveryAgent,
active_dials: BTreeMap<PeerId, ConnectionId>,
}
#[cfg(feature = "portable-mdns")]
impl<T: Transport, E: EntropySource, I: MdnsIo> core::ops::Deref for PortableMdnsEndpoint<T, E, I> {
type Target = PortableEndpoint<T, E>;
fn deref(&self) -> &Self::Target {
&self.endpoint
}
}
#[cfg(feature = "portable-mdns")]
impl<T: Transport, E: EntropySource, I: MdnsIo> core::ops::DerefMut
for PortableMdnsEndpoint<T, E, I>
{
fn deref_mut(&mut self) -> &mut Self::Target {
&mut self.endpoint
}
}
#[cfg(feature = "portable-mdns")]
impl<T: Transport, E: EntropySource, I: MdnsIo> PortableMdnsEndpoint<T, E, I> {
pub fn known_peers(&self) -> Vec<KnownPeer> {
self.discovery.known_peers()
}
pub fn mdns_io(&self) -> Option<&I> {
self.mdns.io()
}
pub fn mdns_io_mut(&mut self) -> Option<&mut I> {
self.mdns.io_mut()
}
pub fn poll(&mut self, now: Now) -> Result<Vec<PortableMdnsEvent>, PortableMdnsError> {
let mut events = Vec::new();
self.poll_endpoint(now, &mut events)?;
let local_addrs = self.endpoint.stats().listen_addresses;
self.mdns.tick(now.monotonic_ms, &local_addrs)?;
self.poll_endpoint(now, &mut events)?;
while let Some(event) = self.mdns.poll_event() {
match event {
MdnsEvent::PeerObserved { peer, addrs } => {
self.discovery.observe_mdns(peer, addrs, now.monotonic_ms);
}
MdnsEvent::ProtocolViolation { peer, reason } => {
self.discovery.report_violation(
peer,
minip2p_discovery::DiscoverySource::Mdns,
&reason,
);
}
}
}
self.discovery.handle_tick(now.monotonic_ms);
self.drive_discovery_actions(now)?;
while let Some(event) = self.discovery.poll_event() {
events.push(PortableMdnsEvent::Discovery(event));
}
Ok(events)
}
fn poll_endpoint(
&mut self,
now: Now,
events: &mut Vec<PortableMdnsEvent>,
) -> Result<(), DriverError> {
for event in self.endpoint.poll(now)? {
match &event {
SwarmEvent::ConnectionEstablished { peer_id, .. } => {
self.active_dials.remove(peer_id);
self.discovery.dial_succeeded(peer_id, now.monotonic_ms);
self.discovery.peer_connected(peer_id, now.monotonic_ms);
}
SwarmEvent::ConnectionClosed { peer_id, .. } => {
self.discovery.peer_disconnected(peer_id, now.monotonic_ms);
}
SwarmEvent::Error(error) => {
if let Some(peer_id) = remove_failed_autodial(
&mut self.active_dials,
error.peer_id.as_ref(),
error.conn_id,
) {
self.discovery
.dial_failed(&peer_id, &error.detail, now.monotonic_ms);
}
}
_ => {}
}
events.push(PortableMdnsEvent::Endpoint(event));
}
Ok(())
}
pub fn next_deadline(&self, now: Now) -> Option<PollDeadline> {
let mdns = self
.mdns
.next_timeout(now.monotonic_ms)
.map(|timeout| PollDeadline::from_millis(now.monotonic_ms.saturating_add(timeout)));
let discovery = self
.discovery
.next_timeout(now.monotonic_ms)
.map(|timeout| PollDeadline::from_millis(now.monotonic_ms.saturating_add(timeout)));
PollDeadline::earliest_opt(
PollDeadline::earliest_opt(self.endpoint.next_deadline(now), mdns),
discovery,
)
}
pub fn shutdown(mut self, now: Now) -> Result<Vec<PortableMdnsEvent>, PortableMdnsError> {
self.mdns.shutdown(now.monotonic_ms)?;
Ok(self
.endpoint
.shutdown(now)?
.into_iter()
.map(PortableMdnsEvent::Endpoint)
.collect())
}
fn drive_discovery_actions(&mut self, now: Now) -> Result<(), PortableMdnsError> {
while let Some(action) = self.discovery.poll_action() {
match action {
minip2p_discovery::DiscoveryAction::Dial { peer, addrs, .. } => {
let mut started = None;
let mut last_error = None;
for address in addrs {
let Ok(target) = PeerAddr::new(address, peer.clone()) else {
continue;
};
match self.endpoint.dial(&target) {
Ok(connection) => {
started = Some(connection);
break;
}
Err(error) => last_error = Some(error),
}
}
if let Some(connection) = started {
self.active_dials.insert(peer, connection);
} else {
let reason = last_error
.map(|error| error.to_string())
.unwrap_or_else(|| "mDNS supplied no dialable address".into());
self.discovery.dial_failed(&peer, &reason, now.monotonic_ms);
}
}
minip2p_discovery::DiscoveryAction::CancelDial { peer } => {
if let Some(connection) = self.active_dials.remove(&peer) {
let close = self
.endpoint
.runtime_mut()
.transport_mut()
.close(connection);
match close {
Ok(())
| Err(
TransportError::ConnectionNotFound { .. }
| TransportError::InvalidState { .. },
) => {}
Err(error) => return Err(DriverError::from(error).into()),
}
}
}
}
}
Ok(())
}
}
#[cfg(feature = "portable-mdns")]
fn remove_failed_autodial(
active: &mut BTreeMap<PeerId, ConnectionId>,
peer_id: Option<&PeerId>,
connection: Option<ConnectionId>,
) -> Option<PeerId> {
let peer = peer_id
.filter(|peer| active.contains_key(*peer))
.cloned()
.or_else(|| {
connection.and_then(|failed| {
active
.iter()
.find_map(|(peer, current)| (*current == failed).then(|| peer.clone()))
})
})?;
active.remove(&peer);
Some(peer)
}
#[cfg(feature = "portable-mdns")]
#[derive(Debug)]
pub enum PortableMdnsError {
Endpoint(DriverError),
Mdns(MdnsError),
}
#[cfg(feature = "portable-mdns")]
impl From<DriverError> for PortableMdnsError {
fn from(error: DriverError) -> Self {
Self::Endpoint(error)
}
}
#[cfg(feature = "portable-mdns")]
impl From<MdnsError> for PortableMdnsError {
fn from(error: MdnsError) -> Self {
Self::Mdns(error)
}
}
#[cfg(feature = "portable-mdns")]
impl core::fmt::Display for PortableMdnsError {
fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
match self {
Self::Endpoint(error) => error.fmt(formatter),
Self::Mdns(error) => error.fmt(formatter),
}
}
}
pub struct PortableEndpointBuilder<E> {
swarm: SwarmBuilder,
#[cfg(feature = "smoltcp")]
identity: Ed25519Keypair,
entropy: E,
}
impl<E: EntropySource> PortableEndpointBuilder<E> {
pub fn new(identity: &Ed25519Keypair, entropy: E) -> Self {
Self {
swarm: SwarmBuilder::new(identity),
#[cfg(feature = "smoltcp")]
identity: identity.clone(),
entropy,
}
}
pub fn agent_version(mut self, value: impl Into<String>) -> Self {
self.swarm = self.swarm.agent_version(value);
self
}
pub fn protocol(mut self, protocol_id: impl Into<String>) -> Self {
self.swarm = self.swarm.protocol(protocol_id);
self
}
pub fn build<T: Transport>(self, transport: T) -> Result<PortableEndpoint<T, E>, SwarmError> {
Ok(PortableEndpoint {
runtime: self.swarm.build_runtime(transport, self.entropy)?,
})
}
#[cfg(feature = "smoltcp")]
pub fn smoltcp<D: smoltcp::phy::Device>(
self,
stack: SmoltcpStack<D>,
) -> SmoltcpEndpointBuilder<D, E> {
SmoltcpEndpointBuilder {
swarm: self.swarm,
identity: self.identity,
entropy: SharedEntropy::new(self.entropy),
stack,
tcp_config: TcpConfig::default(),
provider_config: SmoltcpConfig::default(),
listens: Vec::new(),
mdns: None,
#[cfg(feature = "pubsub")]
pubsub: None,
#[cfg(feature = "pubsub")]
beacon: None,
#[cfg(feature = "portable-autonat")]
nat_config: None,
discovery: PeerDiscoveryConfig::default(),
}
}
}
#[cfg(feature = "smoltcp")]
struct EmbeddedMdnsConfig {
mdns: MdnsConfig,
carrier: SmoltcpMdnsConfig,
}
#[cfg(feature = "smoltcp")]
pub struct SmoltcpEndpointBuilder<D: smoltcp::phy::Device, E: EntropySource> {
swarm: SwarmBuilder,
identity: Ed25519Keypair,
entropy: SharedEntropy<E>,
stack: SmoltcpStack<D>,
tcp_config: TcpConfig,
provider_config: SmoltcpConfig,
listens: Vec<String>,
mdns: Option<EmbeddedMdnsConfig>,
#[cfg(feature = "pubsub")]
pubsub: Option<PubsubConfig>,
#[cfg(feature = "pubsub")]
beacon: Option<BeaconConfig>,
#[cfg(feature = "portable-autonat")]
nat_config: Option<minip2p_nat::NatConfig>,
discovery: PeerDiscoveryConfig,
}
#[cfg(feature = "smoltcp")]
type SmoltcpTcpTransport<D, E> = TcpTransport<SmoltcpTcpProvider<D>, SharedEntropy<E>>;
#[cfg(all(feature = "smoltcp", not(feature = "portable-relay")))]
type SmoltcpComposedTransport<D, E> = SmoltcpTcpTransport<D, E>;
#[cfg(feature = "portable-relay")]
type SmoltcpComposedTransport<D, E> =
minip2p_circuit::CircuitTransport<SmoltcpTcpTransport<D, E>, SharedEntropy<E>>;
#[cfg(feature = "smoltcp")]
pub struct SmoltcpEndpoint<D: smoltcp::phy::Device, E: EntropySource> {
endpoint: PortableEndpoint<SmoltcpComposedTransport<D, E>, SharedEntropy<E>>,
mdns: Option<minip2p_mdns::MdnsDriver<SmoltcpMdnsIo<D>>>,
#[cfg(feature = "pubsub")]
pubsub: Option<minip2p_pubsub::PubsubAgent>,
#[cfg(feature = "pubsub")]
pending_pubsub_events: VecDeque<PubsubEvent>,
#[cfg(feature = "pubsub")]
beacon: Option<minip2p_discovery::BeaconAgent>,
#[cfg(feature = "portable-autonat")]
nat: Option<nat::PortableNatDriver>,
discovery: Option<minip2p_discovery::PeerDiscoveryAgent>,
active_dials: BTreeMap<PeerId, ConnectionId>,
}
#[cfg(feature = "smoltcp")]
#[derive(Clone, Debug)]
pub enum SmoltcpEvent {
Endpoint(SwarmEvent),
#[cfg(feature = "pubsub")]
Pubsub(PubsubEvent),
#[cfg(feature = "portable-autonat")]
Nat(minip2p_nat::NatEvent),
Discovery(DiscoveryEvent),
}
#[cfg(feature = "smoltcp")]
impl<D: smoltcp::phy::Device, E: EntropySource> core::ops::Deref for SmoltcpEndpoint<D, E> {
type Target = PortableEndpoint<SmoltcpComposedTransport<D, E>, SharedEntropy<E>>;
fn deref(&self) -> &Self::Target {
&self.endpoint
}
}
#[cfg(feature = "smoltcp")]
impl<D: smoltcp::phy::Device, E: EntropySource> core::ops::DerefMut for SmoltcpEndpoint<D, E> {
fn deref_mut(&mut self) -> &mut Self::Target {
&mut self.endpoint
}
}
#[cfg(feature = "smoltcp")]
impl<D: smoltcp::phy::Device, E: EntropySource> SmoltcpEndpoint<D, E> {
#[cfg(feature = "portable-relay")]
pub fn connect_relay(
&mut self,
peer: &PeerId,
now: Now,
) -> Result<minip2p_nat::ConnectId, SmoltcpRelayError> {
let nat = self.nat.as_mut().ok_or(SmoltcpRelayError::NotEnabled)?;
if !nat.relay_enabled() {
return Err(SmoltcpRelayError::NotEnabled);
}
let id = nat.connect(peer.clone(), now);
nat.pump(&mut self.endpoint, now);
Ok(id)
}
#[cfg(feature = "portable-relay")]
pub fn cancel_connect(&mut self, id: minip2p_nat::ConnectId, now: Now) {
if let Some(nat) = self.nat.as_mut() {
nat.cancel(id, now);
nat.pump(&mut self.endpoint, now);
}
}
#[cfg(feature = "portable-relay")]
pub fn path(&self, peer: &PeerId) -> Option<minip2p_nat::Path> {
self.nat.as_ref().and_then(|nat| nat.path(peer))
}
#[cfg(feature = "portable-relay")]
pub fn active_reservation(&self) -> Option<minip2p_nat::ReservationInfo> {
self.nat
.as_ref()
.and_then(nat::PortableNatDriver::active_reservation)
}
#[cfg(feature = "portable-autonat")]
pub fn reachability(&self) -> minip2p_nat::ReachabilityState {
self.nat
.as_ref()
.map(|nat| nat.agent.reachability())
.unwrap_or_default()
}
#[cfg(feature = "pubsub")]
pub fn subscribe(&mut self, topic: &str, now: Now) -> Result<bool, SmoltcpPubsubError> {
let agent = self.pubsub.as_mut().ok_or(SmoltcpPubsubError::NotEnabled)?;
let subscribed = agent.subscribe(topic, now.monotonic_ms)?;
pump_embedded_pubsub(agent, &mut self.endpoint, now);
collect_embedded_pubsub_events(agent, &mut self.pending_pubsub_events);
Ok(subscribed)
}
#[cfg(feature = "pubsub")]
pub fn unsubscribe(&mut self, topic: &str, now: Now) -> Result<bool, SmoltcpPubsubError> {
if self
.beacon
.as_ref()
.is_some_and(|beacon| beacon.topic() == topic)
{
return Err(SmoltcpPubsubError::DiscoveryTopicReserved);
}
let agent = self.pubsub.as_mut().ok_or(SmoltcpPubsubError::NotEnabled)?;
let removed = agent.unsubscribe(topic, now.monotonic_ms);
pump_embedded_pubsub(agent, &mut self.endpoint, now);
collect_embedded_pubsub_events(agent, &mut self.pending_pubsub_events);
Ok(removed)
}
#[cfg(feature = "pubsub")]
pub fn publish(
&mut self,
topic: &str,
data: impl Into<Vec<u8>>,
now: Now,
) -> Result<(), SmoltcpPubsubError> {
if self
.beacon
.as_ref()
.is_some_and(|beacon| beacon.topic() == topic)
{
return Err(SmoltcpPubsubError::DiscoveryTopicReserved);
}
let agent = self.pubsub.as_mut().ok_or(SmoltcpPubsubError::NotEnabled)?;
agent.publish(topic, data.into(), now.monotonic_ms)?;
pump_embedded_pubsub(agent, &mut self.endpoint, now);
collect_embedded_pubsub_events(agent, &mut self.pending_pubsub_events);
Ok(())
}
pub fn known_peers(&self) -> Vec<KnownPeer> {
self.discovery
.as_ref()
.map(minip2p_discovery::PeerDiscoveryAgent::known_peers)
.unwrap_or_default()
}
pub fn poll(&mut self, now: Now) -> Result<Vec<SmoltcpEvent>, SmoltcpDriveError> {
let mut output = Vec::new();
self.poll_swarm(now, &mut output)?;
if let Some(mdns) = self.mdns.as_mut() {
mdns.tick(now.monotonic_ms, &self.endpoint.stats().listen_addresses)?;
}
if self.mdns.is_some() {
self.poll_swarm(now, &mut output)?;
}
if let Some(mdns) = self.mdns.as_mut() {
while let Some(event) = mdns.poll_event() {
if let Some(discovery) = self.discovery.as_mut() {
match event {
MdnsEvent::PeerObserved { peer, addrs } => {
discovery.observe_mdns(peer, addrs, now.monotonic_ms);
}
MdnsEvent::ProtocolViolation { peer, reason } => discovery
.report_violation(
peer,
minip2p_discovery::DiscoverySource::Mdns,
&reason,
),
}
}
}
}
#[cfg(feature = "pubsub")]
if let Some(pubsub) = self.pubsub.as_mut() {
if pubsub.next_timeout(now.monotonic_ms) == Some(0) {
pubsub.handle_tick(now.monotonic_ms);
}
pump_embedded_pubsub(pubsub, &mut self.endpoint, now);
}
#[cfg(feature = "pubsub")]
self.drive_beacon(now, &mut output);
#[cfg(feature = "portable-autonat")]
if let Some(nat) = self.nat.as_mut() {
nat.tick(&mut self.endpoint, now);
while let Some(event) = nat.events.pop_front() {
output.push(SmoltcpEvent::Nat(event));
}
}
if let Some(discovery) = self.discovery.as_mut()
&& discovery.next_timeout(now.monotonic_ms) == Some(0)
{
discovery.handle_tick(now.monotonic_ms);
}
self.drive_discovery_actions(now)?;
if let Some(discovery) = self.discovery.as_mut() {
while let Some(event) = discovery.poll_event() {
output.push(SmoltcpEvent::Discovery(event));
}
}
Ok(output)
}
fn poll_swarm(&mut self, now: Now, output: &mut Vec<SmoltcpEvent>) -> Result<(), DriverError> {
for event in self.endpoint.poll(now)? {
if let Some(discovery) = self.discovery.as_mut() {
match &event {
SwarmEvent::ConnectionEstablished { peer_id, .. } => {
self.active_dials.remove(peer_id);
discovery.dial_succeeded(peer_id, now.monotonic_ms);
discovery.peer_connected(peer_id, now.monotonic_ms);
}
SwarmEvent::ConnectionClosed { peer_id, .. } => {
discovery.peer_disconnected(peer_id, now.monotonic_ms);
}
SwarmEvent::Error(error) => {
if let Some(peer) = remove_failed_autodial(
&mut self.active_dials,
error.peer_id.as_ref(),
error.conn_id,
) {
discovery.dial_failed(&peer, &error.detail, now.monotonic_ms);
}
}
_ => {}
}
}
#[cfg(feature = "portable-autonat")]
let claimed = self
.nat
.as_mut()
.is_some_and(|nat| nat.ingest(&event, &mut self.endpoint, now));
#[cfg(not(feature = "portable-autonat"))]
let claimed = false;
#[cfg(feature = "pubsub")]
let claimed = claimed
|| self
.pubsub
.as_mut()
.is_some_and(|agent| agent.handle_event(&event, now.monotonic_ms));
#[cfg(feature = "pubsub")]
if let Some(agent) = self.pubsub.as_mut() {
pump_embedded_pubsub(agent, &mut self.endpoint, now);
}
if !claimed {
output.push(SmoltcpEvent::Endpoint(event));
}
}
Ok(())
}
#[cfg(feature = "pubsub")]
fn drive_beacon(&mut self, now: Now, output: &mut Vec<SmoltcpEvent>) {
let Some(beacon) = self.beacon.as_mut() else {
if let Some(pubsub) = self.pubsub.as_mut() {
collect_embedded_pubsub_events(pubsub, &mut self.pending_pubsub_events);
}
output.extend(
self.pending_pubsub_events
.drain(..)
.map(SmoltcpEvent::Pubsub),
);
return;
};
beacon.set_local_addrs(&self.endpoint.stats().listen_addresses, now.monotonic_ms);
if beacon.next_timeout(now.monotonic_ms) == Some(0) {
beacon.handle_tick(now.monotonic_ms);
}
if let Some(pubsub) = self.pubsub.as_mut() {
collect_embedded_pubsub_events(pubsub, &mut self.pending_pubsub_events);
while let Some(event) = self.pending_pubsub_events.pop_front() {
let consumed = match &event {
PubsubEvent::Message {
from,
topics,
data,
signed,
..
} if topics.iter().any(|topic| topic == beacon.topic()) => {
beacon.handle_beacon(from, data, *signed);
true
}
PubsubEvent::PeerSubscribed { topic, .. }
| PubsubEvent::PeerUnsubscribed { topic, .. }
if topic == beacon.topic() =>
{
true
}
_ => false,
};
if !consumed {
output.push(SmoltcpEvent::Pubsub(event));
}
}
while let Some(minip2p_discovery::BeaconAction::PublishBeacon { topic, payload }) =
beacon.poll_action()
{
if pubsub.publish(&topic, payload, now.monotonic_ms).is_ok() {
pump_embedded_pubsub(pubsub, &mut self.endpoint, now);
}
}
collect_embedded_pubsub_events(pubsub, &mut self.pending_pubsub_events);
while let Some(event) = self.pending_pubsub_events.pop_front() {
output.push(SmoltcpEvent::Pubsub(event));
}
}
if let Some(discovery) = self.discovery.as_mut() {
while let Some(event) = beacon.poll_event() {
match event {
minip2p_discovery::BeaconEvent::Observation(observation) => {
discovery.observe_beacon(observation, now.monotonic_ms);
}
minip2p_discovery::BeaconEvent::ProtocolViolation { peer, reason } => {
discovery.report_violation(
Some(peer),
minip2p_discovery::DiscoverySource::SignedBeacon,
&reason,
);
}
}
}
}
}
fn drive_discovery_actions(&mut self, now: Now) -> Result<(), DriverError> {
let Some(discovery) = self.discovery.as_mut() else {
return Ok(());
};
while let Some(action) = discovery.poll_action() {
match action {
minip2p_discovery::DiscoveryAction::Dial { peer, addrs, .. } => {
let mut started = None;
let mut last_error = None;
for address in addrs {
let Ok(target) = PeerAddr::new(address, peer.clone()) else {
continue;
};
match self.endpoint.dial(&target) {
Ok(connection) => {
started = Some(connection);
break;
}
Err(error) => last_error = Some(error),
}
}
if let Some(connection) = started {
self.active_dials.insert(peer, connection);
} else {
let reason = last_error
.map(|error| error.to_string())
.unwrap_or_else(|| "discovery supplied no dialable address".into());
discovery.dial_failed(&peer, &reason, now.monotonic_ms);
}
}
minip2p_discovery::DiscoveryAction::CancelDial { peer } => {
if let Some(connection) = self.active_dials.remove(&peer) {
match self
.endpoint
.runtime_mut()
.transport_mut()
.close(connection)
{
Ok(())
| Err(
TransportError::ConnectionNotFound { .. }
| TransportError::InvalidState { .. },
) => {}
Err(error) => return Err(error.into()),
}
}
}
}
}
Ok(())
}
pub fn next_deadline(&self, now: Now) -> Option<PollDeadline> {
let mut deadline = self.endpoint.next_deadline(now);
let mut timeouts = Vec::new();
timeouts.extend(
self.mdns
.as_ref()
.and_then(|agent| agent.next_timeout(now.monotonic_ms)),
);
#[cfg(feature = "pubsub")]
{
if !self.pending_pubsub_events.is_empty() {
deadline = PollDeadline::earliest_opt(deadline, Some(PollDeadline::IMMEDIATE));
}
timeouts.extend(
self.pubsub
.as_ref()
.and_then(|agent| agent.next_timeout(now.monotonic_ms)),
);
timeouts.extend(
self.beacon
.as_ref()
.and_then(|agent| agent.next_timeout(now.monotonic_ms)),
);
}
#[cfg(feature = "portable-autonat")]
if let Some(nat) = self.nat.as_ref()
&& let Some(timeout) = nat.agent.next_timeout(now.monotonic_ms)
{
deadline = PollDeadline::earliest_opt(
deadline,
Some(PollDeadline::from_millis(
now.monotonic_ms.saturating_add(timeout),
)),
);
if !nat.events.is_empty() {
deadline = PollDeadline::earliest_opt(deadline, Some(PollDeadline::IMMEDIATE));
}
}
timeouts.extend(
self.discovery
.as_ref()
.and_then(|agent| agent.next_timeout(now.monotonic_ms)),
);
for timeout in timeouts {
deadline = PollDeadline::earliest_opt(
deadline,
Some(PollDeadline::from_millis(
now.monotonic_ms.saturating_add(timeout),
)),
);
}
deadline
}
pub fn shutdown(mut self, now: Now) -> Result<Vec<SmoltcpEvent>, SmoltcpDriveError> {
if let Some(mdns) = self.mdns.as_mut() {
mdns.shutdown(now.monotonic_ms)?;
}
Ok(self
.endpoint
.shutdown(now)?
.into_iter()
.map(SmoltcpEvent::Endpoint)
.collect())
}
}
#[cfg(feature = "smoltcp")]
#[cfg(feature = "pubsub")]
fn pump_embedded_pubsub<D: smoltcp::phy::Device, E: EntropySource>(
agent: &mut minip2p_pubsub::PubsubAgent,
endpoint: &mut PortableEndpoint<SmoltcpComposedTransport<D, E>, SharedEntropy<E>>,
now: Now,
) {
while let Some(action) = agent.poll_action() {
match action {
minip2p_pubsub::PubsubAction::OpenStream {
token,
peer,
protocol_id,
} => {
let result = endpoint
.open_stream(&peer, &protocol_id, now)
.map_err(|e| e.to_string());
agent.stream_open_result(&peer, token, result, now.monotonic_ms);
}
minip2p_pubsub::PubsubAction::SendStream {
token,
peer,
stream_id,
data,
} => {
let result = endpoint
.send_stream(&peer, stream_id, data, now)
.map_err(|e| e.to_string());
agent.send_result(&peer, stream_id, token, result, now.monotonic_ms);
}
minip2p_pubsub::PubsubAction::CloseStreamWrite { peer, stream_id } => {
let _ = endpoint.close_stream_write(&peer, stream_id, now);
}
minip2p_pubsub::PubsubAction::ResetStream { peer, stream_id } => {
let _ = endpoint.reset_stream(&peer, stream_id, now);
}
}
}
}
#[cfg(all(feature = "smoltcp", feature = "pubsub"))]
fn collect_embedded_pubsub_events(
agent: &mut minip2p_pubsub::PubsubAgent,
pending: &mut VecDeque<PubsubEvent>,
) {
while let Some(event) = agent.poll_event() {
pending.push_back(event);
}
}
#[cfg(feature = "smoltcp")]
impl<D: smoltcp::phy::Device, E: EntropySource> SmoltcpEndpointBuilder<D, E> {
pub fn agent_version(mut self, value: impl Into<String>) -> Self {
self.swarm = self.swarm.agent_version(value);
self
}
pub fn protocol(mut self, protocol_id: impl Into<String>) -> Self {
self.swarm = self.swarm.protocol(protocol_id);
self
}
pub fn tcp_config(mut self, config: TcpConfig) -> Self {
self.tcp_config = config;
self
}
pub fn smoltcp_config(mut self, config: SmoltcpConfig) -> Self {
self.provider_config = config;
self
}
pub fn listen(mut self, address: impl Into<String>) -> Self {
self.listens.push(address.into());
self
}
pub fn mdns(self) -> Self {
self.mdns_config(MdnsConfig::default())
}
pub fn mdns_config(mut self, config: MdnsConfig) -> Self {
self.mdns = Some(EmbeddedMdnsConfig {
mdns: config,
carrier: SmoltcpMdnsConfig::default(),
});
self
}
pub fn mdns_carrier_config(mut self, config: SmoltcpMdnsConfig) -> Self {
self.mdns
.get_or_insert_with(|| EmbeddedMdnsConfig {
mdns: MdnsConfig::default(),
carrier: SmoltcpMdnsConfig::default(),
})
.carrier = config;
self
}
#[cfg(feature = "pubsub")]
pub fn pubsub(mut self) -> Self {
self.pubsub.get_or_insert_with(PubsubConfig::default);
self
}
#[cfg(feature = "pubsub")]
pub fn pubsub_config(mut self, config: impl Into<PubsubConfig>) -> Self {
self.pubsub = Some(config.into());
self
}
#[cfg(feature = "pubsub")]
pub fn discovery(mut self) -> Self {
self.pubsub.get_or_insert_with(PubsubConfig::default);
self.beacon.get_or_insert_with(BeaconConfig::default);
self
}
#[cfg(feature = "pubsub")]
pub fn beacon_config(mut self, config: BeaconConfig) -> Self {
self.pubsub.get_or_insert_with(PubsubConfig::default);
self.beacon = Some(config);
self
}
pub fn discovery_config(mut self, config: PeerDiscoveryConfig) -> Self {
self.discovery = config;
self
}
#[cfg(feature = "portable-relay")]
pub fn relay(mut self, relay: PeerAddr) -> Self {
let mut config = match self.nat_config.take() {
Some(config) => config,
None => minip2p_nat::NatConfig {
reservation_policy: minip2p_nat::ReservationPolicy::Never,
..minip2p_nat::NatConfig::default()
},
};
config.relays.push(relay);
config.force_relay = true;
self.nat_config = Some(config);
self
}
#[cfg(feature = "portable-autonat")]
pub fn autonat(mut self, server: PeerAddr) -> Self {
self.nat_config
.get_or_insert_with(|| minip2p_nat::NatConfig {
reservation_policy: minip2p_nat::ReservationPolicy::Never,
..minip2p_nat::NatConfig::default()
})
.autonat_servers
.push(server);
self
}
#[cfg(feature = "portable-autonat")]
pub fn portable_nat_config(mut self, config: minip2p_nat::NatConfig) -> Self {
self.nat_config = Some(config);
self
}
#[cfg(feature = "portable-relay")]
pub fn relay_config(mut self, config: minip2p_nat::NatConfig) -> Self {
self.nat_config = Some(config);
self
}
pub fn build(mut self) -> Result<SmoltcpEndpoint<D, E>, SmoltcpBuildError> {
#[cfg(feature = "portable-autonat")]
if let Some(config) = &self.nat_config {
if config.relays.is_empty() && config.autonat_servers.is_empty() {
return Err(SmoltcpBuildError::NatConfig {
reason: "at least one relay or AutoNAT server is required",
});
}
if config.relays.is_empty()
&& config.reservation_policy != minip2p_nat::ReservationPolicy::Never
{
return Err(SmoltcpBuildError::NatConfig {
reason: "relay reservation policy requires at least one relay address",
});
}
if !config.relays.is_empty() && !config.force_relay {
return Err(SmoltcpBuildError::NatConfig {
reason: "portable TCP relay requires force_relay = true; DCUtR is QUIC-only",
});
}
#[cfg(not(feature = "portable-relay"))]
if !config.relays.is_empty() {
return Err(SmoltcpBuildError::NatConfig {
reason: "relay addresses require the portable-relay feature",
});
}
if !config.autonat_servers.is_empty() {
self.swarm = self.swarm.protocol(minip2p_nat::AUTONAT_PROTOCOL_ID);
}
}
#[cfg(feature = "pubsub")]
if let Some(config) = &self.pubsub {
for protocol in config.protocol_ids() {
self.swarm = self.swarm.protocol(*protocol);
}
}
let transport = TcpTransport::with_config(
SmoltcpTcpProvider::on_stack(self.stack.clone(), self.provider_config),
self.identity.clone(),
self.entropy.clone(),
self.tcp_config,
);
#[cfg(feature = "portable-relay")]
let transport = minip2p_circuit::CircuitTransport::new(
transport,
self.identity.clone(),
self.entropy.clone(),
);
let mut endpoint = PortableEndpoint {
runtime: self.swarm.build_runtime(transport, self.entropy.clone())?,
};
#[cfg(feature = "portable-autonat")]
if self
.nat_config
.as_ref()
.is_some_and(|config| !config.relays.is_empty())
{
endpoint
.runtime
.add_outbound_protocol(minip2p_nat::HOP_PROTOCOL_ID)?;
endpoint
.runtime
.add_inbound_protocol(minip2p_nat::STOP_PROTOCOL_ID)?;
endpoint
.runtime
.add_advertised_protocol(minip2p_nat::STOP_PROTOCOL_ID)?;
}
for address in self.listens {
let parsed =
address
.parse::<Multiaddr>()
.map_err(|error| SmoltcpBuildError::ListenAddress {
address,
reason: error.to_string(),
})?;
endpoint.listen(&parsed)?;
}
let mdns = if let Some(config) = self.mdns {
let mut seed = [0; 32];
self.entropy.fill_bytes(&mut seed)?;
let io = SmoltcpMdnsIo::on_stack(self.stack, config.carrier)?;
let agent =
minip2p_mdns::MdnsAgent::new(endpoint.peer_id().clone(), config.mdns.clone(), seed)
.map_err(|error| {
SmoltcpBuildError::MdnsConfig(PortableMdnsConfigError::Mdns(error))
})?;
Some(minip2p_mdns::MdnsDriver::new(agent, io, &config.mdns))
} else {
None
};
#[cfg(feature = "pubsub")]
let mut pubsub = self
.pubsub
.map(|config| {
minip2p_pubsub::PubsubAgent::new(
self.identity.clone(),
config,
self.entropy.next_u64()?,
self.entropy.next_u64()?,
)
.map_err(SmoltcpBuildError::Pubsub)
})
.transpose()?;
#[cfg(feature = "pubsub")]
let beacon = self
.beacon
.map(|config| minip2p_discovery::BeaconAgent::new(self.identity.public_key(), config))
.transpose()?;
#[cfg(feature = "pubsub")]
if let (Some(pubsub), Some(beacon)) = (&mut pubsub, &beacon) {
pubsub.subscribe(beacon.topic(), 0)?;
}
#[cfg(feature = "pubsub")]
let discovery_enabled = mdns.is_some() || beacon.is_some();
#[cfg(not(feature = "pubsub"))]
let discovery_enabled = mdns.is_some();
let discovery = if discovery_enabled {
Some(minip2p_discovery::PeerDiscoveryAgent::new(
endpoint.peer_id().clone(),
self.discovery,
)?)
} else {
None
};
#[cfg(feature = "portable-autonat")]
let nat = self.nat_config.map(|config| {
let relay_addrs = config
.relays
.iter()
.map(|relay| (relay.peer_id().clone(), relay.transport().clone()))
.collect();
nat::PortableNatDriver::new(
minip2p_nat::NatAgent::new(endpoint.peer_id().clone(), config),
relay_addrs,
)
});
Ok(SmoltcpEndpoint {
endpoint,
mdns,
#[cfg(feature = "pubsub")]
pubsub,
#[cfg(feature = "pubsub")]
pending_pubsub_events: VecDeque::new(),
#[cfg(feature = "pubsub")]
beacon,
#[cfg(feature = "portable-autonat")]
nat,
discovery,
active_dials: BTreeMap::new(),
})
}
}
#[cfg(feature = "smoltcp")]
#[derive(Debug)]
pub enum SmoltcpBuildError {
Entropy(minip2p_platform::EntropyError),
Swarm(SwarmError),
Driver(DriverError),
ListenAddress { address: String, reason: String },
Mdns(MdnsError),
MdnsConfig(PortableMdnsConfigError),
#[cfg(feature = "pubsub")]
Pubsub(PubsubConfigError),
#[cfg(feature = "pubsub")]
Topic(TopicError),
Discovery(minip2p_discovery::DiscoveryConfigError),
#[cfg(feature = "portable-autonat")]
NatConfig { reason: &'static str },
}
#[cfg(feature = "smoltcp")]
impl From<minip2p_platform::EntropyError> for SmoltcpBuildError {
fn from(error: minip2p_platform::EntropyError) -> Self {
Self::Entropy(error)
}
}
#[cfg(feature = "smoltcp")]
impl From<SwarmError> for SmoltcpBuildError {
fn from(error: SwarmError) -> Self {
Self::Swarm(error)
}
}
#[cfg(feature = "smoltcp")]
impl From<DriverError> for SmoltcpBuildError {
fn from(error: DriverError) -> Self {
Self::Driver(error)
}
}
#[cfg(feature = "smoltcp")]
impl From<MdnsError> for SmoltcpBuildError {
fn from(error: MdnsError) -> Self {
Self::Mdns(error)
}
}
#[cfg(feature = "smoltcp")]
impl From<PortableMdnsConfigError> for SmoltcpBuildError {
fn from(error: PortableMdnsConfigError) -> Self {
Self::MdnsConfig(error)
}
}
#[cfg(feature = "smoltcp")]
#[cfg(feature = "pubsub")]
impl From<TopicError> for SmoltcpBuildError {
fn from(error: TopicError) -> Self {
Self::Topic(error)
}
}
#[cfg(feature = "smoltcp")]
impl From<minip2p_discovery::DiscoveryConfigError> for SmoltcpBuildError {
fn from(error: minip2p_discovery::DiscoveryConfigError) -> Self {
Self::Discovery(error)
}
}
#[cfg(feature = "smoltcp")]
impl core::fmt::Display for SmoltcpBuildError {
fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
match self {
Self::Entropy(error) => error.fmt(formatter),
Self::Swarm(error) => error.fmt(formatter),
Self::Driver(error) => error.fmt(formatter),
Self::ListenAddress { address, reason } => {
write!(
formatter,
"invalid smoltcp listen address {address}: {reason}"
)
}
Self::Mdns(error) => error.fmt(formatter),
Self::MdnsConfig(error) => error.fmt(formatter),
#[cfg(feature = "pubsub")]
Self::Pubsub(error) => error.fmt(formatter),
#[cfg(feature = "pubsub")]
Self::Topic(error) => error.fmt(formatter),
Self::Discovery(error) => error.fmt(formatter),
#[cfg(feature = "portable-autonat")]
Self::NatConfig { reason } => {
write!(formatter, "invalid portable NAT config: {reason}")
}
}
}
}
#[cfg(feature = "smoltcp")]
#[derive(Debug)]
pub enum SmoltcpDriveError {
Endpoint(DriverError),
Mdns(MdnsError),
}
#[cfg(feature = "smoltcp")]
impl From<DriverError> for SmoltcpDriveError {
fn from(error: DriverError) -> Self {
Self::Endpoint(error)
}
}
#[cfg(feature = "smoltcp")]
impl From<MdnsError> for SmoltcpDriveError {
fn from(error: MdnsError) -> Self {
Self::Mdns(error)
}
}
#[cfg(feature = "smoltcp")]
impl core::fmt::Display for SmoltcpDriveError {
fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
match self {
Self::Endpoint(error) => error.fmt(formatter),
Self::Mdns(error) => error.fmt(formatter),
}
}
}
#[cfg(feature = "portable-relay")]
#[derive(Clone, Debug, Eq, PartialEq)]
pub enum SmoltcpRelayError {
NotEnabled,
}
#[cfg(feature = "portable-relay")]
impl core::fmt::Display for SmoltcpRelayError {
fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
match self {
Self::NotEnabled => write!(formatter, "relay is not enabled; call .relay()"),
}
}
}
#[cfg(all(feature = "smoltcp", feature = "pubsub"))]
#[derive(Debug, Eq, PartialEq)]
pub enum SmoltcpPubsubError {
NotEnabled,
DiscoveryTopicReserved,
Topic(TopicError),
Publish(PublishError),
}
#[cfg(all(feature = "smoltcp", feature = "pubsub"))]
impl From<TopicError> for SmoltcpPubsubError {
fn from(error: TopicError) -> Self {
Self::Topic(error)
}
}
#[cfg(all(feature = "smoltcp", feature = "pubsub"))]
impl From<PublishError> for SmoltcpPubsubError {
fn from(error: PublishError) -> Self {
Self::Publish(error)
}
}
#[cfg(all(feature = "smoltcp", feature = "pubsub"))]
impl core::fmt::Display for SmoltcpPubsubError {
fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
match self {
Self::NotEnabled => write!(
formatter,
"pubsub is not enabled; call .pubsub() or .discovery()"
),
Self::DiscoveryTopicReserved => {
write!(formatter, "signed discovery owns this topic subscription")
}
Self::Topic(error) => error.fmt(formatter),
Self::Publish(error) => error.fmt(formatter),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::Endpoint;
use alloc::{rc::Rc, vec, vec::Vec};
use core::cell::Cell;
use minip2p_platform::EntropyError;
use minip2p_test_support::InMemoryTransport;
use minip2p_transport::{ConnectionEndpoint, TransportEvent};
struct NoopTransport;
impl Transport for NoopTransport {
fn dial(&mut self, _: &PeerAddr) -> Result<ConnectionId, TransportError> {
unreachable!()
}
fn listen(&mut self, _: &Multiaddr) -> Result<Multiaddr, TransportError> {
unreachable!()
}
fn open_stream(&mut self, _: ConnectionId) -> Result<StreamId, TransportError> {
unreachable!()
}
fn send_stream(
&mut self,
_: ConnectionId,
_: StreamId,
_: Vec<u8>,
) -> Result<(), TransportError> {
unreachable!()
}
fn close_stream_write(
&mut self,
_: ConnectionId,
_: StreamId,
) -> Result<(), TransportError> {
unreachable!()
}
fn reset_stream(&mut self, _: ConnectionId, _: StreamId) -> Result<(), TransportError> {
unreachable!()
}
fn close(&mut self, _: ConnectionId) -> Result<(), TransportError> {
unreachable!()
}
fn poll(&mut self, _: Now) -> Result<Vec<TransportEvent>, TransportError> {
Ok(Vec::new())
}
}
struct ZeroEntropy;
impl EntropySource for ZeroEntropy {
fn fill_bytes(&mut self, output: &mut [u8]) -> Result<(), EntropyError> {
output.fill(0);
Ok(())
}
}
struct RecordingTransport {
events: Vec<TransportEvent>,
local_address: Multiaddr,
next_stream_id: u64,
closed: Rc<Cell<bool>>,
}
impl Transport for RecordingTransport {
fn dial(&mut self, _: &PeerAddr) -> Result<ConnectionId, TransportError> {
unreachable!()
}
fn listen(&mut self, address: &Multiaddr) -> Result<Multiaddr, TransportError> {
Ok(address.clone())
}
fn open_stream(&mut self, _: ConnectionId) -> Result<StreamId, TransportError> {
self.next_stream_id += 1;
Ok(StreamId::new(self.next_stream_id))
}
fn send_stream(
&mut self,
_: ConnectionId,
_: StreamId,
_: Vec<u8>,
) -> Result<(), TransportError> {
Ok(())
}
fn close_stream_write(
&mut self,
_: ConnectionId,
_: StreamId,
) -> Result<(), TransportError> {
Ok(())
}
fn reset_stream(&mut self, _: ConnectionId, _: StreamId) -> Result<(), TransportError> {
Ok(())
}
fn close(&mut self, id: ConnectionId) -> Result<(), TransportError> {
self.closed.set(true);
self.events.push(TransportEvent::Closed { id });
Ok(())
}
fn poll(&mut self, _: Now) -> Result<Vec<TransportEvent>, TransportError> {
Ok(core::mem::take(&mut self.events))
}
fn local_addresses(&self) -> Vec<Multiaddr> {
vec![self.local_address.clone()]
}
}
#[test]
fn portable_builder_constructs_a_protocol_ready_endpoint() {
let identity = Ed25519Keypair::from_secret_key_bytes([7; 32]);
let mut endpoint = Endpoint::portable(&identity, ZeroEntropy)
.protocol("/example/1.0.0")
.build(NoopTransport)
.expect("portable endpoint configuration is valid");
let remote = Ed25519Keypair::from_secret_key_bytes([8; 32]).peer_id();
assert_eq!(endpoint.peer_id(), &identity.peer_id());
assert!(matches!(
endpoint
.runtime_mut()
.open_stream(&remote, "/example/1.0.0", 0),
Err(DriverError::Swarm(SwarmError::NotConnected { .. }))
));
assert!(matches!(
endpoint
.runtime_mut()
.open_stream(&remote, "/missing/1.0.0", 0),
Err(DriverError::Swarm(SwarmError::ProtocolNotRegistered { .. }))
));
let error = Endpoint::portable(&identity, ZeroEntropy)
.protocol(minip2p_swarm::RESERVED_PROTOCOL_IDS[0])
.build(NoopTransport)
.err()
.expect("reserved protocols must be rejected");
assert!(matches!(error, SwarmError::ReservedProtocol { .. }));
}
#[test]
fn portable_stats_count_identify_ready_peers() {
let a_identity = Ed25519Keypair::from_secret_key_bytes([11; 32]);
let b_identity = Ed25519Keypair::from_secret_key_bytes([12; 32]);
let (a_transport, b_transport) =
InMemoryTransport::pair(a_identity.peer_id(), b_identity.peer_id());
let mut a = Endpoint::portable(&a_identity, ZeroEntropy)
.build(a_transport)
.expect("a endpoint configuration is valid");
let mut b = Endpoint::portable(&b_identity, ZeroEntropy)
.build(b_transport)
.expect("b endpoint configuration is valid");
for now_ms in 0..64 {
a.poll(Now::from_millis(now_ms)).expect("drive a");
b.poll(Now::from_millis(now_ms)).expect("drive b");
if a.is_peer_ready(&b_identity.peer_id()) && b.is_peer_ready(&a_identity.peer_id()) {
break;
}
}
assert_eq!(a.stats().connected_peers, 1);
assert_eq!(a.stats().ready_peers, 1);
assert_eq!(b.stats().connected_peers, 1);
assert_eq!(b.stats().ready_peers, 1);
}
#[test]
fn portable_shutdown_closes_established_peers_and_releases_the_endpoint() {
let identity = Ed25519Keypair::from_secret_key_bytes([9; 32]);
let remote = Ed25519Keypair::from_secret_key_bytes([10; 32]).peer_id();
let address: Multiaddr = "/ip4/127.0.0.1/tcp/4001".parse().expect("address");
let conn_id = ConnectionId::new(1);
let closed = Rc::new(Cell::new(false));
let transport = RecordingTransport {
events: vec![TransportEvent::Connected {
id: conn_id,
endpoint: ConnectionEndpoint::with_peer_id(address.clone(), remote.clone()),
}],
local_address: address.clone(),
next_stream_id: 0,
closed: closed.clone(),
};
let mut endpoint = Endpoint::portable(&identity, ZeroEntropy)
.build(transport)
.expect("portable endpoint configuration is valid");
endpoint
.poll(Now::from_millis(1))
.expect("connection event is accepted");
assert_eq!(
endpoint.stats(),
PortableEndpointStats {
connected_peers: 1,
ready_peers: 0,
listen_addresses: vec![address],
}
);
let events = endpoint
.shutdown(Now::from_millis(2))
.expect("shutdown succeeds");
assert!(
closed.get(),
"the established transport connection was closed"
);
assert!(events.iter().any(|event| matches!(
event,
SwarmEvent::ConnectionClosed { peer_id, conn_id: closed_id, .. }
if peer_id == &remote && closed_id == &conn_id
)));
}
#[cfg(feature = "portable-mdns")]
#[test]
fn failed_autodial_is_correlated_by_connection_without_a_peer_id() {
let peer = Ed25519Keypair::from_secret_key_bytes([42; 32]).peer_id();
let connection = ConnectionId::new(77);
let mut active = BTreeMap::from([(peer.clone(), connection)]);
assert_eq!(
remove_failed_autodial(&mut active, None, Some(connection)),
Some(peer)
);
assert!(active.is_empty(), "a failed dial must become retryable");
}
}