From 35c98149b80ee018e2e3b238b988eed5739af488 Mon Sep 17 00:00:00 2001 From: manasnagaraj Date: Wed, 18 Dec 2024 17:49:15 -0500 Subject: [PATCH 01/10] emit newexternaladdr with discovered --- protocols/mdns/src/behaviour.rs | 16 +++++++++++++++- 1 file changed, 15 insertions(+), 1 deletion(-) diff --git a/protocols/mdns/src/behaviour.rs b/protocols/mdns/src/behaviour.rs index b6dde8f4487..c8da7226caa 100644 --- a/protocols/mdns/src/behaviour.rs +++ b/protocols/mdns/src/behaviour.rs @@ -34,7 +34,7 @@ use std::{ task::{Context, Poll}, time::Instant, }; - +use std::collections::VecDeque; use futures::{channel::mpsc, Stream, StreamExt}; use if_watch::IfEvent; use libp2p_core::{transport::PortUse, Endpoint, Multiaddr}; @@ -188,6 +188,9 @@ where listen_addresses: Arc>, local_peer_id: PeerId, + + /// Pending behaviour events to be emitted. + pending_events: VecDeque, } impl

Behaviour

@@ -208,6 +211,7 @@ where closest_expiration: Default::default(), listen_addresses: Default::default(), local_peer_id, + pending_events: Default::default(), }) } @@ -304,6 +308,11 @@ where &mut self, cx: &mut Context<'_>, ) -> Poll>> { + // Checking for pending events and emit them + if let Some(event) = self.pending_events.pop_front() { + return Poll::Ready(ToSwarm::GenerateEvent(event)); + } + // Poll ifwatch. while let Poll::Ready(Some(event)) = Pin::new(&mut self.if_watch).poll_next(cx) { match event { @@ -360,6 +369,8 @@ where tracing::info!(%peer, address=%addr, "discovered peer on address"); self.discovered_nodes.push((peer, addr.clone(), expiration)); discovered.push((peer, addr)); + + self.pending_events.push_back(Event::NewExternalAddr(addr.clone())); } } @@ -400,6 +411,9 @@ pub enum Event { /// Discovered nodes through mDNS. Discovered(Vec<(PeerId, Multiaddr)>), + /// The multiaddress is reachable externally. + NewExternalAddr(Multiaddr), + /// The given combinations of `PeerId` and `Multiaddr` have expired. /// /// Each discovered record has a time-to-live. When this TTL expires and the address hasn't From 2f19be28138cc4551439fba9f08cfc23f1efcc76 Mon Sep 17 00:00:00 2001 From: manasnagaraj Date: Wed, 18 Dec 2024 19:56:12 -0500 Subject: [PATCH 02/10] some fixes --- protocols/mdns/src/behaviour.rs | 39 +++++++++++++++++---------------- 1 file changed, 20 insertions(+), 19 deletions(-) diff --git a/protocols/mdns/src/behaviour.rs b/protocols/mdns/src/behaviour.rs index c8da7226caa..2eadef9c088 100644 --- a/protocols/mdns/src/behaviour.rs +++ b/protocols/mdns/src/behaviour.rs @@ -22,6 +22,16 @@ mod iface; mod socket; mod timer; +use futures::{channel::mpsc, Stream, StreamExt}; +use if_watch::IfEvent; +use libp2p_core::{transport::PortUse, Endpoint, Multiaddr}; +use libp2p_identity::PeerId; +use libp2p_swarm::{ + behaviour::FromSwarm, dummy, ConnectionDenied, ConnectionId, ListenAddresses, NetworkBehaviour, + THandler, THandlerInEvent, THandlerOutEvent, ToSwarm, +}; +use smallvec::SmallVec; +use std::collections::VecDeque; use std::{ cmp, collections::hash_map::{Entry, HashMap}, @@ -34,17 +44,7 @@ use std::{ task::{Context, Poll}, time::Instant, }; -use std::collections::VecDeque; -use futures::{channel::mpsc, Stream, StreamExt}; -use if_watch::IfEvent; -use libp2p_core::{transport::PortUse, Endpoint, Multiaddr}; -use libp2p_identity::PeerId; -use libp2p_swarm::{ - behaviour::FromSwarm, dummy, ConnectionDenied, ConnectionId, ListenAddresses, NetworkBehaviour, - THandler, THandlerInEvent, THandlerOutEvent, ToSwarm, -}; -use smallvec::SmallVec; - +use std::convert::Infallible; use self::iface::InterfaceState; use crate::{ behaviour::{socket::AsyncSocket, timer::Builder}, @@ -190,7 +190,7 @@ where local_peer_id: PeerId, /// Pending behaviour events to be emitted. - pending_events: VecDeque, + pending_events: VecDeque>, } impl

Behaviour

@@ -310,7 +310,7 @@ where ) -> Poll>> { // Checking for pending events and emit them if let Some(event) = self.pending_events.pop_front() { - return Poll::Ready(ToSwarm::GenerateEvent(event)); + return Poll::Ready(event); } // Poll ifwatch. @@ -368,15 +368,19 @@ where } else { tracing::info!(%peer, address=%addr, "discovered peer on address"); self.discovered_nodes.push((peer, addr.clone(), expiration)); - discovered.push((peer, addr)); + discovered.push((peer, addr.clone())); - self.pending_events.push_back(Event::NewExternalAddr(addr.clone())); + self.pending_events + .push_back(ToSwarm::NewExternalAddrOfPeer { + peer_id: peer, + address: addr, + }); } } if !discovered.is_empty() { let event = Event::Discovered(discovered); - return Poll::Ready(ToSwarm::GenerateEvent(event)); + self.pending_events.push_back(ToSwarm::GenerateEvent(event)); } // Emit expired event. let now = Instant::now(); @@ -411,9 +415,6 @@ pub enum Event { /// Discovered nodes through mDNS. Discovered(Vec<(PeerId, Multiaddr)>), - /// The multiaddress is reachable externally. - NewExternalAddr(Multiaddr), - /// The given combinations of `PeerId` and `Multiaddr` have expired. /// /// Each discovered record has a time-to-live. When this TTL expires and the address hasn't From 50c407db9c87513802e79efbc6af39276d7b5e5b Mon Sep 17 00:00:00 2001 From: manasnagaraj Date: Thu, 19 Dec 2024 12:48:58 -0500 Subject: [PATCH 03/10] requested changes --- protocols/mdns/src/behaviour.rs | 33 +++++++++++++++++++-------------- 1 file changed, 19 insertions(+), 14 deletions(-) diff --git a/protocols/mdns/src/behaviour.rs b/protocols/mdns/src/behaviour.rs index 2eadef9c088..093d30954b8 100644 --- a/protocols/mdns/src/behaviour.rs +++ b/protocols/mdns/src/behaviour.rs @@ -22,19 +22,13 @@ mod iface; mod socket; mod timer; -use futures::{channel::mpsc, Stream, StreamExt}; -use if_watch::IfEvent; -use libp2p_core::{transport::PortUse, Endpoint, Multiaddr}; -use libp2p_identity::PeerId; -use libp2p_swarm::{ - behaviour::FromSwarm, dummy, ConnectionDenied, ConnectionId, ListenAddresses, NetworkBehaviour, - THandler, THandlerInEvent, THandlerOutEvent, ToSwarm, -}; -use smallvec::SmallVec; -use std::collections::VecDeque; use std::{ cmp, - collections::hash_map::{Entry, HashMap}, + collections::{ + hash_map::{Entry, HashMap}, + VecDeque, + }, + convert::Infallible, fmt, future::Future, io, @@ -44,7 +38,17 @@ use std::{ task::{Context, Poll}, time::Instant, }; -use std::convert::Infallible; + +use futures::{channel::mpsc, Stream, StreamExt}; +use if_watch::IfEvent; +use libp2p_core::{transport::PortUse, Endpoint, Multiaddr}; +use libp2p_identity::PeerId; +use libp2p_swarm::{ + behaviour::FromSwarm, dummy, ConnectionDenied, ConnectionId, ListenAddresses, NetworkBehaviour, + THandler, THandlerInEvent, THandlerOutEvent, ToSwarm, +}; +use smallvec::SmallVec; + use self::iface::InterfaceState; use crate::{ behaviour::{socket::AsyncSocket, timer::Builder}, @@ -380,7 +384,8 @@ where if !discovered.is_empty() { let event = Event::Discovered(discovered); - self.pending_events.push_back(ToSwarm::GenerateEvent(event)); + self.pending_events + .push_front(ToSwarm::GenerateEvent(event)); } // Emit expired event. let now = Instant::now(); @@ -397,7 +402,7 @@ where }); if !expired.is_empty() { let event = Event::Expired(expired); - return Poll::Ready(ToSwarm::GenerateEvent(event)); + self.pending_events.push_back(ToSwarm::GenerateEvent(event)); } if let Some(closest_expiration) = closest_expiration { let mut timer = P::Timer::at(closest_expiration); From 1cd6ee58804fc14778f079570cca9668206b2e0c Mon Sep 17 00:00:00 2001 From: manasnagaraj Date: Thu, 19 Dec 2024 15:38:04 -0500 Subject: [PATCH 04/10] added CHANGELOG.md --- protocols/mdns/CHANGELOG.md | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/protocols/mdns/CHANGELOG.md b/protocols/mdns/CHANGELOG.md index 61290703c34..0b626d4fc9b 100644 --- a/protocols/mdns/CHANGELOG.md +++ b/protocols/mdns/CHANGELOG.md @@ -1,3 +1,8 @@ +## 0.46.2 + +- Emit `ToSwarm::NewExternalAddrOfPeer` on discovery + See [PR 5753](https://github.com/libp2p/rust-libp2p/pull/5753) + ## 0.46.1 - Upgrade `hickory-proto`. From 442c6eaf5dd7ba41529326b723db93ab269b7e7c Mon Sep 17 00:00:00 2001 From: manasnagaraj Date: Thu, 19 Dec 2024 15:54:31 -0500 Subject: [PATCH 05/10] version bump --- Cargo.lock | 2 +- Cargo.toml | 2 +- protocols/mdns/Cargo.toml | 2 +- 3 files changed, 3 insertions(+), 3 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 31df58e8ec4..43bcd4c8689 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2985,7 +2985,7 @@ dependencies = [ [[package]] name = "libp2p-mdns" -version = "0.46.1" +version = "0.46.2" dependencies = [ "async-io", "async-std", diff --git a/Cargo.toml b/Cargo.toml index c77768db311..7a116949bcc 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -84,7 +84,7 @@ libp2p-gossipsub = { version = "0.48.0", path = "protocols/gossipsub" } libp2p-identify = { version = "0.46.1", path = "protocols/identify" } libp2p-identity = { version = "0.2.10" } libp2p-kad = { version = "0.47.1", path = "protocols/kad" } -libp2p-mdns = { version = "0.46.1", path = "protocols/mdns" } +libp2p-mdns = { version = "0.46.2", path = "protocols/mdns" } libp2p-memory-connection-limits = { version = "0.3.1", path = "misc/memory-connection-limits" } libp2p-metrics = { version = "0.15.0", path = "misc/metrics" } libp2p-mplex = { version = "0.42.0", path = "muxers/mplex" } diff --git a/protocols/mdns/Cargo.toml b/protocols/mdns/Cargo.toml index 16436848efe..618d41e9b9d 100644 --- a/protocols/mdns/Cargo.toml +++ b/protocols/mdns/Cargo.toml @@ -2,7 +2,7 @@ name = "libp2p-mdns" edition = "2021" rust-version = { workspace = true } -version = "0.46.1" +version = "0.46.2" description = "Implementation of the libp2p mDNS discovery method" authors = ["Parity Technologies "] license = "MIT" From 3a9e1fddfeccea9ead5783bb10bb8e45afabde57 Mon Sep 17 00:00:00 2001 From: manasnagaraj Date: Thu, 19 Dec 2024 19:40:32 -0500 Subject: [PATCH 06/10] connecting the dots --- protocols/mdns/CHANGELOG.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/protocols/mdns/CHANGELOG.md b/protocols/mdns/CHANGELOG.md index 0b626d4fc9b..98dc3d55454 100644 --- a/protocols/mdns/CHANGELOG.md +++ b/protocols/mdns/CHANGELOG.md @@ -1,6 +1,6 @@ ## 0.46.2 -- Emit `ToSwarm::NewExternalAddrOfPeer` on discovery +- Emit `ToSwarm::NewExternalAddrOfPeer` on discovery. See [PR 5753](https://github.com/libp2p/rust-libp2p/pull/5753) ## 0.46.1 From 6cf77425091a3e68b5d26c3a87923eb95848cc61 Mon Sep 17 00:00:00 2001 From: manasnagaraj Date: Sun, 22 Dec 2024 17:29:47 -0500 Subject: [PATCH 07/10] fixing sequence of events --- protocols/mdns/src/behaviour.rs | 175 +++++++++++++++++--------------- 1 file changed, 93 insertions(+), 82 deletions(-) diff --git a/protocols/mdns/src/behaviour.rs b/protocols/mdns/src/behaviour.rs index 093d30954b8..a8cec119938 100644 --- a/protocols/mdns/src/behaviour.rs +++ b/protocols/mdns/src/behaviour.rs @@ -313,104 +313,115 @@ where cx: &mut Context<'_>, ) -> Poll>> { // Checking for pending events and emit them - if let Some(event) = self.pending_events.pop_front() { - return Poll::Ready(event); - } - // Poll ifwatch. - while let Poll::Ready(Some(event)) = Pin::new(&mut self.if_watch).poll_next(cx) { - match event { - Ok(IfEvent::Up(inet)) => { - let addr = inet.addr(); - if addr.is_loopback() { - continue; - } - if addr.is_ipv4() && self.config.enable_ipv6 - || addr.is_ipv6() && !self.config.enable_ipv6 - { - continue; - } - if let Entry::Vacant(e) = self.if_tasks.entry(addr) { - match InterfaceState::::new( - addr, - self.config.clone(), - self.local_peer_id, - self.listen_addresses.clone(), - self.query_response_sender.clone(), - ) { - Ok(iface_state) => { - e.insert(P::spawn(iface_state)); - } - Err(err) => { - tracing::error!("failed to create `InterfaceState`: {}", err) + loop { + if let Some(event) = self.pending_events.pop_front() { + return Poll::Ready(event); + } + + // Poll ifwatch. + while let Poll::Ready(Some(event)) = Pin::new(&mut self.if_watch).poll_next(cx) { + match event { + Ok(IfEvent::Up(inet)) => { + let addr = inet.addr(); + if addr.is_loopback() { + continue; + } + if addr.is_ipv4() && self.config.enable_ipv6 + || addr.is_ipv6() && !self.config.enable_ipv6 + { + continue; + } + if let Entry::Vacant(e) = self.if_tasks.entry(addr) { + match InterfaceState::::new( + addr, + self.config.clone(), + self.local_peer_id, + self.listen_addresses.clone(), + self.query_response_sender.clone(), + ) { + Ok(iface_state) => { + e.insert(P::spawn(iface_state)); + } + Err(err) => { + tracing::error!("failed to create `InterfaceState`: {}", err) + } } } } - } - Ok(IfEvent::Down(inet)) => { - if let Some(handle) = self.if_tasks.remove(&inet.addr()) { - tracing::info!(instance=%inet.addr(), "dropping instance"); + Ok(IfEvent::Down(inet)) => { + if let Some(handle) = self.if_tasks.remove(&inet.addr()) { + tracing::info!(instance=%inet.addr(), "dropping instance"); - handle.abort(); + handle.abort(); + } } + Err(err) => tracing::error!("if watch returned an error: {}", err), } - Err(err) => tracing::error!("if watch returned an error: {}", err), } - } - // Emit discovered event. - let mut discovered = Vec::new(); - - while let Poll::Ready(Some((peer, addr, expiration))) = - self.query_response_receiver.poll_next_unpin(cx) - { - if let Some((_, _, cur_expires)) = self - .discovered_nodes - .iter_mut() - .find(|(p, a, _)| *p == peer && *a == addr) + // Emit discovered event. + let mut discovered = Vec::new(); + + while let Poll::Ready(Some((peer, addr, expiration))) = + self.query_response_receiver.poll_next_unpin(cx) { - *cur_expires = cmp::max(*cur_expires, expiration); - } else { - tracing::info!(%peer, address=%addr, "discovered peer on address"); - self.discovered_nodes.push((peer, addr.clone(), expiration)); - discovered.push((peer, addr.clone())); + if let Some((_, _, cur_expires)) = self + .discovered_nodes + .iter_mut() + .find(|(p, a, _)| *p == peer && *a == addr) + { + *cur_expires = cmp::max(*cur_expires, expiration); + } else { + tracing::info!(%peer, address=%addr, "discovered peer on address"); + self.discovered_nodes.push((peer, addr.clone(), expiration)); + discovered.push((peer, addr.clone())); + + self.pending_events + .push_back(ToSwarm::NewExternalAddrOfPeer { + peer_id: peer, + address: addr, + }); + } + } + if !discovered.is_empty() { + let event = Event::Discovered(discovered); + /// Push to the front of the queue so that the behavior event is reported before + /// the individual discovered addresses. self.pending_events - .push_back(ToSwarm::NewExternalAddrOfPeer { - peer_id: peer, - address: addr, - }); + .push_front(ToSwarm::GenerateEvent(event)); + continue; } - } + // Emit expired event. + let now = Instant::now(); + let mut closest_expiration = None; + let mut expired = Vec::new(); + self.discovered_nodes.retain(|(peer, addr, expiration)| { + if *expiration <= now { + tracing::info!(%peer, address=%addr, "expired peer on address"); + expired.push((*peer, addr.clone())); + return false; + } + closest_expiration = + Some(closest_expiration.unwrap_or(*expiration).min(*expiration)); + true + }); + if !expired.is_empty() { + let event = Event::Expired(expired); + self.pending_events.push_back(ToSwarm::GenerateEvent(event)); + continue; + } + if let Some(closest_expiration) = closest_expiration { + let mut timer = P::Timer::at(closest_expiration); + let _ = Pin::new(&mut timer).poll_next(cx); - if !discovered.is_empty() { - let event = Event::Discovered(discovered); - self.pending_events - .push_front(ToSwarm::GenerateEvent(event)); - } - // Emit expired event. - let now = Instant::now(); - let mut closest_expiration = None; - let mut expired = Vec::new(); - self.discovered_nodes.retain(|(peer, addr, expiration)| { - if *expiration <= now { - tracing::info!(%peer, address=%addr, "expired peer on address"); - expired.push((*peer, addr.clone())); - return false; + self.closest_expiration = Some(timer); } - closest_expiration = Some(closest_expiration.unwrap_or(*expiration).min(*expiration)); - true - }); - if !expired.is_empty() { - let event = Event::Expired(expired); - self.pending_events.push_back(ToSwarm::GenerateEvent(event)); - } - if let Some(closest_expiration) = closest_expiration { - let mut timer = P::Timer::at(closest_expiration); - let _ = Pin::new(&mut timer).poll_next(cx); - self.closest_expiration = Some(timer); + if self.pending_events.is_empty() { + return Poll::Pending; + } } - Poll::Pending } } From 725b3bf987c6f175ab06af9a2ebaab14638b48f8 Mon Sep 17 00:00:00 2001 From: manasnagaraj Date: Sun, 22 Dec 2024 17:35:48 -0500 Subject: [PATCH 08/10] fixing some logic --- protocols/mdns/src/behaviour.rs | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/protocols/mdns/src/behaviour.rs b/protocols/mdns/src/behaviour.rs index a8cec119938..ab15868f05c 100644 --- a/protocols/mdns/src/behaviour.rs +++ b/protocols/mdns/src/behaviour.rs @@ -418,9 +418,7 @@ where self.closest_expiration = Some(timer); } - if self.pending_events.is_empty() { - return Poll::Pending; - } + return Poll::Pending; } } } From 9285a5cc3ebfcce0f2d7bf760ba29cf8dc0712e9 Mon Sep 17 00:00:00 2001 From: Elena Frank Date: Mon, 23 Dec 2024 14:15:14 +0700 Subject: [PATCH 09/10] Fix comment --- protocols/mdns/src/behaviour.rs | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/protocols/mdns/src/behaviour.rs b/protocols/mdns/src/behaviour.rs index ab15868f05c..5016fc4760d 100644 --- a/protocols/mdns/src/behaviour.rs +++ b/protocols/mdns/src/behaviour.rs @@ -312,9 +312,8 @@ where &mut self, cx: &mut Context<'_>, ) -> Poll>> { - // Checking for pending events and emit them - loop { + // Check for pending events and emit them. if let Some(event) = self.pending_events.pop_front() { return Poll::Ready(event); } From 8ac6928ca261cab2085b4f25206938a843693868 Mon Sep 17 00:00:00 2001 From: Elena Frank Date: Mon, 23 Dec 2024 14:24:12 +0700 Subject: [PATCH 10/10] Fix unused doc comment. --- protocols/mdns/src/behaviour.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/protocols/mdns/src/behaviour.rs b/protocols/mdns/src/behaviour.rs index 5016fc4760d..68e28cf3d63 100644 --- a/protocols/mdns/src/behaviour.rs +++ b/protocols/mdns/src/behaviour.rs @@ -385,8 +385,8 @@ where if !discovered.is_empty() { let event = Event::Discovered(discovered); - /// Push to the front of the queue so that the behavior event is reported before - /// the individual discovered addresses. + // Push to the front of the queue so that the behavior event is reported before + // the individual discovered addresses. self.pending_events .push_front(ToSwarm::GenerateEvent(event)); continue;