From 6000ab09c4ed4cf1a56c9d2de4e6ebf72f1463f0 Mon Sep 17 00:00:00 2001 From: ChenYing Kuo Date: Wed, 17 Dec 2025 16:38:36 +0800 Subject: [PATCH 01/11] Support multiple connections in the client mode. Signed-off-by: ChenYing Kuo --- commons/zenoh-config/src/lib.rs | 12 +- commons/zenoh-config/src/mode_dependent.rs | 40 +++- commons/zenoh-protocol/src/core/endpoint.rs | 116 +++++++++++ zenoh-ext/tests/advanced.rs | 28 +-- zenoh-ext/tests/liveliness.rs | 17 +- zenoh/src/net/runtime/orchestrator.rs | 148 +++++++++----- zenoh/tests/adminspace.rs | 3 +- zenoh/tests/authentication.rs | 18 +- zenoh/tests/liveliness.rs | 213 ++++++++++---------- zenoh/tests/open_time.rs | 4 +- zenoh/tests/unicity.rs | 8 +- zenohd/src/main.rs | 4 +- 12 files changed, 404 insertions(+), 207 deletions(-) diff --git a/commons/zenoh-config/src/lib.rs b/commons/zenoh-config/src/lib.rs index f77317fc19..924fadbe2e 100644 --- a/commons/zenoh-config/src/lib.rs +++ b/commons/zenoh-config/src/lib.rs @@ -51,7 +51,7 @@ use validated_struct::ValidatedMapAssociatedTypes; pub use validated_struct::{GetError, ValidatedMap}; pub use wrappers::ZenohId; pub use zenoh_protocol::core::{ - whatami, EndPoint, Locator, WhatAmI, WhatAmIMatcher, WhatAmIMatcherVisitor, + whatami, EndPoint, EndPoints, Locator, WhatAmI, WhatAmIMatcher, WhatAmIMatcherVisitor, }; use zenoh_protocol::{ core::{ @@ -446,8 +446,12 @@ pub fn peer() -> Config { pub fn client, T: Into>(peers: I) -> Config { let mut config = Config::default(); config.set_mode(Some(WhatAmI::Client)).unwrap(); - config.connect.endpoints = - ModeDependentValue::Unique(peers.into_iter().map(|t| t.into()).collect()); + config.connect.endpoints = ModeDependentValue::Unique( + peers + .into_iter() + .map(|t| EndPoints::Single(t.into())) + .collect(), + ); config } @@ -478,7 +482,7 @@ validated_struct::validator! { /// global timeout for full connect cycle pub timeout_ms: Option>, /// The list of endpoints to connect to - pub endpoints: ModeDependentValue>, + pub endpoints: ModeDependentValue>, /// if connection timeout exceed, exit from application pub exit_on_failure: Option>, pub retry: Option, diff --git a/commons/zenoh-config/src/mode_dependent.rs b/commons/zenoh-config/src/mode_dependent.rs index 56c4c67554..871d4909d9 100644 --- a/commons/zenoh-config/src/mode_dependent.rs +++ b/commons/zenoh-config/src/mode_dependent.rs @@ -18,7 +18,7 @@ use serde::{ de::{self, IntoDeserializer, MapAccess, Visitor}, Deserialize, Serialize, }; -use zenoh_protocol::core::{EndPoint, WhatAmI, WhatAmIMatcher, WhatAmIMatcherVisitor}; +use zenoh_protocol::core::{EndPoint, EndPoints, WhatAmI, WhatAmIMatcher, WhatAmIMatcherVisitor}; use crate::AutoConnectStrategy; @@ -331,6 +331,44 @@ impl<'a> serde::Deserialize<'a> for ModeDependentValue> { } } +impl<'a> serde::Deserialize<'a> for ModeDependentValue> { + fn deserialize(deserializer: D) -> Result + where + D: serde::Deserializer<'a>, + { + struct UniqueOrDependent(PhantomData U>); + + impl<'de> Visitor<'de> for UniqueOrDependent>> { + type Value = ModeDependentValue>; + + fn expecting(&self, formatter: &mut fmt::Formatter) -> fmt::Result { + formatter.write_str("list of endpoints or mode dependent list of endpoints") + } + + fn visit_seq(self, mut seq: A) -> Result + where + A: de::SeqAccess<'de>, + { + let mut v = seq.size_hint().map_or_else(Vec::new, Vec::with_capacity); + + while let Some(s) = seq.next_element()? { + v.push(s); + } + Ok(ModeDependentValue::Unique(v)) + } + + fn visit_map(self, map: M) -> Result + where + M: MapAccess<'de>, + { + ModeValues::deserialize(de::value::MapAccessDeserializer::new(map)) + .map(ModeDependentValue::Dependent) + } + } + deserializer.deserialize_any(UniqueOrDependent(PhantomData)) + } +} + impl ModeDependent for Option> { #[inline] fn router(&self) -> Option<&T> { diff --git a/commons/zenoh-protocol/src/core/endpoint.rs b/commons/zenoh-protocol/src/core/endpoint.rs index aa5ac4bafe..3c5c6acc60 100644 --- a/commons/zenoh-protocol/src/core/endpoint.rs +++ b/commons/zenoh-protocol/src/core/endpoint.rs @@ -696,6 +696,122 @@ impl EndPoint { } } +#[derive(Clone, Debug, serde::Serialize)] +pub enum EndPoints { + Single(EndPoint), + Vec(Vec), +} +impl EndPoints { + pub fn flatten(self) -> Vec { + match self { + EndPoints::Single(ep) => vec![ep], + EndPoints::Vec(eps) => eps, + } + } + + pub fn as_vec(&self) -> Vec { + match self { + EndPoints::Single(ep) => vec![ep.clone()], + EndPoints::Vec(eps) => eps.clone(), + } + } + + // Helper function to determine if all EndPoint in an EndPoints enum have the same proto/address + pub fn all_endpoints_have_same_proto_addr(&self) -> bool { + match &self { + EndPoints::Single(_) => true, + EndPoints::Vec(endpoints_vec) => { + if endpoints_vec.is_empty() { + return true; + } + let first_ep_proto_addr = { + let p = endpoints_vec[0].protocol().as_str(); + let a = endpoints_vec[0].address().as_str(); + format!("{}/{}", p, a) + }; + endpoints_vec.iter().all(|ep| { + let p = ep.protocol().as_str(); + let a = ep.address().as_str(); + format!("{}/{}", p, a) == first_ep_proto_addr + }) + } + } + } +} + +impl From for EndPoints { + fn from(ep: EndPoint) -> EndPoints { + EndPoints::Single(ep) + } +} + +impl<'de> serde::Deserialize<'de> for EndPoints { + fn deserialize(deserializer: D) -> Result + where + D: serde::Deserializer<'de>, + { + struct EndPointsVisitor; + + impl<'de> serde::de::Visitor<'de> for EndPointsVisitor { + type Value = EndPoints; + + fn expecting(&self, formatter: &mut fmt::Formatter) -> fmt::Result { + formatter.write_str("a single endpoint string or a list of endpoint strings") + } + + fn visit_str(self, v: &str) -> Result + where + E: serde::de::Error, + { + EndPoint::from_str(v) + .map(EndPoints::Single) + .map_err(serde::de::Error::custom) + } + + fn visit_seq(self, mut seq: A) -> Result + where + A: serde::de::SeqAccess<'de>, + { + let mut vec = Vec::new(); + while let Some(elem) = seq.next_element::()? { + vec.push(elem); + } + Ok(EndPoints::Vec(vec)) + } + } + + deserializer.deserialize_any(EndPointsVisitor) + } +} + +impl TryFrom for EndPoints { + type Error = ZError; + + fn try_from(s: String) -> Result { + const ERR: &str = + "Endpoints must be of the form or [, , ...]"; + if s.starts_with('[') && s.ends_with(']') { + let eps: ZResult> = s[1..s.len() - 1] + .split(',') + .map(|x| EndPoint::from_str(x.trim())) + .collect(); + eps.map(EndPoints::Vec) + } else { + EndPoint::from_str(s.as_str()) + .map(EndPoints::Single) + .map_err(|e| zerror!("{}: {}", ERR, e).into()) + } + } +} + +impl FromStr for EndPoints { + type Err = ZError; + + fn from_str(s: &str) -> Result { + Self::try_from(s.to_owned()) + } +} + #[test] fn endpoints() { assert!(EndPoint::from_str("/").is_err()); diff --git a/zenoh-ext/tests/advanced.rs b/zenoh-ext/tests/advanced.rs index 6d18f18ae1..3bbcb1260c 100644 --- a/zenoh-ext/tests/advanced.rs +++ b/zenoh-ext/tests/advanced.rs @@ -13,7 +13,7 @@ // use zenoh::{key_expr::OwnedNonWildKeyExpr, sample::SampleKind}; -use zenoh_config::{EndPoint, ModeDependentValue, WhatAmI}; +use zenoh_config::{EndPoint, EndPoints, ModeDependentValue, WhatAmI}; use zenoh_ext::{ AdvancedPublisherBuilderExt, AdvancedSubscriberBuilderExt, CacheConfig, HistoryConfig, MissDetectionConfig, RecoveryConfig, @@ -68,7 +68,7 @@ async fn test_advanced_history_inner( let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![endpoint.parse::().unwrap()]) + .set(vec![endpoint.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); c.namespace = sub_namespace; @@ -182,7 +182,7 @@ async fn test_advanced_retransmission_inner( let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![endpoint.parse::().unwrap()]) + .set(vec![endpoint.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); c.namespace = pub_namespace; @@ -196,7 +196,7 @@ async fn test_advanced_retransmission_inner( let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![endpoint.parse::().unwrap()]) + .set(vec![endpoint.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); c.namespace = sub_namespace; @@ -351,7 +351,7 @@ async fn test_advanced_retransmission_periodic_inner( let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![endpoint.parse::().unwrap()]) + .set(vec![endpoint.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); c.namespace = pub_namespace; @@ -365,7 +365,7 @@ async fn test_advanced_retransmission_periodic_inner( let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![endpoint.parse::().unwrap()]) + .set(vec![endpoint.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); c.namespace = sub_namespace; @@ -513,7 +513,7 @@ async fn test_advanced_sample_miss_inner( let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![endpoint.parse::().unwrap()]) + .set(vec![endpoint.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); c.namespace = pub_namespace; @@ -527,7 +527,7 @@ async fn test_advanced_sample_miss_inner( let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![endpoint.parse::().unwrap()]) + .set(vec![endpoint.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); c.namespace = sub_namespace; @@ -671,7 +671,7 @@ async fn test_advanced_retransmission_sample_miss_inner( let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![endpoint.parse::().unwrap()]) + .set(vec![endpoint.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); c.namespace = pub_namespace; @@ -685,7 +685,7 @@ async fn test_advanced_retransmission_sample_miss_inner( let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![endpoint.parse::().unwrap()]) + .set(vec![endpoint.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); c.namespace = sub_namespace; @@ -826,7 +826,7 @@ async fn test_advanced_late_joiner_inner( let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![endpoint.parse::().unwrap()]) + .set(vec![endpoint.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); c.timestamping @@ -843,7 +843,7 @@ async fn test_advanced_late_joiner_inner( let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![endpoint.parse::().unwrap()]) + .set(vec![endpoint.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); c.namespace = sub_namespace; @@ -986,7 +986,7 @@ async fn test_advanced_retransmission_heartbeat_inner( let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![endpoint.parse::().unwrap()]) + .set(vec![endpoint.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); c.namespace = pub_namespace; @@ -1000,7 +1000,7 @@ async fn test_advanced_retransmission_heartbeat_inner( let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![endpoint.parse::().unwrap()]) + .set(vec![endpoint.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); c.namespace = sub_namespace; diff --git a/zenoh-ext/tests/liveliness.rs b/zenoh-ext/tests/liveliness.rs index cd5a776f55..888e26fd31 100644 --- a/zenoh-ext/tests/liveliness.rs +++ b/zenoh-ext/tests/liveliness.rs @@ -17,6 +17,7 @@ use zenoh::{ sample::SampleKind, Wait, }; +use zenoh_config::EndPoints; #[tokio::test(flavor = "multi_thread", worker_threads = 4)] #[allow(deprecated)] @@ -53,7 +54,7 @@ async fn test_liveliness_querying_subscriber_clique() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![PEER1_ENDPOINT.parse::().unwrap()]) + .set(vec![PEER1_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -132,7 +133,7 @@ async fn test_liveliness_querying_subscriber_brokered() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -145,7 +146,7 @@ async fn test_liveliness_querying_subscriber_brokered() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -158,7 +159,7 @@ async fn test_liveliness_querying_subscriber_brokered() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -239,7 +240,7 @@ async fn test_liveliness_fetching_subscriber_clique() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![PEER1_ENDPOINT.parse::().unwrap()]) + .set(vec![PEER1_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -322,7 +323,7 @@ async fn test_liveliness_fetching_subscriber_brokered() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -335,7 +336,7 @@ async fn test_liveliness_fetching_subscriber_brokered() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -348,7 +349,7 @@ async fn test_liveliness_fetching_subscriber_brokered() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); diff --git a/zenoh/src/net/runtime/orchestrator.rs b/zenoh/src/net/runtime/orchestrator.rs index 3fd98356dd..52877b65c9 100644 --- a/zenoh/src/net/runtime/orchestrator.rs +++ b/zenoh/src/net/runtime/orchestrator.rs @@ -37,7 +37,10 @@ use zenoh_config::{ }; use zenoh_link::{Locator, LocatorInspector}; use zenoh_protocol::{ - core::{whatami::WhatAmIMatcher, EndPoint, Metadata, PriorityRange, WhatAmI, ZenohIdProto}, + core::{ + whatami::WhatAmIMatcher, EndPoint, EndPoints, Metadata, PriorityRange, WhatAmI, + ZenohIdProto, + }, scouting::{HelloProto, Scout, ScoutingBody, ScoutingMessage}, }; use zenoh_result::{bail, zerror, ZResult}; @@ -318,7 +321,7 @@ impl Runtime { Ok(()) } - async fn connect_peers(&self, peers: &[EndPoint], single_link: bool) -> ZResult<()> { + async fn connect_peers(&self, peers: &[EndPoints], single_link: bool) -> ZResult<()> { let timeout = self.get_global_connect_timeout(); if timeout.is_zero() { self.connect_peers_impl(peers, single_link).await @@ -338,7 +341,7 @@ impl Runtime { } } - async fn connect_peers_impl(&self, peers: &[EndPoint], single_link: bool) -> ZResult<()> { + async fn connect_peers_impl(&self, peers: &[EndPoints], single_link: bool) -> ZResult<()> { if single_link { self.connect_peers_single_link(peers).await } else { @@ -346,64 +349,87 @@ impl Runtime { } } - async fn connect_peers_single_link(&self, peers: &[EndPoint]) -> ZResult<()> { - let mut peers_to_retry = Vec::new(); - for peer in peers { - let endpoint = peer.clone(); - let retry_config = self.get_connect_retry_config(&endpoint); - if retry_config.timeout().is_zero() || self.get_global_connect_timeout().is_zero() { - tracing::debug!( - "Try to connect: {:?}: global timeout: {:?}, retry: {:?}", - endpoint, - self.get_global_connect_timeout(), - retry_config - ); - // try to connect and exit immediately without retry - if self.peer_connector(endpoint).await.is_ok() { - return Ok(()); + async fn connect_peers_single_link(&self, peers: &[EndPoints]) -> ZResult<()> { + let mut success_flag = false; + for peer_group in peers { + // Check the peer_group has the same proto/host:port. If not, just ignore the group + if !peer_group.all_endpoints_have_same_proto_addr() { + continue; + } + + // Try to connect to each peer in the group + let mut peers_to_retry = Vec::new(); + for peer in peer_group.as_vec() { + let endpoint = peer.clone(); + let retry_config = self.get_connect_retry_config(&endpoint); + if retry_config.timeout().is_zero() || self.get_global_connect_timeout().is_zero() { + tracing::debug!( + "Try to connect: {:?}: global timeout: {:?}, retry: {:?}", + endpoint, + self.get_global_connect_timeout(), + retry_config + ); + // try to connect directly when there is no timeout configuration + if self.peer_connector(endpoint).await.is_ok() { + success_flag = true; + } + } else { + peers_to_retry.push(endpoint); } - } else { - peers_to_retry.push(endpoint); } - } - // sequentially try to connect to one of the remaining peers - // respecting connection retry delays - match self.peers_connector_retry(peers_to_retry, true).await { - Ok(_) => Ok(()), - Err(_) => { - let e = zerror!("Unable to connect to any of {:?}! ", peers); - tracing::warn!("{}", &e); - Err(e.into()) + // sequentially try to connect to one of the remaining peers + // respecting connection retry delays + if self + .peers_connector_retry(peers_to_retry, false) + .await + .is_ok() + { + success_flag = true; + } + // Any endpoint in the group is available, it's marked as success and break + if success_flag { + break; } } + + // Return error if none of them succeeded + if success_flag { + Ok(()) + } else { + let e = zerror!("Unable to connect to any of {:?}! ", peers); + tracing::warn!("{}", &e); + Err(e.into()) + } } - async fn connect_peers_multiply_links(&self, peers: &[EndPoint]) -> ZResult<()> { - for peer in peers { - let endpoint = peer.clone(); - let retry_config = self.get_connect_retry_config(&endpoint); - tracing::debug!( - "Try to connect: {:?}: global timeout: {:?}, retry: {:?}", - endpoint, - self.get_global_connect_timeout(), - retry_config - ); - if retry_config.timeout().is_zero() || self.get_global_connect_timeout().is_zero() { - // try to connect and exit immediately without retry - if let Err(e) = self.peer_connector(endpoint).await { - if retry_config.exit_on_failure { + async fn connect_peers_multiply_links(&self, peers: &[EndPoints]) -> ZResult<()> { + for peer_group in peers { + for peer in peer_group.as_vec() { + let endpoint = peer.clone(); + let retry_config = self.get_connect_retry_config(&endpoint); + tracing::debug!( + "Try to connect: {:?}: global timeout: {:?}, retry: {:?}", + endpoint, + self.get_global_connect_timeout(), + retry_config + ); + if retry_config.timeout().is_zero() || self.get_global_connect_timeout().is_zero() { + // try to connect and exit immediately without retry + if let Err(e) = self.peer_connector(endpoint).await { + if retry_config.exit_on_failure { + return Err(e); + } + } + } else if retry_config.exit_on_failure { + // try to connect with retry waiting + let _ = self.peer_connector_retry(endpoint).await; + } else { + // try to connect in background + if let Err(e) = self.spawn_peer_connector(endpoint.clone()).await { + tracing::warn!("Error connecting to {}: {}", endpoint, e); return Err(e); } } - } else if retry_config.exit_on_failure { - // try to connect with retry waiting - let _ = self.peer_connector_retry(endpoint).await; - } else { - // try to connect in background - if let Err(e) = self.spawn_peer_connector(endpoint.clone()).await { - tracing::warn!("Error connecting to {}: {}", endpoint, e); - return Err(e); - } } } Ok(()) @@ -946,8 +972,10 @@ impl Runtime { .connect() .endpoints() .get(self.whatami()) - .iter() - .flat_map(|e| e.iter().map(EndPoint::to_locator)) + .unwrap_or(&vec![]) + .into_iter() + .flat_map(|e| e.as_vec()) + .map(|e| e.to_locator()) .collect::>(); let locators = scouted_locators @@ -1210,7 +1238,7 @@ impl Runtime { if zread!(session.endpoints).is_empty() { return; } - let mut peers = session + let endpoints = session .runtime .state .config @@ -1221,6 +1249,10 @@ impl Runtime { .get(session.runtime.state.whatami) .unwrap_or(&vec![]) .clone(); + let mut peers = vec![]; + for peer in endpoints { + peers.extend(peer.flatten()); + } if session.runtime.whatami() != WhatAmI::Client { let endpoints = std::mem::take(zwrite!(session.endpoints).deref_mut()); @@ -1246,7 +1278,7 @@ impl Runtime { if session.runtime.is_closed() { return; } - let peers = session + let endpoints = session .runtime .state .config @@ -1257,6 +1289,10 @@ impl Runtime { .get(session.runtime.state.whatami) .unwrap_or(&vec![]) .clone(); + let mut peers = vec![]; + for peer in endpoints { + peers.extend(peer.flatten()); + } if peers.contains(&endpoint) && zwrite!(session.endpoints).remove(&endpoint) { let runtime = session.runtime.clone(); diff --git a/zenoh/tests/adminspace.rs b/zenoh/tests/adminspace.rs index 51ea163aba..d21b1bd97c 100644 --- a/zenoh/tests/adminspace.rs +++ b/zenoh/tests/adminspace.rs @@ -18,6 +18,7 @@ use std::time::Duration; use zenoh_config::WhatAmI; use zenoh_core::ztimeout; use zenoh_link::EndPoint; +use zenoh_protocol::core::EndPoints; #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn test_adminspace_wonly() { @@ -86,7 +87,7 @@ async fn test_adminspace_read() { c.listen.endpoints.set(vec![]).unwrap(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); let s = ztimeout!(zenoh::open(c)).unwrap(); s diff --git a/zenoh/tests/authentication.rs b/zenoh/tests/authentication.rs index 227b87b686..7d88737ebb 100644 --- a/zenoh/tests/authentication.rs +++ b/zenoh/tests/authentication.rs @@ -29,7 +29,7 @@ mod test { config::{WhatAmI, ZenohId}, Config, Session, }; - use zenoh_config::{EndPoint, ModeDependentValue}; + use zenoh_config::{EndPoints, ModeDependentValue}; use zenoh_core::{zlock, ztimeout}; const TIMEOUT: Duration = Duration::from_secs(60); @@ -491,7 +491,7 @@ client2name:client2passwd"; .set_endpoints(ModeDependentValue::Unique(vec![format!( "tls/127.0.0.1:{port}" ) - .parse::() + .parse::() .unwrap()])) .unwrap(); config @@ -544,7 +544,7 @@ client2name:client2passwd"; .set_endpoints(ModeDependentValue::Unique(vec![format!( "tls/127.0.0.1:{port}" ) - .parse::() + .parse::() .unwrap()])) .unwrap(); config @@ -602,7 +602,7 @@ client2name:client2passwd"; .set_endpoints(ModeDependentValue::Unique(vec![format!( "quic/127.0.0.1:{port}" ) - .parse::() + .parse::() .unwrap()])) .unwrap(); config @@ -647,7 +647,7 @@ client2name:client2passwd"; .set_endpoints(ModeDependentValue::Unique(vec![format!( "quic/127.0.0.1:{port}" ) - .parse::() + .parse::() .unwrap()])) .unwrap(); config @@ -697,7 +697,7 @@ client2name:client2passwd"; .set_endpoints(ModeDependentValue::Unique(vec![format!( "tcp/127.0.0.1:{port}" ) - .parse::() + .parse::() .unwrap()])) .unwrap(); config @@ -721,7 +721,7 @@ client2name:client2passwd"; .set_endpoints(ModeDependentValue::Unique(vec![format!( "tcp/127.0.0.1:{port}" ) - .parse::() + .parse::() .unwrap()])) .unwrap(); config @@ -751,7 +751,7 @@ client2name:client2passwd"; .set_endpoints(ModeDependentValue::Unique(vec![format!( "quic/127.0.0.1:{port}" ) - .parse::() + .parse::() .unwrap()])) .unwrap(); config @@ -803,7 +803,7 @@ client2name:client2passwd"; .set_endpoints(ModeDependentValue::Unique(vec![format!( "quic/127.0.0.1:{port}" ) - .parse::() + .parse::() .unwrap()])) .unwrap(); config diff --git a/zenoh/tests/liveliness.rs b/zenoh/tests/liveliness.rs index 971c948a51..444f0a0cea 100644 --- a/zenoh/tests/liveliness.rs +++ b/zenoh/tests/liveliness.rs @@ -13,6 +13,7 @@ // #![cfg(feature = "internal_config")] +use zenoh_config::EndPoints; use zenoh_core::ztimeout; #[tokio::test(flavor = "multi_thread", worker_threads = 4)] @@ -45,7 +46,7 @@ async fn test_liveliness_subscriber_clique() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![PEER1_ENDPOINT.parse::().unwrap()]) + .set(vec![PEER1_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -108,7 +109,7 @@ async fn test_liveliness_query_clique() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![PEER1_ENDPOINT.parse::().unwrap()]) + .set(vec![PEER1_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -165,7 +166,7 @@ async fn test_liveliness_subscriber_brokered() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -178,7 +179,7 @@ async fn test_liveliness_subscriber_brokered() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -242,7 +243,7 @@ async fn test_liveliness_query_brokered() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -255,7 +256,7 @@ async fn test_liveliness_query_brokered() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -387,7 +388,7 @@ async fn test_liveliness_after_close() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![PEER1_ENDPOINT.parse::().unwrap()]) + .set(vec![PEER1_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -447,7 +448,7 @@ async fn test_liveliness_subscriber_double_client_before() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -463,7 +464,7 @@ async fn test_liveliness_subscriber_double_client_before() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -541,7 +542,7 @@ async fn test_liveliness_subscriber_double_client_middle() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -560,7 +561,7 @@ async fn test_liveliness_subscriber_double_client_middle() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -638,7 +639,7 @@ async fn test_liveliness_subscriber_double_client_after() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -662,7 +663,7 @@ async fn test_liveliness_subscriber_double_client_after() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -737,7 +738,7 @@ async fn test_liveliness_subscriber_double_client_history_before() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -753,7 +754,7 @@ async fn test_liveliness_subscriber_double_client_history_before() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -839,7 +840,7 @@ async fn test_liveliness_subscriber_double_client_history_middle() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -859,7 +860,7 @@ async fn test_liveliness_subscriber_double_client_history_middle() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -941,7 +942,7 @@ async fn test_liveliness_subscriber_double_client_history_after() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -967,7 +968,7 @@ async fn test_liveliness_subscriber_double_client_history_after() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -1045,7 +1046,7 @@ async fn test_liveliness_subscriber_double_peer_before() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -1061,7 +1062,7 @@ async fn test_liveliness_subscriber_double_peer_before() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -1133,7 +1134,7 @@ async fn test_liveliness_subscriber_double_peer_middle() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -1149,7 +1150,7 @@ async fn test_liveliness_subscriber_double_peer_middle() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -1224,7 +1225,7 @@ async fn test_liveliness_subscriber_double_peer_after() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -1242,7 +1243,7 @@ async fn test_liveliness_subscriber_double_peer_after() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -1317,7 +1318,7 @@ async fn test_liveliness_subscriber_double_peer_history_before() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -1333,7 +1334,7 @@ async fn test_liveliness_subscriber_double_peer_history_before() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -1419,7 +1420,7 @@ async fn test_liveliness_subscriber_double_peer_history_middle() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -1439,7 +1440,7 @@ async fn test_liveliness_subscriber_double_peer_history_middle() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -1521,7 +1522,7 @@ async fn test_liveliness_subscriber_double_peer_history_after() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -1547,7 +1548,7 @@ async fn test_liveliness_subscriber_double_peer_history_after() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -1626,7 +1627,7 @@ async fn test_liveliness_subscriber_double_router_before() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -1646,7 +1647,7 @@ async fn test_liveliness_subscriber_double_router_before() { .unwrap(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Router)); @@ -1729,7 +1730,7 @@ async fn test_liveliness_subscriber_double_router_middle() { .unwrap(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Router)); @@ -1748,7 +1749,7 @@ async fn test_liveliness_subscriber_double_router_middle() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -1831,7 +1832,7 @@ async fn test_liveliness_subscriber_double_router_after() { .unwrap(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Router)); @@ -1855,7 +1856,7 @@ async fn test_liveliness_subscriber_double_router_after() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -1931,7 +1932,7 @@ async fn test_liveliness_subscriber_double_router_history_before() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -1951,7 +1952,7 @@ async fn test_liveliness_subscriber_double_router_history_before() { .unwrap(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Router)); @@ -2042,7 +2043,7 @@ async fn test_liveliness_subscriber_double_router_history_middle() { .unwrap(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Router)); @@ -2062,7 +2063,7 @@ async fn test_liveliness_subscriber_double_router_history_middle() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -2149,7 +2150,7 @@ async fn test_liveliness_subscriber_double_router_history_after() { .unwrap(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Router)); @@ -2175,7 +2176,7 @@ async fn test_liveliness_subscriber_double_router_history_after() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -2258,7 +2259,7 @@ async fn test_liveliness_subscriber_double_clientviapeer_before() { .unwrap(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -2271,7 +2272,7 @@ async fn test_liveliness_subscriber_double_clientviapeer_before() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -2287,7 +2288,7 @@ async fn test_liveliness_subscriber_double_clientviapeer_before() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![PEER_DUMMY_ENDPOINT.parse::().unwrap()]) + .set(vec![PEER_DUMMY_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -2371,7 +2372,7 @@ async fn test_liveliness_subscriber_double_clientviapeer_middle() { .unwrap(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -2384,7 +2385,7 @@ async fn test_liveliness_subscriber_double_clientviapeer_middle() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![PEER_DUMMY_ENDPOINT.parse::().unwrap()]) + .set(vec![PEER_DUMMY_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -2403,7 +2404,7 @@ async fn test_liveliness_subscriber_double_clientviapeer_middle() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -2487,7 +2488,7 @@ async fn test_liveliness_subscriber_double_clientviapeer_after() { .unwrap(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -2500,7 +2501,7 @@ async fn test_liveliness_subscriber_double_clientviapeer_after() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![PEER_DUMMY_ENDPOINT.parse::().unwrap()]) + .set(vec![PEER_DUMMY_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -2524,7 +2525,7 @@ async fn test_liveliness_subscriber_double_clientviapeer_after() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -2606,7 +2607,7 @@ async fn test_liveliness_subscriber_double_clientviapeer_history_before() { .unwrap(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -2619,7 +2620,7 @@ async fn test_liveliness_subscriber_double_clientviapeer_history_before() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -2635,7 +2636,7 @@ async fn test_liveliness_subscriber_double_clientviapeer_history_before() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![PEER_DUMMY_ENDPOINT.parse::().unwrap()]) + .set(vec![PEER_DUMMY_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -2728,7 +2729,7 @@ async fn test_liveliness_subscriber_double_clientviapeer_history_middle() { .unwrap(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -2741,7 +2742,7 @@ async fn test_liveliness_subscriber_double_clientviapeer_history_middle() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![PEER_DUMMY_ENDPOINT.parse::().unwrap()]) + .set(vec![PEER_DUMMY_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -2761,7 +2762,7 @@ async fn test_liveliness_subscriber_double_clientviapeer_history_middle() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -2850,7 +2851,7 @@ async fn test_liveliness_subscriber_double_clientviapeer_history_after() { .unwrap(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -2863,7 +2864,7 @@ async fn test_liveliness_subscriber_double_clientviapeer_history_after() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![PEER_DUMMY_ENDPOINT.parse::().unwrap()]) + .set(vec![PEER_DUMMY_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -2889,7 +2890,7 @@ async fn test_liveliness_subscriber_double_clientviapeer_history_after() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -2968,7 +2969,7 @@ async fn test_liveliness_subget_client_before() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -2984,7 +2985,7 @@ async fn test_liveliness_subget_client_before() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -3060,7 +3061,7 @@ async fn test_liveliness_subget_client_middle() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -3081,7 +3082,7 @@ async fn test_liveliness_subget_client_middle() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -3157,7 +3158,7 @@ async fn test_liveliness_subget_client_history_before() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -3173,7 +3174,7 @@ async fn test_liveliness_subget_client_history_before() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -3253,7 +3254,7 @@ async fn test_liveliness_subget_client_history_middle() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -3275,7 +3276,7 @@ async fn test_liveliness_subget_client_history_middle() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -3354,7 +3355,7 @@ async fn test_liveliness_subget_peer_before() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -3370,7 +3371,7 @@ async fn test_liveliness_subget_peer_before() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -3446,7 +3447,7 @@ async fn test_liveliness_subget_peer_middle() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -3467,7 +3468,7 @@ async fn test_liveliness_subget_peer_middle() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -3543,7 +3544,7 @@ async fn test_liveliness_subget_peer_history_before() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -3559,7 +3560,7 @@ async fn test_liveliness_subget_peer_history_before() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -3639,7 +3640,7 @@ async fn test_liveliness_subget_peer_history_middle() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -3661,7 +3662,7 @@ async fn test_liveliness_subget_peer_history_middle() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -3741,7 +3742,7 @@ async fn test_liveliness_subget_router_before() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -3761,7 +3762,7 @@ async fn test_liveliness_subget_router_before() { .unwrap(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Router)); @@ -3842,7 +3843,7 @@ async fn test_liveliness_subget_router_middle() { .unwrap(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Router)); @@ -3863,7 +3864,7 @@ async fn test_liveliness_subget_router_middle() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -3940,7 +3941,7 @@ async fn test_liveliness_subget_router_history_before() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -3960,7 +3961,7 @@ async fn test_liveliness_subget_router_history_before() { .unwrap(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Router)); @@ -4045,7 +4046,7 @@ async fn test_liveliness_subget_router_history_middle() { .unwrap(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Router)); @@ -4067,7 +4068,7 @@ async fn test_liveliness_subget_router_history_middle() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -4148,7 +4149,7 @@ async fn test_liveliness_regression_1() { .unwrap(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -4165,8 +4166,8 @@ async fn test_liveliness_regression_1() { c.connect .endpoints .set(vec![ - ROUTER_ENDPOINT.parse::().unwrap(), - PEER_TOK_ENDPOINT.parse::().unwrap(), + ROUTER_ENDPOINT.parse::().unwrap(), + PEER_TOK_ENDPOINT.parse::().unwrap(), ]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); @@ -4234,7 +4235,7 @@ async fn test_liveliness_regression_2() { .unwrap(); c.connect .endpoints - .set(vec![PEER_TOK1_ENDPOINT.parse::().unwrap()]) + .set(vec![PEER_TOK1_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -4253,8 +4254,8 @@ async fn test_liveliness_regression_2() { c.connect .endpoints .set(vec![ - PEER_TOK1_ENDPOINT.parse::().unwrap(), - PEER_SUB_ENDPOINT.parse::().unwrap(), + PEER_TOK1_ENDPOINT.parse::().unwrap(), + PEER_SUB_ENDPOINT.parse::().unwrap(), ]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); @@ -4327,7 +4328,7 @@ async fn test_liveliness_regression_2_history() { .unwrap(); c.connect .endpoints - .set(vec![PEER_TOK1_ENDPOINT.parse::().unwrap()]) + .set(vec![PEER_TOK1_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -4353,8 +4354,8 @@ async fn test_liveliness_regression_2_history() { c.connect .endpoints .set(vec![ - PEER_TOK1_ENDPOINT.parse::().unwrap(), - PEER_SUB_ENDPOINT.parse::().unwrap(), + PEER_TOK1_ENDPOINT.parse::().unwrap(), + PEER_SUB_ENDPOINT.parse::().unwrap(), ]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); @@ -4424,7 +4425,7 @@ async fn test_liveliness_regression_3() { .unwrap(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -4440,7 +4441,7 @@ async fn test_liveliness_regression_3() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -4457,8 +4458,8 @@ async fn test_liveliness_regression_3() { c.connect .endpoints .set(vec![ - ROUTER_ENDPOINT.parse::().unwrap(), - PEER_TOK_ENDPOINT.parse::().unwrap(), + ROUTER_ENDPOINT.parse::().unwrap(), + PEER_TOK_ENDPOINT.parse::().unwrap(), ]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); @@ -4540,7 +4541,7 @@ async fn test_liveliness_issue_1470() { .unwrap(); c.connect .endpoints - .set(vec![ROUTER0_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER0_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Router)); @@ -4563,8 +4564,8 @@ async fn test_liveliness_issue_1470() { c.connect .endpoints .set(vec![ - ROUTER0_ENDPOINT.parse::().unwrap(), - ROUTER1_ENDPOINT.parse::().unwrap(), + ROUTER0_ENDPOINT.parse::().unwrap(), + ROUTER1_ENDPOINT.parse::().unwrap(), ]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); @@ -4580,7 +4581,7 @@ async fn test_liveliness_issue_1470() { c.set_id(Some(ZenohId::from_str("c0").unwrap())).unwrap(); c.connect .endpoints - .set(vec![PEER_ENDPOINT.parse::().unwrap()]) + .set(vec![PEER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -4626,7 +4627,7 @@ async fn test_liveliness_issue_1470() { c.set_id(Some(ZenohId::from_str("c1").unwrap())).unwrap(); c.connect .endpoints - .set(vec![PEER_ENDPOINT.parse::().unwrap()]) + .set(vec![PEER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); @@ -4705,7 +4706,7 @@ async fn test_liveliness_double_undeclare_clique() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![PEER1_ENDPOINT.parse::().unwrap()]) + .set(vec![PEER1_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -4780,7 +4781,7 @@ async fn test_liveliness_sub_history_conflict() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Peer)); @@ -4796,7 +4797,7 @@ async fn test_liveliness_sub_history_conflict() { let mut c = zenoh::Config::default(); c.connect .endpoints - .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) + .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); let _ = c.set_mode(Some(WhatAmI::Client)); diff --git a/zenoh/tests/open_time.rs b/zenoh/tests/open_time.rs index 15e96dbe3a..9a35fbf9e7 100644 --- a/zenoh/tests/open_time.rs +++ b/zenoh/tests/open_time.rs @@ -19,7 +19,7 @@ use std::{ time::{Duration, Instant}, }; -use zenoh_config::Config; +use zenoh_config::{Config, EndPoints}; use zenoh_link::EndPoint; use zenoh_protocol::core::WhatAmI; @@ -74,7 +74,7 @@ async fn time_open( app_config .connect .endpoints - .set(vec![connect_endpoint.clone()]) + .set(vec![connect_endpoint.clone().into()]) .unwrap(); app_config .transport diff --git a/zenoh/tests/unicity.rs b/zenoh/tests/unicity.rs index eefead014e..76c27bbdf6 100644 --- a/zenoh/tests/unicity.rs +++ b/zenoh/tests/unicity.rs @@ -24,7 +24,7 @@ use std::{ use tokio::runtime::Handle; use zenoh::{config::WhatAmI, key_expr::KeyExpr, qos::CongestionControl, Session}; -use zenoh_config::{EndPoint, ModeDependentValue}; +use zenoh_config::{EndPoints, ModeDependentValue}; use zenoh_core::ztimeout; const TIMEOUT: Duration = Duration::from_secs(60); @@ -101,7 +101,7 @@ async fn open_client_sessions() -> (Session, Session, Session) { config .connect .set_endpoints(ModeDependentValue::Unique(vec!["tcp/127.0.0.1:30447" - .parse::() + .parse::() .unwrap()])) .unwrap(); println!("[ ][01a] Opening s01 session"); @@ -112,7 +112,7 @@ async fn open_client_sessions() -> (Session, Session, Session) { config .connect .set_endpoints(ModeDependentValue::Unique(vec!["tcp/127.0.0.1:30447" - .parse::() + .parse::() .unwrap()])) .unwrap(); println!("[ ][02a] Opening s02 session"); @@ -123,7 +123,7 @@ async fn open_client_sessions() -> (Session, Session, Session) { config .connect .set_endpoints(ModeDependentValue::Unique(vec!["tcp/127.0.0.1:30447" - .parse::() + .parse::() .unwrap()])) .unwrap(); println!("[ ][03a] Opening s03 session"); diff --git a/zenohd/src/main.rs b/zenohd/src/main.rs index 17e52ec665..4d43a11325 100644 --- a/zenohd/src/main.rs +++ b/zenohd/src/main.rs @@ -15,7 +15,7 @@ use clap::Parser; use git_version::git_version; use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt, EnvFilter}; use zenoh::{config::WhatAmI, Config, Result, Wait}; -use zenoh_config::{EndPoint, ModeDependentValue, PermissionsConf}; +use zenoh_config::{EndPoint, EndPoints, ModeDependentValue, PermissionsConf}; use zenoh_util::LibSearchDirs; const GIT_VERSION: &str = git_version!(prefix = "v", cargo_prefix = "v"); @@ -167,7 +167,7 @@ fn config_from_args(args: &Args) -> Config { .set( args.connect .iter() - .map(|v| match v.parse::() { + .map(|v| match v.parse::() { Ok(v) => v, Err(e) => { panic!("Couldn't parse option --peer={v} into Locator: {e}"); From d76465801df9b636f98978efb8c0e67d8a9451fb Mon Sep 17 00:00:00 2001 From: ChenYing Kuo Date: Wed, 17 Dec 2025 17:10:24 +0800 Subject: [PATCH 02/11] Fix clippy issue. Signed-off-by: ChenYing Kuo --- zenoh/src/net/runtime/orchestrator.rs | 2 +- zenoh/tests/adminspace.rs | 18 ++++++------------ 2 files changed, 7 insertions(+), 13 deletions(-) diff --git a/zenoh/src/net/runtime/orchestrator.rs b/zenoh/src/net/runtime/orchestrator.rs index 52877b65c9..ff82c444db 100644 --- a/zenoh/src/net/runtime/orchestrator.rs +++ b/zenoh/src/net/runtime/orchestrator.rs @@ -973,7 +973,7 @@ impl Runtime { .endpoints() .get(self.whatami()) .unwrap_or(&vec![]) - .into_iter() + .iter() .flat_map(|e| e.as_vec()) .map(|e| e.to_locator()) .collect::>(); diff --git a/zenoh/tests/adminspace.rs b/zenoh/tests/adminspace.rs index d21b1bd97c..dfe7eceb50 100644 --- a/zenoh/tests/adminspace.rs +++ b/zenoh/tests/adminspace.rs @@ -38,8 +38,7 @@ async fn test_adminspace_wonly() { .peer .set_mode(Some("linkstate".to_string())) .unwrap(); - let s = ztimeout!(zenoh::open(c)).unwrap(); - s + ztimeout!(zenoh::open(c)).unwrap() }; let zid = router.zid(); let root = router @@ -77,8 +76,7 @@ async fn test_adminspace_read() { .peer .set_mode(Some("linkstate".to_string())) .unwrap(); - let s = ztimeout!(zenoh::open(c)).unwrap(); - s + ztimeout!(zenoh::open(c)).unwrap() }; let zid = router.zid(); let router2 = { @@ -89,8 +87,7 @@ async fn test_adminspace_read() { .endpoints .set(vec![ROUTER_ENDPOINT.parse::().unwrap()]) .unwrap(); - let s = ztimeout!(zenoh::open(c)).unwrap(); - s + ztimeout!(zenoh::open(c)).unwrap() }; let zid2 = router2.zid(); let peer = { @@ -101,8 +98,7 @@ async fn test_adminspace_read() { .set(vec![MULTICAST_ENDPOINT.parse::().unwrap()]) .unwrap(); c.scouting.multicast.set_enabled(Some(false)).unwrap(); - let s = ztimeout!(zenoh::open(c)).unwrap(); - s + ztimeout!(zenoh::open(c)).unwrap() }; let root = router @@ -292,8 +288,7 @@ async fn test_adminspace_ronly() { .peer .set_mode(Some("linkstate".to_string())) .unwrap(); - let s = ztimeout!(zenoh::open(c)).unwrap(); - s + ztimeout!(zenoh::open(c)).unwrap() }; let zid = router.zid(); @@ -325,8 +320,7 @@ async fn test_adminspace_write() { .peer .set_mode(Some("linkstate".to_string())) .unwrap(); - let s = ztimeout!(zenoh::open(c)).unwrap(); - s + ztimeout!(zenoh::open(c)).unwrap() }; let zid = router.zid(); From cfe752d607861d9038c0319a3cd0fbfb0544df89 Mon Sep 17 00:00:00 2001 From: ChenYing Kuo Date: Wed, 17 Dec 2025 17:47:12 +0800 Subject: [PATCH 03/11] Support no-std Signed-off-by: ChenYing Kuo --- commons/zenoh-protocol/src/core/endpoint.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/commons/zenoh-protocol/src/core/endpoint.rs b/commons/zenoh-protocol/src/core/endpoint.rs index 3c5c6acc60..36c350dd4e 100644 --- a/commons/zenoh-protocol/src/core/endpoint.rs +++ b/commons/zenoh-protocol/src/core/endpoint.rs @@ -11,7 +11,7 @@ // Contributors: // ZettaScale Zenoh Team, // -use alloc::{borrow::ToOwned, format, string::String}; +use alloc::{borrow::ToOwned, format, string::String, vec, vec::Vec}; use core::{borrow::Borrow, convert::TryFrom, fmt, str::FromStr}; use zenoh_result::{bail, zerror, Error as ZError, ZResult}; From d8ef67d9099512430dfa425fc0c091678d23570e Mon Sep 17 00:00:00 2001 From: ChenYing Kuo Date: Thu, 18 Dec 2025 14:54:34 +0800 Subject: [PATCH 04/11] Don't check the address. Signed-off-by: ChenYing Kuo --- commons/zenoh-protocol/src/core/endpoint.rs | 22 --------------------- zenoh/src/net/runtime/orchestrator.rs | 5 ----- 2 files changed, 27 deletions(-) diff --git a/commons/zenoh-protocol/src/core/endpoint.rs b/commons/zenoh-protocol/src/core/endpoint.rs index 36c350dd4e..0447b7bbb0 100644 --- a/commons/zenoh-protocol/src/core/endpoint.rs +++ b/commons/zenoh-protocol/src/core/endpoint.rs @@ -715,28 +715,6 @@ impl EndPoints { EndPoints::Vec(eps) => eps.clone(), } } - - // Helper function to determine if all EndPoint in an EndPoints enum have the same proto/address - pub fn all_endpoints_have_same_proto_addr(&self) -> bool { - match &self { - EndPoints::Single(_) => true, - EndPoints::Vec(endpoints_vec) => { - if endpoints_vec.is_empty() { - return true; - } - let first_ep_proto_addr = { - let p = endpoints_vec[0].protocol().as_str(); - let a = endpoints_vec[0].address().as_str(); - format!("{}/{}", p, a) - }; - endpoints_vec.iter().all(|ep| { - let p = ep.protocol().as_str(); - let a = ep.address().as_str(); - format!("{}/{}", p, a) == first_ep_proto_addr - }) - } - } - } } impl From for EndPoints { diff --git a/zenoh/src/net/runtime/orchestrator.rs b/zenoh/src/net/runtime/orchestrator.rs index ff82c444db..57bb1b7a68 100644 --- a/zenoh/src/net/runtime/orchestrator.rs +++ b/zenoh/src/net/runtime/orchestrator.rs @@ -352,11 +352,6 @@ impl Runtime { async fn connect_peers_single_link(&self, peers: &[EndPoints]) -> ZResult<()> { let mut success_flag = false; for peer_group in peers { - // Check the peer_group has the same proto/host:port. If not, just ignore the group - if !peer_group.all_endpoints_have_same_proto_addr() { - continue; - } - // Try to connect to each peer in the group let mut peers_to_retry = Vec::new(); for peer in peer_group.as_vec() { From 940852f5844f10100fed2c027d2392861fcf702c Mon Sep 17 00:00:00 2001 From: ChenYing Kuo Date: Thu, 18 Dec 2025 14:54:57 +0800 Subject: [PATCH 05/11] Able to accept multiple connections from CLI Signed-off-by: ChenYing Kuo --- examples/src/lib.rs | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/examples/src/lib.rs b/examples/src/lib.rs index ab9d212dbe..0d36f67c47 100644 --- a/examples/src/lib.rs +++ b/examples/src/lib.rs @@ -55,8 +55,14 @@ impl From<&CommonArgs> for Config { } if !args.connect.is_empty() { + // Able to parse multiple endpoints, e.g. "tcp/127.0.0.1:7447?rel=0,tcp/127.0.0.1:7447?rel=1" + let endpoints: Vec> = args + .connect + .iter() + .map(|s| s.split(',').map(String::from).collect()) + .collect(); config - .insert_json5("connect/endpoints", &json!(args.connect).to_string()) + .insert_json5("connect/endpoints", &json!(endpoints).to_string()) .unwrap(); } if !args.listen.is_empty() { From 6688c80b635e1a3006ccab7b65e8d9afa4b5b4fb Mon Sep 17 00:00:00 2001 From: ChenYing Kuo Date: Thu, 18 Dec 2025 16:09:39 +0800 Subject: [PATCH 06/11] Add tests for the client mode config. Signed-off-by: ChenYing Kuo --- commons/zenoh-config/src/lib.rs | 24 +++++++++++++++++++++ commons/zenoh-protocol/src/core/endpoint.rs | 19 +++++++++++++++- zenoh/src/net/runtime/orchestrator.rs | 6 +++--- 3 files changed, 45 insertions(+), 4 deletions(-) diff --git a/commons/zenoh-config/src/lib.rs b/commons/zenoh-config/src/lib.rs index 924fadbe2e..123c568b54 100644 --- a/commons/zenoh-config/src/lib.rs +++ b/commons/zenoh-config/src/lib.rs @@ -28,6 +28,8 @@ pub mod wrappers; #[allow(unused_imports)] use std::convert::TryFrom; +#[allow(unused_imports)] +use std::str::FromStr; // This is a false positive from the rust analyser use std::{ any::Any, @@ -1201,6 +1203,28 @@ fn config_deser() { }) ); + let config = Config::from_deserializer( + &mut json5::Deserializer::from_str( + r#"{ + mode: "client", + connect: { + endpoints: [ + ["tcp/127.0.0.1:7447?rel=0","tcp/127.0.0.1:7448?rel=1"], + ] + } + }"#, + ) + .unwrap(), + ) + .unwrap(); + assert_eq!(*config.mode(), Some(WhatAmI::Client)); + let endpoints = config.connect().endpoints().client().unwrap(); + assert_eq!(endpoints.len(), 1); + assert_eq!( + endpoints[0], + EndPoints::from_str("[tcp/127.0.0.1:7447?rel=0,tcp/127.0.0.1:7448?rel=1]").unwrap(), + ); + dbg!(Config::from_file("../../DEFAULT_CONFIG.json5").unwrap()); } diff --git a/commons/zenoh-protocol/src/core/endpoint.rs b/commons/zenoh-protocol/src/core/endpoint.rs index 0447b7bbb0..5b046cb90d 100644 --- a/commons/zenoh-protocol/src/core/endpoint.rs +++ b/commons/zenoh-protocol/src/core/endpoint.rs @@ -696,7 +696,7 @@ impl EndPoint { } } -#[derive(Clone, Debug, serde::Serialize)] +#[derive(Clone, Debug, PartialEq, Eq, serde::Serialize)] pub enum EndPoints { Single(EndPoint), Vec(Vec), @@ -792,6 +792,23 @@ impl FromStr for EndPoints { #[test] fn endpoints() { + // Single + assert_eq!( + EndPoints::from_str("udp/127.0.0.1:7447").unwrap(), + EndPoints::Single(EndPoint::from_str("udp/127.0.0.1:7447").unwrap()) + ); + // Vec + assert_eq!( + EndPoints::from_str("[udp/127.0.0.1:7447?rel=0,udp/127.0.0.1:7447?rel=1]").unwrap(), + EndPoints::Vec(vec![ + EndPoint::from_str("udp/127.0.0.1:7447?rel=0").unwrap(), + EndPoint::from_str("udp/127.0.0.1:7447?rel=1").unwrap() + ]) + ); +} + +#[test] +fn endpoint() { assert!(EndPoint::from_str("/").is_err()); assert!(EndPoint::from_str("?").is_err()); assert!(EndPoint::from_str("#").is_err()); diff --git a/zenoh/src/net/runtime/orchestrator.rs b/zenoh/src/net/runtime/orchestrator.rs index 57bb1b7a68..0a1d2a66a5 100644 --- a/zenoh/src/net/runtime/orchestrator.rs +++ b/zenoh/src/net/runtime/orchestrator.rs @@ -352,7 +352,7 @@ impl Runtime { async fn connect_peers_single_link(&self, peers: &[EndPoints]) -> ZResult<()> { let mut success_flag = false; for peer_group in peers { - // Try to connect to each peer in the group + // try to connect to each peer in the group let mut peers_to_retry = Vec::new(); for peer in peer_group.as_vec() { let endpoint = peer.clone(); @@ -381,13 +381,13 @@ impl Runtime { { success_flag = true; } - // Any endpoint in the group is available, it's marked as success and break + // any endpoint in the group is available, it's marked as success and break if success_flag { break; } } - // Return error if none of them succeeded + // return error if none of them succeeded if success_flag { Ok(()) } else { From d9095ba891e03b85dfcc45c661b6ac4b497c1a95 Mon Sep 17 00:00:00 2001 From: ChenYing Kuo Date: Tue, 6 Jan 2026 13:38:51 +0800 Subject: [PATCH 07/11] Add the usage inside the config file. Signed-off-by: ChenYing Kuo --- DEFAULT_CONFIG.json5 | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/DEFAULT_CONFIG.json5 b/DEFAULT_CONFIG.json5 index 97bb269265..60c12b0494 100644 --- a/DEFAULT_CONFIG.json5 +++ b/DEFAULT_CONFIG.json5 @@ -44,6 +44,13 @@ /// Accepts a single list (e.g. endpoints: ["tcp/10.10.10.10:7447", "tcp/11.11.11.11:7447"]) /// or different lists for router, peer and client (e.g. endpoints: { router: ["tcp/10.10.10.10:7447"], peer: ["tcp/11.11.11.11:7447"] }). /// + /// Note that every element in the list can also be a list: + /// E.g. endpoints: [["tcp/10.10.10.10:7447?rel=0", "tcp/10.10.10.10:7447?rel=1"], "tcp/11.11.11.11:7447"] + /// This can be used in the client mode. + /// Since the client mode only allows connecting to a single endpoint, this indicates that we want to build multiple links to the same endpoint. + /// It doesn't make any difference for the peer or router mode. + /// endpoints: [["tcp/10.10.10.10:7447?rel=0", "tcp/10.10.10.10:7447?rel=1"]] is equivalent to endpoints: ["tcp/10.10.10.10:7447?rel=0", "tcp/10.10.10.10:7447?rel=1"]. + /// /// See https://docs.rs/zenoh/latest/zenoh/config/struct.EndPoint.html endpoints: [ // "/
" From d1abb0d318a71a9980a090d5152b48b0a11c800e Mon Sep 17 00:00:00 2001 From: ChenYing Kuo Date: Thu, 15 Jan 2026 13:51:26 +0800 Subject: [PATCH 08/11] Revert the change in examples behaviors. Signed-off-by: ChenYing Kuo --- examples/src/lib.rs | 8 +------- 1 file changed, 1 insertion(+), 7 deletions(-) diff --git a/examples/src/lib.rs b/examples/src/lib.rs index 0d36f67c47..ab9d212dbe 100644 --- a/examples/src/lib.rs +++ b/examples/src/lib.rs @@ -55,14 +55,8 @@ impl From<&CommonArgs> for Config { } if !args.connect.is_empty() { - // Able to parse multiple endpoints, e.g. "tcp/127.0.0.1:7447?rel=0,tcp/127.0.0.1:7447?rel=1" - let endpoints: Vec> = args - .connect - .iter() - .map(|s| s.split(',').map(String::from).collect()) - .collect(); config - .insert_json5("connect/endpoints", &json!(endpoints).to_string()) + .insert_json5("connect/endpoints", &json!(args.connect).to_string()) .unwrap(); } if !args.listen.is_empty() { From 4316d5550c6b230a24d12436d871068ed973349e Mon Sep 17 00:00:00 2001 From: ChenYing Kuo Date: Thu, 15 Jan 2026 15:00:17 +0800 Subject: [PATCH 09/11] Update the connect config format. Support object instead of array. Signed-off-by: ChenYing Kuo --- Cargo.lock | 1 + DEFAULT_CONFIG.json5 | 4 +- commons/zenoh-config/src/lib.rs | 10 ++- commons/zenoh-protocol/Cargo.toml | 1 + commons/zenoh-protocol/src/core/endpoint.rs | 78 +++++++++++++-------- 5 files changed, 62 insertions(+), 32 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 9b55bdf640..0afc3e0bac 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -6444,6 +6444,7 @@ dependencies = [ "lazy_static", "rand 0.8.5", "serde", + "serde_json", "uhlc", "zenoh-buffers", "zenoh-keyexpr", diff --git a/DEFAULT_CONFIG.json5 b/DEFAULT_CONFIG.json5 index 60c12b0494..697c159232 100644 --- a/DEFAULT_CONFIG.json5 +++ b/DEFAULT_CONFIG.json5 @@ -45,11 +45,11 @@ /// or different lists for router, peer and client (e.g. endpoints: { router: ["tcp/10.10.10.10:7447"], peer: ["tcp/11.11.11.11:7447"] }). /// /// Note that every element in the list can also be a list: - /// E.g. endpoints: [["tcp/10.10.10.10:7447?rel=0", "tcp/10.10.10.10:7447?rel=1"], "tcp/11.11.11.11:7447"] + /// E.g. endpoints: [{"strategy": "allOf", "locators": ["tcp/10.10.10.10:7447?rel=0", "tcp/10.10.10.10:7447?rel=1"]}, "tcp/11.11.11.11:7447"] /// This can be used in the client mode. /// Since the client mode only allows connecting to a single endpoint, this indicates that we want to build multiple links to the same endpoint. /// It doesn't make any difference for the peer or router mode. - /// endpoints: [["tcp/10.10.10.10:7447?rel=0", "tcp/10.10.10.10:7447?rel=1"]] is equivalent to endpoints: ["tcp/10.10.10.10:7447?rel=0", "tcp/10.10.10.10:7447?rel=1"]. + /// endpoints: [{"strategy": "allOf", "locators": ["tcp/10.10.10.10:7447?rel=0", "tcp/10.10.10.10:7447?rel=1"]}] is equivalent to endpoints: ["tcp/10.10.10.10:7447?rel=0", "tcp/10.10.10.10:7447?rel=1"]. /// /// See https://docs.rs/zenoh/latest/zenoh/config/struct.EndPoint.html endpoints: [ diff --git a/commons/zenoh-config/src/lib.rs b/commons/zenoh-config/src/lib.rs index 123c568b54..e85bfd0313 100644 --- a/commons/zenoh-config/src/lib.rs +++ b/commons/zenoh-config/src/lib.rs @@ -1209,7 +1209,7 @@ fn config_deser() { mode: "client", connect: { endpoints: [ - ["tcp/127.0.0.1:7447?rel=0","tcp/127.0.0.1:7448?rel=1"], + { strategy: "allOf", locators: ["tcp/127.0.0.1:7447?rel=0", "tcp/127.0.0.1:7448?rel=1"] }, ] } }"#, @@ -1222,7 +1222,13 @@ fn config_deser() { assert_eq!(endpoints.len(), 1); assert_eq!( endpoints[0], - EndPoints::from_str("[tcp/127.0.0.1:7447?rel=0,tcp/127.0.0.1:7448?rel=1]").unwrap(), + EndPoints::Locators(zenoh_protocol::core::Locators { + strategy: zenoh_protocol::core::LocatorsStrategy::AllOf, + locators: vec![ + EndPoint::from_str("tcp/127.0.0.1:7447?rel=0").unwrap(), + EndPoint::from_str("tcp/127.0.0.1:7448?rel=1").unwrap() + ] + }) ); dbg!(Config::from_file("../../DEFAULT_CONFIG.json5").unwrap()); diff --git a/commons/zenoh-protocol/Cargo.toml b/commons/zenoh-protocol/Cargo.toml index 2fb510c6f9..07c452e406 100644 --- a/commons/zenoh-protocol/Cargo.toml +++ b/commons/zenoh-protocol/Cargo.toml @@ -52,3 +52,4 @@ zenoh-result = { workspace = true } [dev-dependencies] lazy_static = { workspace = true } rand = { workspace = true, features = ["default"] } +serde_json = { workspace = true } diff --git a/commons/zenoh-protocol/src/core/endpoint.rs b/commons/zenoh-protocol/src/core/endpoint.rs index 5b046cb90d..a2f67278f4 100644 --- a/commons/zenoh-protocol/src/core/endpoint.rs +++ b/commons/zenoh-protocol/src/core/endpoint.rs @@ -696,23 +696,37 @@ impl EndPoint { } } +#[derive(Clone, Copy, Debug, PartialEq, Eq, serde::Serialize, serde::Deserialize)] +#[serde(rename_all = "camelCase")] +pub enum LocatorsStrategy { + AllOf, + OneOf, +} + +#[derive(Clone, Debug, PartialEq, Eq, serde::Serialize, serde::Deserialize)] +pub struct Locators { + pub strategy: LocatorsStrategy, + pub locators: Vec, +} + #[derive(Clone, Debug, PartialEq, Eq, serde::Serialize)] +#[serde(untagged)] pub enum EndPoints { Single(EndPoint), - Vec(Vec), + Locators(Locators), } impl EndPoints { pub fn flatten(self) -> Vec { match self { EndPoints::Single(ep) => vec![ep], - EndPoints::Vec(eps) => eps, + EndPoints::Locators(l) => l.locators, } } pub fn as_vec(&self) -> Vec { match self { EndPoints::Single(ep) => vec![ep.clone()], - EndPoints::Vec(eps) => eps.clone(), + EndPoints::Locators(l) => l.locators.clone(), } } } @@ -734,7 +748,9 @@ impl<'de> serde::Deserialize<'de> for EndPoints { type Value = EndPoints; fn expecting(&self, formatter: &mut fmt::Formatter) -> fmt::Result { - formatter.write_str("a single endpoint string or a list of endpoint strings") + formatter.write_str( + "a single endpoint string or an object with 'strategy' and 'locators'", + ) } fn visit_str(self, v: &str) -> Result @@ -746,15 +762,24 @@ impl<'de> serde::Deserialize<'de> for EndPoints { .map_err(serde::de::Error::custom) } - fn visit_seq(self, mut seq: A) -> Result + fn visit_map(self, map: A) -> Result where - A: serde::de::SeqAccess<'de>, + A: serde::de::MapAccess<'de>, { - let mut vec = Vec::new(); - while let Some(elem) = seq.next_element::()? { - vec.push(elem); + #[derive(serde::Deserialize)] + struct LocatorsHelper { + strategy: LocatorsStrategy, + locators: Vec, } - Ok(EndPoints::Vec(vec)) + + let s = serde::Deserialize::deserialize(serde::de::value::MapAccessDeserializer::new( + map, + ))?; + let helper: LocatorsHelper = s; + Ok(EndPoints::Locators(Locators { + strategy: helper.strategy, + locators: helper.locators, + })) } } @@ -767,18 +792,10 @@ impl TryFrom for EndPoints { fn try_from(s: String) -> Result { const ERR: &str = - "Endpoints must be of the form or [, , ...]"; - if s.starts_with('[') && s.ends_with(']') { - let eps: ZResult> = s[1..s.len() - 1] - .split(',') - .map(|x| EndPoint::from_str(x.trim())) - .collect(); - eps.map(EndPoints::Vec) - } else { - EndPoint::from_str(s.as_str()) - .map(EndPoints::Single) - .map_err(|e| zerror!("{}: {}", ERR, e).into()) - } + "Endpoints must be of the form "; + EndPoint::from_str(s.as_str()) + .map(EndPoints::Single) + .map_err(|e| zerror!("{}: {}", ERR, e).into()) } } @@ -797,13 +814,18 @@ fn endpoints() { EndPoints::from_str("udp/127.0.0.1:7447").unwrap(), EndPoints::Single(EndPoint::from_str("udp/127.0.0.1:7447").unwrap()) ); - // Vec + // Locators + let json = r#"{"strategy": "allOf", "locators": ["udp/127.0.0.1:7447?rel=0", "udp/127.0.0.1:7447?rel=1"]}"#; + let eps: EndPoints = serde_json::from_str(json).unwrap(); assert_eq!( - EndPoints::from_str("[udp/127.0.0.1:7447?rel=0,udp/127.0.0.1:7447?rel=1]").unwrap(), - EndPoints::Vec(vec![ - EndPoint::from_str("udp/127.0.0.1:7447?rel=0").unwrap(), - EndPoint::from_str("udp/127.0.0.1:7447?rel=1").unwrap() - ]) + eps, + EndPoints::Locators(Locators { + strategy: LocatorsStrategy::AllOf, + locators: vec![ + EndPoint::from_str("udp/127.0.0.1:7447?rel=0").unwrap(), + EndPoint::from_str("udp/127.0.0.1:7447?rel=1").unwrap() + ] + }) ); } From 83cd651fdf4133cacea4a05a6a5c53526febbc69 Mon Sep 17 00:00:00 2001 From: ChenYing Kuo Date: Thu, 15 Jan 2026 15:12:24 +0800 Subject: [PATCH 10/11] Fix the Rust format issue Signed-off-by: ChenYing Kuo --- commons/zenoh-protocol/src/core/endpoint.rs | 9 ++++----- 1 file changed, 4 insertions(+), 5 deletions(-) diff --git a/commons/zenoh-protocol/src/core/endpoint.rs b/commons/zenoh-protocol/src/core/endpoint.rs index a2f67278f4..74aad18356 100644 --- a/commons/zenoh-protocol/src/core/endpoint.rs +++ b/commons/zenoh-protocol/src/core/endpoint.rs @@ -772,9 +772,9 @@ impl<'de> serde::Deserialize<'de> for EndPoints { locators: Vec, } - let s = serde::Deserialize::deserialize(serde::de::value::MapAccessDeserializer::new( - map, - ))?; + let s = serde::Deserialize::deserialize( + serde::de::value::MapAccessDeserializer::new(map), + )?; let helper: LocatorsHelper = s; Ok(EndPoints::Locators(Locators { strategy: helper.strategy, @@ -791,8 +791,7 @@ impl TryFrom for EndPoints { type Error = ZError; fn try_from(s: String) -> Result { - const ERR: &str = - "Endpoints must be of the form "; + const ERR: &str = "Endpoints must be of the form "; EndPoint::from_str(s.as_str()) .map(EndPoints::Single) .map_err(|e| zerror!("{}: {}", ERR, e).into()) From 2405e55e727f23a2b5275be5031f9127a305d665 Mon Sep 17 00:00:00 2001 From: ChenYing Kuo Date: Mon, 4 May 2026 15:45:58 +0800 Subject: [PATCH 11/11] test(connectivity): cover allOf locator groups Signed-off-by: ChenYing Kuo --- commons/zenoh-protocol/src/core/endpoint.rs | 3 + zenoh/src/net/runtime/orchestrator.rs | 17 +++- zenoh/tests/connectivity.rs | 102 ++++++++++++++++++++ 3 files changed, 120 insertions(+), 2 deletions(-) diff --git a/commons/zenoh-protocol/src/core/endpoint.rs b/commons/zenoh-protocol/src/core/endpoint.rs index 30fd5c0e19..86b63edd93 100644 --- a/commons/zenoh-protocol/src/core/endpoint.rs +++ b/commons/zenoh-protocol/src/core/endpoint.rs @@ -709,7 +709,10 @@ impl EndPoint { #[derive(Clone, Copy, Debug, PartialEq, Eq, serde::Serialize, serde::Deserialize)] #[serde(rename_all = "camelCase")] pub enum LocatorsStrategy { + /// Open links to all locators in the group. AllOf, + /// Reserved for future support. The runtime currently only implements + /// `AllOf` semantics. OneOf, } diff --git a/zenoh/src/net/runtime/orchestrator.rs b/zenoh/src/net/runtime/orchestrator.rs index 98f7b1cfec..d5d0cd2ebe 100644 --- a/zenoh/src/net/runtime/orchestrator.rs +++ b/zenoh/src/net/runtime/orchestrator.rs @@ -38,8 +38,8 @@ use zenoh_config::{ use zenoh_link::{Locator, LocatorInspector}; use zenoh_protocol::{ core::{ - whatami::WhatAmIMatcher, EndPoint, EndPoints, Metadata, PriorityRange, WhatAmI, - ZenohIdProto, + whatami::WhatAmIMatcher, EndPoint, EndPoints, LocatorsStrategy, Metadata, PriorityRange, + WhatAmI, ZenohIdProto, }, scouting::{HelloProto, Scout, ScoutingBody, ScoutingMessage}, }; @@ -126,6 +126,17 @@ impl StartConditions { } impl Runtime { + fn warn_if_oneof(peer_group: &EndPoints) { + if let EndPoints::Locators(group) = peer_group { + if matches!(group.strategy, LocatorsStrategy::OneOf) { + tracing::warn!( + "connect.endpoints locator groups with strategy=oneOf are not implemented yet; \ + falling back to current allOf behavior" + ); + } + } + } + pub async fn start(&mut self) -> ZResult<()> { match self.whatami() { WhatAmI::Client => self.start_client().await, @@ -365,6 +376,7 @@ impl Runtime { async fn connect_peers_single_link(&self, peers: &[EndPoints]) -> ZResult<()> { let mut success_flag = false; for peer_group in peers { + Self::warn_if_oneof(peer_group); // try to connect to each peer in the group let mut peers_to_retry = Vec::new(); for peer in peer_group.as_vec() { @@ -412,6 +424,7 @@ impl Runtime { async fn connect_peers_multiply_links(&self, peers: &[EndPoints]) -> ZResult<()> { for peer_group in peers { + Self::warn_if_oneof(peer_group); for peer in peer_group.as_vec() { let endpoint = peer.clone(); let retry_config = self.get_connect_retry_config(&endpoint); diff --git a/zenoh/tests/connectivity.rs b/zenoh/tests/connectivity.rs index 9e317bbffe..5decd8752b 100644 --- a/zenoh/tests/connectivity.rs +++ b/zenoh/tests/connectivity.rs @@ -21,6 +21,14 @@ mod tests { }; use zenoh::sample::SampleKind; + #[cfg(feature = "transport_multilink")] + use zenoh::{config::WhatAmI, Session}; + #[cfg(feature = "transport_multilink")] + use zenoh_config::EndPoint; + #[cfg(feature = "transport_multilink")] + use zenoh_protocol::core::{EndPoints, Locators, LocatorsStrategy}; + #[cfg(feature = "transport_multilink")] + use zenoh_test::get_free_tcp_port; use zenoh_test::TestSessions; async fn collect_events(events: &flume::Receiver, timeout: Duration) -> Vec { @@ -35,6 +43,47 @@ mod tests { const SLEEP: Duration = Duration::from_millis(100); + #[cfg(feature = "transport_multilink")] + async fn wait_for_connection_counts( + session: &Session, + expected_transports: usize, + expected_links: usize, + ) { + tokio::time::timeout(Duration::from_secs(5), async { + loop { + let transports = session.info().transports().await.count(); + let links = session.info().links().await.count(); + if transports == expected_transports && links == expected_links { + break; + } + tokio::time::sleep(SLEEP).await; + } + }) + .await + .expect("Timed out waiting for expected transport/link counts"); + } + + #[cfg(feature = "transport_multilink")] + fn get_client_allof_config(locators: Vec) -> zenoh_config::Config { + let mut config = zenoh_config::Config::default(); + config.set_mode(Some(WhatAmI::Client)).unwrap(); + config.scouting.multicast.set_enabled(Some(false)).unwrap(); + config + .transport + .unicast + .set_max_links(locators.len()) + .unwrap(); + config + .connect + .endpoints + .set(vec![EndPoints::Locators(Locators { + strategy: LocatorsStrategy::AllOf, + locators, + })]) + .unwrap(); + config + } + /// Test that transports() returns an iterator of Transport objects #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn test_info_transports() { @@ -248,6 +297,59 @@ mod tests { session1.close().await.unwrap(); } + /// Test that a client configured with an allOf locator group opens + /// multiple links to the same remote session. + #[cfg(feature = "transport_multilink")] + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + async fn client_connect_allof_opens_multiple_links() { + zenoh_util::init_log_from_env_or("error"); + + let mut test_context = TestSessions::new(); + let listener_config = test_context.get_listener_config("tcp/127.0.0.1:0", 2); + let listener = test_context.open_listener_with_cfg(listener_config).await; + let locators = test_context.locators(); + assert_eq!(locators.len(), 2, "Listener should expose 2 locators"); + + let client = test_context + .open_connector_with_cfg(get_client_allof_config(locators)) + .await; + + wait_for_connection_counts(&listener, 1, 2).await; + wait_for_connection_counts(&client, 1, 2).await; + + test_context.close().await; + } + + /// Test that an allOf locator group is best effort today: if one locator + /// is unreachable, the session still opens and the reachable locator + /// yields one link. + #[cfg(feature = "transport_multilink")] + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + async fn client_connect_allof_survives_partial_failure() { + zenoh_util::init_log_from_env_or("error"); + + let mut test_context = TestSessions::new(); + let listener = test_context.open_listener().await; + let mut locators = test_context.locators(); + assert_eq!(locators.len(), 1, "Listener should expose 1 locator"); + + let unreachable_port = get_free_tcp_port(); + locators.push( + format!("tcp/127.0.0.1:{unreachable_port}") + .parse::() + .unwrap(), + ); + + let client = test_context + .open_connector_with_cfg(get_client_allof_config(locators)) + .await; + + wait_for_connection_counts(&listener, 1, 1).await; + wait_for_connection_counts(&client, 1, 1).await; + + test_context.close().await; + } + /// Test that event history works correctly - sends existing transports/links as Put events #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn test_event_history() {