Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion commons/zenoh-test/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -312,7 +312,7 @@ impl TestSessions {
listener_runtime.start().await.unwrap();

let locators = listener_runtime
.get_locators()
.get_locators_noloopback()
.into_iter()
.map(|l| l.to_endpoint())
.collect();
Expand Down
22 changes: 18 additions & 4 deletions io/zenoh-link-commons/src/listener.rs
Original file line number Diff line number Diff line change
Expand Up @@ -138,7 +138,14 @@ impl ListenersUnicastIP {
.collect()
}

fn get_locators_impl(&self, noloopback: bool) -> Vec<Locator> {
/// Returns the set of listener locators across all listener endpoints.
///
/// When resolving unspecified addresses, loopback addresses are excluded if
/// and only if `exclude_unspecified_lo` is set to `true`.
///
/// Note that `exclude_unspecified_lo` does not impact "explicit" loopback locators
/// (e.g. `tcp/127.0.0.1:9000` or `udp/localhost:0`.)
fn get_locators_impl(&self, exclude_unspecified_lo: bool) -> Vec<Locator> {
let mut locators = vec![];

let guard = zread!(self.listeners);
Expand All @@ -150,8 +157,12 @@ impl ListenersUnicastIP {
// Either ipv4/0.0.0.0 or ipv6/[::]
if kip.is_unspecified() {
let mut addrs = match kip {
IpAddr::V4(_) => zenoh_util::net::get_ipv4_ipaddrs(iface, noloopback),
IpAddr::V6(_) => zenoh_util::net::get_ipv6_ipaddrs(iface, noloopback),
IpAddr::V4(_) => {
zenoh_util::net::get_ipv4_ipaddrs(iface, exclude_unspecified_lo)
}
IpAddr::V6(_) => {
zenoh_util::net::get_ipv6_ipaddrs(iface, exclude_unspecified_lo)
}
};
let iter = addrs.drain(..).map(|x| {
Locator::new(
Expand All @@ -162,18 +173,21 @@ impl ListenersUnicastIP {
.unwrap()
});
locators.extend(iter);
} else if !noloopback || !kip.is_loopback() {
} else {
locators.push(value.endpoint.to_locator());
}
}

locators
}

/// Returns the set of listener locators across all listener endpoints.
pub fn get_locators(&self) -> Vec<Locator> {
self.get_locators_impl(false)
}

/// Returns the set of listener locators across all listener endpoints;
/// excludes loopback locators obtained by resolving unspecified listener endpoints.
pub fn get_locators_noloopback(&self) -> Vec<Locator> {
self.get_locators_impl(true)
}
Expand Down
2 changes: 2 additions & 0 deletions io/zenoh-links/zenoh-link-quic/src/unicast.rs
Original file line number Diff line number Diff line change
Expand Up @@ -340,10 +340,12 @@ impl LinkManagerUnicastTrait for LinkManagerUnicastQuic {
self.listeners.get_endpoints()
}

/// See [`zenoh_link_commons::ListenersUnicastIP::get_locators`].
async fn get_locators(&self) -> Vec<Locator> {
self.listeners.get_locators()
}

/// See [`zenoh_link_commons::ListenersUnicastIP::get_locators_noloopback`].
async fn get_locators_noloopback(&self) -> Vec<Locator> {
self.listeners.get_locators_noloopback()
}
Expand Down
2 changes: 2 additions & 0 deletions io/zenoh-links/zenoh-link-quic_datagram/src/unicast.rs
Original file line number Diff line number Diff line change
Expand Up @@ -324,10 +324,12 @@ impl LinkManagerUnicastTrait for LinkManagerUnicastQuicDatagram {
self.listeners.get_endpoints()
}

/// See [`zenoh_link_commons::ListenersUnicastIP::get_locators`].
async fn get_locators(&self) -> Vec<Locator> {
self.listeners.get_locators()
}

/// See [`zenoh_link_commons::ListenersUnicastIP::get_locators_noloopback`].
async fn get_locators_noloopback(&self) -> Vec<Locator> {
self.listeners.get_locators_noloopback()
}
Expand Down
2 changes: 2 additions & 0 deletions io/zenoh-links/zenoh-link-tcp/src/unicast.rs
Original file line number Diff line number Diff line change
Expand Up @@ -387,10 +387,12 @@ impl LinkManagerUnicastTrait for LinkManagerUnicastTcp {
self.listeners.get_endpoints()
}

/// See [`zenoh_link_commons::ListenersUnicastIP::get_locators`].
async fn get_locators(&self) -> Vec<Locator> {
self.listeners.get_locators()
}

/// See [`zenoh_link_commons::ListenersUnicastIP::get_locators_noloopback`].
async fn get_locators_noloopback(&self) -> Vec<Locator> {
self.listeners.get_locators_noloopback()
}
Expand Down
2 changes: 2 additions & 0 deletions io/zenoh-links/zenoh-link-tls/src/unicast.rs
Original file line number Diff line number Diff line change
Expand Up @@ -482,10 +482,12 @@ impl LinkManagerUnicastTrait for LinkManagerUnicastTls {
self.listeners.get_endpoints()
}

/// See [`zenoh_link_commons::ListenersUnicastIP::get_locators`].
async fn get_locators(&self) -> Vec<Locator> {
self.listeners.get_locators()
}

/// See [`zenoh_link_commons::ListenersUnicastIP::get_locators_noloopback`].
async fn get_locators_noloopback(&self) -> Vec<Locator> {
self.listeners.get_locators_noloopback()
}
Expand Down
2 changes: 2 additions & 0 deletions io/zenoh-links/zenoh-link-udp/src/unicast.rs
Original file line number Diff line number Diff line change
Expand Up @@ -606,10 +606,12 @@ impl LinkManagerUnicastTrait for LinkManagerUnicastUdp {
self.listeners.get_endpoints()
}

/// See [`zenoh_link_commons::ListenersUnicastIP::get_locators`].
async fn get_locators(&self) -> Vec<Locator> {
self.listeners.get_locators()
}

/// See [`zenoh_link_commons::ListenersUnicastIP::get_locators_noloopback`].
async fn get_locators_noloopback(&self) -> Vec<Locator> {
self.listeners.get_locators_noloopback()
}
Expand Down
21 changes: 11 additions & 10 deletions io/zenoh-transport/src/unicast/manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -512,19 +512,19 @@ impl TransportManager {
);
tracing::trace!("{}", e);
let (l, asl) = link.fail();
return Err(InitTransportError::Link((
return Err(InitTransportError::Link(Box::new((
e.into(),
l,
asl,
close::reason::INVALID,
)));
))));
}

// Add the link to the transport
let (start_tx, start_rx, ack, add_link_guard) = transport
.add_link(link, other_initial_sn, other_lease)
.await
.map_err(InitTransportError::Link)?;
.map_err(|e| InitTransportError::Link(Box::new(e)))?;

// complete establish procedure
let c_link = ack.link();
Expand Down Expand Up @@ -605,7 +605,7 @@ impl TransportManager {
Ok(output) => output,
Err(e) => {
let (l, asl) = link.fail();
return Err(InitTransportError::Link((e, l, asl, $reason)));
return Err(InitTransportError::Link(Box::new((e, l, asl, $reason))));
}
}
};
Expand All @@ -616,12 +616,12 @@ impl TransportManager {
let e = zerror!("{} Attempt to establish transport to itself", self.zid());
tracing::warn!("{e}");
let (l, asl) = link.fail();
return Err(InitTransportError::Link((
return Err(InitTransportError::Link(Box::new((
e.into(),
l,
asl,
close::reason::CONNECTION_TO_SELF,
)));
))));
}

// Verify that we haven't reached the transport number limit
Expand All @@ -633,12 +633,12 @@ impl TransportManager {
);
tracing::trace!("{e}");
let (l, asl) = link.fail();
return Err(InitTransportError::Link((
return Err(InitTransportError::Link(Box::new((
e.into(),
l,
asl,
close::reason::INVALID,
)));
))));
}

// Create the transport
Expand Down Expand Up @@ -712,7 +712,7 @@ impl TransportManager {
Ok(val) => val,
Err(e) => {
let _ = t.close(e.3).await;
return Err(InitTransportError::Link(e));
return Err(InitTransportError::Link(Box::new(e)));
}
};

Expand Down Expand Up @@ -825,7 +825,8 @@ impl TransportManager {

match init_result {
Ok(transport) => Ok(TransportUnicast(Arc::downgrade(&transport))),
Err(InitTransportError::Link((e, link, associated_link, reason))) => {
Err(InitTransportError::Link(error)) => {
let (e, link, associated_link, reason) = *error;
let _ = link.close(Some(reason)).await;
if let Some(asl) = associated_link {
let _ = asl.close(Some(reason)).await;
Expand Down
2 changes: 1 addition & 1 deletion io/zenoh-transport/src/unicast/transport_unicast_inner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,7 @@ pub(crate) type LinkError = (
);
pub(crate) type TransportError = (zenoh_result::Error, Arc<dyn TransportUnicastTrait>, u8);
pub(crate) enum InitTransportError {
Link(LinkError),
Link(Box<LinkError>),
Transport(TransportError),
}

Expand Down
2 changes: 1 addition & 1 deletion zenoh/src/api/info.rs
Original file line number Diff line number Diff line change
Expand Up @@ -121,7 +121,7 @@ impl SessionInfo {
/// ```
#[zenoh_macros::unstable]
pub fn locators(&self) -> impl Resolve<Vec<Locator>> + '_ {
ResolveClosure::new(|| self.session.runtime().get_locators())
ResolveClosure::new(|| self.session.runtime().get_locators_noloopback())
}

/// Return information about currently opened transport sessions. Transport session is a connection to another zenoh node.
Expand Down
2 changes: 1 addition & 1 deletion zenoh/src/net/runtime/adminspace.rs
Original file line number Diff line number Diff line change
Expand Up @@ -769,7 +769,7 @@ fn metrics(prefix: &keyexpr, context: &AdminContext, query: Query) {
),
zid = context.runtime.state.zid,
whatami = context.runtime.state.whatami,
version = &*LONG_VERSION,
version = *LONG_VERSION,
);
#[cfg(feature = "stats")]
let mut metrics = String::new();
Expand Down
86 changes: 85 additions & 1 deletion zenoh/tests/scouting.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ use zenoh::config::WhatAmI;
use zenoh::sample::SampleKind;
use zenoh_config::{Config, ModeDependentValue, WhatAmIMatcher};
use zenoh_link::EndPoint;
use zenoh_test::get_free_udp_port;
use zenoh_test::{get_free_tcp_port, get_free_udp_port};

#[tokio::test(flavor = "multi_thread")]
async fn multicast_scouting_works_on_loopback() {
Expand Down Expand Up @@ -99,6 +99,90 @@ async fn multicast_scouting_works_on_loopback() {
responder.close().await.unwrap();
}

#[cfg(all(feature = "unstable", feature = "transport_tcp"))]
#[tokio::test(flavor = "multi_thread")]
async fn gossip_autoconnect_works_on_loopback() {
zenoh::init_log_from_env_or("error");

let port_a = get_free_tcp_port();
let port_b = get_free_tcp_port();
let port_c = get_free_tcp_port();
let endpoint_a = format!("tcp/127.0.0.1:{port_a}");
let endpoint_b = format!("tcp/127.0.0.1:{port_b}");
let endpoint_c = format!("tcp/127.0.0.1:{port_c}");

let peer_config = |listen: &str, connect: Option<&str>| {
let mut config = Config::default();
config.set_mode(Some(WhatAmI::Peer)).unwrap();
config.scouting.multicast.set_enabled(Some(false)).unwrap();
config.scouting.gossip.set_enabled(Some(true)).unwrap();
config
.listen
.endpoints
.set(vec![listen.parse::<EndPoint>().unwrap()])
.unwrap();
if let Some(connect) = connect {
config
.connect
.endpoints
.set(vec![connect.parse::<zenoh_config::EndPoints>().unwrap()])
.unwrap();
}
config
};

let peer_a = zenoh::open(peer_config(&endpoint_a, None)).await.unwrap();
let peer_a_events = peer_a
.info()
.transport_events_listener()
.with(flume::bounded(32))
.await
.unwrap();

let peer_b = zenoh::open(peer_config(&endpoint_b, Some(&endpoint_a)))
.await
.unwrap();
let peer_b_zid = peer_b.zid();

timeout(Duration::from_secs(5), async {
loop {
let event = peer_a_events.recv_async().await.unwrap();
if event.kind() == SampleKind::Put && event.transport().zid() == &peer_b_zid {
break;
}
}
})
.await
.expect("timed out waiting for the initial gossip connection");

// C only connects to A. It must discover B through gossip, including B's
// loopback listener address.
let peer_c = zenoh::open(peer_config(&endpoint_c, Some(&endpoint_a)))
.await
.unwrap();
let peer_c_events = peer_c
.info()
.transport_events_listener()
.history(true)
.await
.unwrap();

timeout(Duration::from_secs(5), async {
loop {
let event = peer_c_events.recv_async().await.unwrap();
if event.kind() == SampleKind::Put && event.transport().zid() == &peer_b_zid {
break;
}
}
})
.await
.expect("timed out waiting for gossip discovery on loopback");

peer_c.close().await.unwrap();
peer_b.close().await.unwrap();
peer_a.close().await.unwrap();
}

#[cfg(all(feature = "unstable", feature = "transport_tcp"))]
#[tokio::test(flavor = "multi_thread")]
async fn multicast_autoconnect_works_on_loopback() {
Expand Down
Loading