use std::net::{IpAddr, Ipv4Addr, SocketAddr};
pub const LOCALHOST: IpAddr = IpAddr::V4(Ipv4Addr::new(127, 0, 0, 1));
pub const DORA_DAEMON_LOCAL_LISTEN_PORT_DEFAULT: u16 = 53291;
pub const DORA_DAEMON_LOCAL_LISTEN_PORT_ENV: &str = "DORA_DAEMON_LOCAL_LISTEN_PORT";
pub const DORA_COORDINATOR_PORT_WS_DEFAULT: u16 = 6013;
pub const DORA_ZENOH_CONNECT_ENV: &str = "DORA_ZENOH_CONNECT";
pub const DORA_ZENOH_LISTEN_ENV: &str = "DORA_ZENOH_LISTEN";
pub const DORA_ZENOH_LISTEN_EXTRA_ENV: &str = "DORA_ZENOH_LISTEN_EXTRA";
pub const DORA_ZENOH_MULTICAST_ENV: &str = "DORA_ZENOH_MULTICAST";
pub const DORA_RUN_PARENT_PID_ENV: &str = "DORA_RUN_PARENT_PID";
#[cfg(feature = "zenoh")]
pub const ZENOH_CONFIG_PATH_ENV: &str = zenoh::Config::DEFAULT_CONFIG_PATH_ENV;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum MulticastScouting {
#[default]
Allowed,
Disabled,
}
#[cfg(feature = "zenoh")]
fn multicast_disabled_by_env() -> bool {
multicast_disabled_by_value(std::env::var(DORA_ZENOH_MULTICAST_ENV).ok().as_deref())
}
#[cfg(feature = "zenoh")]
pub fn multicast_disabled(requested_off: bool) -> bool {
requested_off || multicast_disabled_by_env()
}
#[cfg(feature = "zenoh")]
fn multicast_disabled_by_value(value: Option<&str>) -> bool {
matches!(
value.map(|v| v.trim().to_ascii_lowercase()).as_deref(),
Some("off" | "0" | "false" | "no")
)
}
#[cfg(feature = "zenoh")]
fn split_endpoints(value: &str) -> impl Iterator<Item = String> + '_ {
value
.split(',')
.map(str::trim)
.filter(|s| !s.is_empty())
.map(String::from)
}
#[cfg(feature = "zenoh")]
fn listen_endpoints_from_env(listen: Option<&str>, extra: Option<&str>) -> Vec<String> {
listen
.into_iter()
.chain(extra)
.flat_map(split_endpoints)
.collect()
}
#[cfg(feature = "zenoh")]
pub const DORA_ZENOH_CONFIG_OVERLAY_ENV: &str = "DORA_ZENOH_CONFIG_OVERLAY";
#[cfg(feature = "zenoh")]
#[derive(Debug, Default, Clone)]
pub struct ZenohOverlay {
pub connect_endpoints: Vec<String>,
pub listen_endpoints: Vec<String>,
rest: serde_json::Map<String, serde_json::Value>,
}
#[cfg(feature = "zenoh")]
impl ZenohOverlay {
pub fn from_env() -> eyre::Result<Option<Self>> {
use eyre::Context;
let Some(path) = std::env::var_os(DORA_ZENOH_CONFIG_OVERLAY_ENV) else {
return Ok(None);
};
let path = std::path::PathBuf::from(path);
let contents = std::fs::read_to_string(&path).wrap_err_with(|| {
format!(
"failed to read the zenoh config overlay at `{}` named by {}",
path.display(),
DORA_ZENOH_CONFIG_OVERLAY_ENV
)
})?;
Self::parse(&contents)
.wrap_err_with(|| format!("invalid zenoh config overlay at `{}`", path.display()))
.map(Some)
}
pub fn parse(json5_str: &str) -> eyre::Result<Self> {
use eyre::{Context, eyre};
let value: serde_json::Value = json5::from_str(json5_str)
.map_err(|err| eyre!("{err}"))
.wrap_err(
"overlay must be a JSON5 object of zenoh config keys, e.g. \
`{ connect: { endpoints: [\"tcp/10.0.0.1:7447\"] } }`",
)?;
let serde_json::Value::Object(mut object) = value else {
eyre::bail!("overlay must be a JSON5 object, not a bare value");
};
let mut overlay = Self::default();
for (key, endpoints) in [
("connect", &mut overlay.connect_endpoints),
("listen", &mut overlay.listen_endpoints),
] {
let Some(serde_json::Value::Object(section)) = object.get_mut(key) else {
continue;
};
let Some(value) = section.remove("endpoints") else {
continue;
};
let serde_json::Value::Array(list) = value else {
eyre::bail!("`{key}.endpoints` must be an array of endpoint strings");
};
for entry in list {
match entry {
serde_json::Value::String(endpoint) => endpoints.push(endpoint),
other => eyre::bail!(
"`{key}.endpoints` must contain endpoint strings, found `{other}`"
),
}
}
if section.is_empty() {
object.remove(key);
}
}
overlay.rest = object;
Ok(overlay)
}
fn apply(&self, config: &mut zenoh::Config) -> eyre::Result<()> {
for (key, value) in &self.rest {
insert_overlay_value(config, key, value)?;
}
Ok(())
}
}
#[cfg(feature = "zenoh")]
fn insert_overlay_value(
config: &mut zenoh::Config,
path: &str,
value: &serde_json::Value,
) -> eyre::Result<()> {
if let serde_json::Value::Object(fields) = value
&& !fields.is_empty()
{
let mut written = true;
for (field, child) in fields {
if insert_overlay_value(config, &format!("{path}/{field}"), child).is_err() {
written = false;
break;
}
}
if written {
return Ok(());
}
}
config
.insert_json5(path, &value.to_string())
.map_err(|err| eyre::eyre!("failed to apply zenoh config overlay key `{path}`: {err}"))
}
#[cfg(feature = "zenoh")]
pub async fn open_zenoh_session(coordinator_addr: Option<IpAddr>) -> eyre::Result<zenoh::Session> {
let (session, _) = open_zenoh_session_with_listen(ZenohSessionParams {
coordinator_addr,
..Default::default()
})
.await?;
Ok(session)
}
#[cfg(feature = "zenoh")]
#[derive(Debug, Default, Clone, Copy)]
pub struct ZenohSessionParams<'a> {
pub coordinator_addr: Option<IpAddr>,
pub listen_endpoint: Option<&'a str>,
pub inter_daemon_peer: Option<&'a str>,
pub connect_endpoints: &'a [String],
pub discovered_connect_endpoints: &'a [String],
pub multicast: MulticastScouting,
}
#[cfg(feature = "zenoh")]
fn coordinator_connect_endpoints(addr: IpAddr) -> String {
let peer = SocketAddr::new(addr, 5456);
format!(r#"{{ router: ["tcp/[::]:7447"], peer: ["tcp/{peer}"] }}"#)
}
#[cfg(feature = "zenoh")]
pub async fn open_zenoh_session_with_listen(
params: ZenohSessionParams<'_>,
) -> eyre::Result<(zenoh::Session, Option<String>)> {
use eyre::{Context, eyre};
use tracing::warn;
let ZenohSessionParams {
coordinator_addr,
listen_endpoint,
inter_daemon_peer,
connect_endpoints,
discovered_connect_endpoints,
multicast,
} = params;
let mut effective_listen_endpoint: Option<String> = None;
let overlay = ZenohOverlay::from_env()?;
let zenoh_session = match std::env::var(zenoh::Config::DEFAULT_CONFIG_PATH_ENV) {
Ok(path) => {
if overlay.is_some() {
eyre::bail!(
"both {} and {} are set. {} replaces the whole zenoh \
configuration, while {} layers onto the one dora computes — \
keep one. Prefer the overlay: it adds your endpoints \
without discarding the direct node-to-node links the daemon \
plans for this dataflow.",
zenoh::Config::DEFAULT_CONFIG_PATH_ENV,
DORA_ZENOH_CONFIG_OVERLAY_ENV,
zenoh::Config::DEFAULT_CONFIG_PATH_ENV,
DORA_ZENOH_CONFIG_OVERLAY_ENV,
);
}
let zenoh_config = zenoh::Config::from_file(&path)
.map_err(|e| eyre!(e))
.wrap_err_with(|| format!("failed to read zenoh config from {path}"))?;
zenoh::open(zenoh_config)
.await
.map_err(|e| eyre!(e))
.context("failed to open zenoh session")?
}
Err(std::env::VarError::NotPresent) => {
let mut zenoh_config = zenoh::Config::default();
let mut connect_eps: Vec<String> = Vec::new();
if let Ok(eps) = std::env::var(DORA_ZENOH_CONNECT_ENV) {
connect_eps.extend(split_endpoints(&eps));
}
connect_eps.extend(connect_endpoints.iter().cloned());
let authoritative_connect_eps = connect_eps.len();
connect_eps.extend(discovered_connect_endpoints.iter().filter_map(|ep| {
match validate_zenoh_endpoint(ep) {
Ok(()) => Some(ep.clone()),
Err(err) => {
warn!("ignoring discovered zenoh endpoint: {err}");
None
}
}
}));
if let Some(peer) = inter_daemon_peer {
connect_eps.push(peer.to_string());
}
if let Some(overlay) = &overlay {
connect_eps.extend(overlay.connect_endpoints.iter().cloned());
}
let mut seen_connect = std::collections::HashSet::new();
let has_authoritative_connect = authoritative_connect_eps > 0
|| inter_daemon_peer.is_some()
|| overlay
.as_ref()
.is_some_and(|o| !o.connect_endpoints.is_empty());
connect_eps.retain(|ep| seen_connect.insert(ep.clone()));
let mut connect_inserted = false;
if !connect_eps.is_empty() {
let json = endpoint_array_json(&connect_eps);
match zenoh_config.insert_json5("connect/endpoints", &json) {
Ok(()) => connect_inserted = true,
Err(err) => {
warn!(
"failed to set zenoh connect/endpoints to {json} ({err}); leaving multicast scouting enabled as fallback"
);
}
}
}
let mut listen_inserted_into_configured: Option<String> = None;
let mut listen_configured = false;
let env_listen_endpoints = listen_endpoints_from_env(
std::env::var(DORA_ZENOH_LISTEN_ENV).ok().as_deref(),
std::env::var(DORA_ZENOH_LISTEN_EXTRA_ENV).ok().as_deref(),
);
let mut listen_eps: Vec<String> = Vec::new();
if let Some(ep) = listen_endpoint {
listen_eps.push(ep.to_string());
}
listen_eps.extend(env_listen_endpoints.iter().cloned());
if let Some(peer) = inter_daemon_peer {
listen_eps.push(peer.to_string());
}
if let Some(overlay) = &overlay {
listen_eps.extend(overlay.listen_endpoints.iter().cloned());
}
if !listen_eps.is_empty() {
let json = endpoint_array_json(&listen_eps);
let listen_inserted = match zenoh_config.insert_json5("listen/endpoints", &json) {
Ok(()) => {
listen_inserted_into_configured = listen_endpoint.map(String::from);
listen_configured = true;
true
}
Err(err) => {
warn!("failed to set zenoh listen/endpoints to {json}: {err}");
false
}
};
if listen_inserted
&& let Err(err) = zenoh_config.insert_json5("listen/exit_on_failure", "false")
{
warn!("failed to set zenoh listen/exit_on_failure: {err}");
}
}
let requested_off =
multicast_disabled(matches!(multicast, MulticastScouting::Disabled));
let multicast_scouting_off = (connect_inserted && has_authoritative_connect)
|| (requested_off && listen_configured);
if multicast_scouting_off
&& let Err(err) = zenoh_config.insert_json5("scouting/multicast/enabled", "false")
{
warn!("failed to disable zenoh scouting/multicast: {err}");
}
if let Some(addr) = coordinator_addr
&& let Err(err) = zenoh_config
.insert_json5("connect/endpoints", &coordinator_connect_endpoints(addr))
{
warn!("failed to set zenoh connect/endpoints for coordinator {addr}: {err}");
}
if let Some(overlay) = &overlay {
overlay.apply(&mut zenoh_config)?;
}
match zenoh::open(zenoh_config).await {
Ok(zenoh_session) => {
let mut verify: Vec<&str> = listen_inserted_into_configured
.iter()
.map(String::as_str)
.collect();
if listen_configured {
verify.extend(env_listen_endpoints.iter().map(String::as_str));
}
if !verify.is_empty() {
let bound_locators: Vec<String> = zenoh_session
.info()
.locators()
.await
.into_iter()
.map(|l| l.as_str().to_string())
.collect();
for requested in verify {
let bound = bound_locators
.iter()
.any(|l| l.split(['?', '#']).next() == Some(requested));
if bound {
if listen_inserted_into_configured.as_deref() == Some(requested) {
effective_listen_endpoint = Some(requested.to_string());
}
} else if multicast_scouting_off {
warn!(
"zenoh session opened but listener for `{requested}` \
did not bind (actually bound: {bound_locators:?}); \
multicast scouting is disabled for this session \
(explicit connect endpoints are set), so peers told \
to dial `{requested}` have no fallback path to reach \
it (#2762)"
);
} else {
warn!(
"zenoh session opened but listener for `{requested}` \
did not bind (actually bound: {bound_locators:?}); \
falling back to multicast scouting for discovery"
);
}
}
}
zenoh_session
}
Err(err) => {
warn!(
"failed to open tuned zenoh session ({err}), retrying with default config"
);
let zenoh_config = zenoh::Config::default();
zenoh::open(zenoh_config)
.await
.map_err(|e| eyre!(e))
.context("failed to open zenoh session")?
}
}
}
Err(std::env::VarError::NotUnicode(_)) => eyre::bail!(
"{} env variable is not valid unicode",
zenoh::Config::DEFAULT_CONFIG_PATH_ENV
),
};
Ok((zenoh_session, effective_listen_endpoint))
}
#[cfg(feature = "zenoh")]
fn endpoint_array_json(endpoints: &[String]) -> String {
serde_json::Value::Array(
endpoints
.iter()
.map(|e| serde_json::Value::String(e.clone()))
.collect(),
)
.to_string()
}
pub fn validate_zenoh_endpoint(endpoint: &str) -> Result<(), String> {
const MAX_LEN: usize = 256;
if endpoint.is_empty() {
return Err("zenoh endpoint must not be empty".to_string());
}
if endpoint.len() > MAX_LEN {
return Err(format!(
"zenoh endpoint is {} bytes, over the {MAX_LEN}-byte limit",
endpoint.len()
));
}
if let Some(bad) = endpoint
.chars()
.find(|c| !matches!(c, 'a'..='z' | 'A'..='Z' | '0'..='9' | '.' | ':' | '/' | '-' | '_' | '[' | ']' | '%' | '?' | '#' | '=' | '&' | '+' | '*' | ','))
{
return Err(format!(
"zenoh endpoint contains the disallowed character {bad:?}"
));
}
Ok(())
}
pub const DORA_ZENOH_LISTEN_PORT_DEFAULT: u16 = 5456;
pub fn zenoh_endpoint(addr: IpAddr, port: u16) -> String {
format!("tcp/{}", SocketAddr::new(addr, port))
}
pub fn reserve_zenoh_endpoint(bind: IpAddr) -> std::io::Result<String> {
let listener = std::net::TcpListener::bind((bind, 0))?;
let port = listener.local_addr()?.port();
drop(listener);
Ok(zenoh_endpoint(bind, port))
}
pub fn reserve_loopback_zenoh_endpoint() -> std::io::Result<String> {
reserve_zenoh_endpoint(LOCALHOST)
}
pub fn zenoh_bind_address_for(coordinator_addr: SocketAddr) -> IpAddr {
if coordinator_addr.ip().is_loopback() {
return LOCALHOST;
}
local_address_toward(coordinator_addr).unwrap_or(LOCALHOST)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct ZenohListen {
pub addr: IpAddr,
pub port: Option<u16>,
}
impl ZenohListen {
pub fn endpoint(&self) -> std::io::Result<String> {
match self.port {
Some(port) => Ok(zenoh_endpoint(self.addr, port)),
None => reserve_zenoh_endpoint(self.addr),
}
}
}
impl std::fmt::Display for ZenohListen {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self.port {
Some(port) => write!(f, "{}", SocketAddr::new(self.addr, port)),
None => write!(f, "{}", self.addr),
}
}
}
impl std::str::FromStr for ZenohListen {
type Err = String;
fn from_str(s: &str) -> Result<Self, Self::Err> {
let listen = if let Ok(socket) = s.parse::<SocketAddr>() {
Self {
addr: socket.ip(),
port: Some(socket.port()),
}
} else if let Ok(addr) = s.parse::<IpAddr>() {
Self { addr, port: None }
} else {
return Err(format!(
"`{s}` is neither an IP address (`10.0.2.100`) nor an address \
with a port (`10.0.2.100:5456`, `[fd7a:1::2]:5456`)"
));
};
if listen.port == Some(0) {
return Err(format!(
"`{s}` names port 0; omit the port to let the OS pick one \
(dora then advertises the concrete port it reserved)"
));
}
validate_zenoh_listen(listen.addr).map_err(|err| format!("{err}"))?;
Ok(listen)
}
}
pub fn validate_zenoh_listen(bind: IpAddr) -> eyre::Result<()> {
if bind.is_unspecified() {
eyre::bail!(
"zenoh listen address must be concrete, not the wildcard `{bind}`. \
Zenoh advertises the address it binds and remote daemons dial exactly \
that, so a wildcard has nothing to advertise. Pass the address other \
daemons should use to reach this host (e.g. its LAN or VPN address)."
);
}
Ok(())
}
fn usable_source(local: IpAddr) -> Option<IpAddr> {
if local.is_unspecified() || local.is_loopback() {
return None;
}
Some(local)
}
fn local_address_toward(target: SocketAddr) -> Option<IpAddr> {
let bind: SocketAddr = if target.is_ipv4() {
(Ipv4Addr::UNSPECIFIED, 0).into()
} else {
(std::net::Ipv6Addr::UNSPECIFIED, 0).into()
};
let socket = std::net::UdpSocket::bind(bind).ok()?;
socket.connect(target).ok()?;
usable_source(socket.local_addr().ok()?.ip())
}
#[cfg(feature = "zenoh")]
pub fn zenoh_output_publish_topic(
dataflow_id: uuid::Uuid,
node_id: &dora_message::id::NodeId,
output_id: &dora_message::id::DataId,
) -> String {
let network_id = "default";
format!("dora/{network_id}/{dataflow_id}/output/{node_id}/{output_id}")
}
#[cfg(feature = "zenoh")]
fn hex_key_segment(id: &dora_message::id::DataId) -> String {
use std::fmt::Write;
let s: &str = id.as_ref();
let mut out = String::with_capacity(s.len() * 2);
for b in s.bytes() {
let _ = write!(out, "{b:02x}");
}
out
}
#[cfg(feature = "zenoh")]
pub fn zenoh_output_schema_topic(
dataflow_id: uuid::Uuid,
node_id: &dora_message::id::NodeId,
output_id: &dora_message::id::DataId,
) -> String {
let network_id = "default";
let output = hex_key_segment(output_id);
format!("dora/{network_id}/{dataflow_id}/schema/{node_id}/{output}")
}
#[cfg(feature = "zenoh")]
pub fn zenoh_output_ack_topic(
dataflow_id: uuid::Uuid,
node_id: &dora_message::id::NodeId,
output_id: &dora_message::id::DataId,
) -> String {
format!(
"{}/@ack",
zenoh_output_publish_topic(dataflow_id, node_id, output_id)
)
}
#[cfg(feature = "zenoh")]
pub fn zenoh_daemon_control_topic(
dataflow_id: uuid::Uuid,
node_id: &dora_message::id::NodeId,
output_id: &dora_message::id::DataId,
) -> String {
let network_id = "default";
format!("dora/{network_id}/{dataflow_id}/control/{node_id}/{output_id}")
}
pub fn dataflow_extension_topic(dataflow_id: &uuid::Uuid, namespace: &str) -> String {
let network_id = "default";
format!("dora/{network_id}/{dataflow_id}/ext/{namespace}")
}
#[cfg(test)]
mod endpoint_validation_tests {
use super::validate_zenoh_endpoint;
#[test]
fn ordinary_locators_are_accepted() {
for ok in [
"tcp/127.0.0.1:7447",
"tcp/10.0.2.100:5456",
"tcp/[fd7a:1::2]:5456",
"udp/192.168.1.1:7447",
"tcp/host.example:7447",
"tcp/10.0.0.1:7447?prio=high",
"tcp/10.0.0.1:7447#iface=eth0",
] {
assert!(validate_zenoh_endpoint(ok).is_ok(), "rejected `{ok}`");
}
}
#[test]
fn quotes_and_escapes_are_rejected() {
for bad in [
r#"tcp/1.2.3.4:1"#.to_string() + "\"",
r#"tcp/1.2.3.4:1","tcp/evil:7447"#.to_string(),
r"tcp/1.2.3.4:1\u0022".to_string(),
"tcp/1.2.3.4:1\n".to_string(),
"tcp/1.2.3.4:1 ".to_string(),
] {
assert!(validate_zenoh_endpoint(&bad).is_err(), "accepted `{bad}`");
}
}
#[test]
fn empty_and_oversized_endpoints_are_rejected() {
assert!(validate_zenoh_endpoint("").is_err());
assert!(validate_zenoh_endpoint(&"a".repeat(257)).is_err());
assert!(validate_zenoh_endpoint(&"a".repeat(256)).is_ok());
}
}
#[cfg(test)]
mod tests {
use super::*;
#[cfg(feature = "zenoh")]
#[test]
fn listen_endpoints_keep_loopback_first_and_merge_the_extra_var() {
assert_eq!(
listen_endpoints_from_env(Some("tcp/127.0.0.1:41000"), Some("tcp/10.0.0.2:41001")),
vec![
"tcp/127.0.0.1:41000".to_string(),
"tcp/10.0.0.2:41001".to_string()
],
"loopback must stay first — order decides which locator a \
same-machine consumer picks, and only loopback carries shared memory"
);
}
#[cfg(feature = "zenoh")]
#[test]
fn listen_endpoints_tolerate_an_absent_or_empty_extra_var() {
assert_eq!(
listen_endpoints_from_env(Some("tcp/127.0.0.1:41000"), None),
vec!["tcp/127.0.0.1:41000".to_string()]
);
assert_eq!(
listen_endpoints_from_env(Some("tcp/127.0.0.1:41000"), Some("")),
vec!["tcp/127.0.0.1:41000".to_string()]
);
assert!(listen_endpoints_from_env(None, None).is_empty());
}
#[cfg(feature = "zenoh")]
#[test]
fn listen_endpoints_still_accept_a_list_in_either_var() {
assert_eq!(
listen_endpoints_from_env(Some("tcp/127.0.0.1:41000, tcp/10.0.0.2:41001"), None),
vec![
"tcp/127.0.0.1:41000".to_string(),
"tcp/10.0.0.2:41001".to_string()
]
);
}
#[cfg(feature = "zenoh")]
#[test]
fn multicast_disable_spellings_are_recognized() {
for value in ["off", "0", "false", "no", "OFF", "False", " off "] {
assert!(
multicast_disabled_by_value(Some(value)),
"{value:?} should disable multicast scouting"
);
}
}
#[cfg(feature = "zenoh")]
#[test]
fn unset_or_unrecognized_multicast_value_keeps_default() {
assert!(!multicast_disabled_by_value(None));
for value in ["on", "1", "true", "yes", "", "maybe"] {
assert!(
!multicast_disabled_by_value(Some(value)),
"{value:?} must not disable multicast scouting"
);
}
}
#[cfg(feature = "zenoh")]
#[test]
fn an_explicit_multicast_disable_is_never_lost() {
assert!(multicast_disabled(true));
}
#[test]
fn reserve_loopback_endpoint_returns_loopback_tcp() {
let endpoint = reserve_loopback_zenoh_endpoint().expect("reservation succeeds");
assert!(
endpoint.starts_with("tcp/127.0.0.1:"),
"expected loopback tcp endpoint, got {endpoint}"
);
let port: u16 = endpoint
.rsplit(':')
.next()
.and_then(|p| p.parse().ok())
.expect("endpoint has a numeric port");
assert!(port > 0, "kernel must hand out a non-zero ephemeral port");
}
#[test]
fn local_coordinator_keeps_the_listener_on_loopback() {
for addr in ["127.0.0.1:6013", "[::1]:6013"] {
let addr: SocketAddr = addr.parse().unwrap();
assert_eq!(
zenoh_bind_address_for(addr),
LOCALHOST,
"a loopback coordinator at {addr} must not move the listener off loopback"
);
}
}
#[test]
fn ipv6_endpoints_are_bracketed() {
match reserve_zenoh_endpoint(IpAddr::V6(std::net::Ipv6Addr::LOCALHOST)) {
Ok(endpoint) => assert!(
endpoint.starts_with("tcp/[::1]:"),
"expected a bracketed IPv6 endpoint, got {endpoint}"
),
Err(e)
if matches!(
e.kind(),
std::io::ErrorKind::AddrNotAvailable | std::io::ErrorKind::Unsupported
) || e.raw_os_error() == Some(97) => {}
Err(e) => panic!("unexpected error reserving ::1 endpoint: {e}"),
}
}
#[cfg(feature = "zenoh")]
#[test]
fn coordinator_connect_endpoints_bracket_ipv6() {
let v6 = coordinator_connect_endpoints(IpAddr::V6(std::net::Ipv6Addr::LOCALHOST));
assert!(
v6.contains(r#"peer: ["tcp/[::1]:5456"]"#),
"IPv6 coordinator peer must be bracketed, got {v6}"
);
assert!(
!v6.contains("tcp/::1:5456"),
"unbracketed IPv6 peer locator is malformed, got {v6}"
);
let v4 = coordinator_connect_endpoints(IpAddr::V4(Ipv4Addr::new(127, 0, 0, 1)));
assert!(
v4.contains(r#"peer: ["tcp/127.0.0.1:5456"]"#),
"IPv4 coordinator peer must be unbracketed, got {v4}"
);
}
#[test]
fn only_concrete_routable_sources_are_advertisable() {
for ok in ["10.0.2.100", "192.168.1.7", "100.64.0.3"] {
let addr: IpAddr = ok.parse().unwrap();
assert_eq!(
usable_source(addr),
Some(addr),
"{ok} is a concrete routable address and must be advertisable"
);
}
for rejected in ["0.0.0.0", "127.0.0.1", "::", "::1"] {
let addr: IpAddr = rejected.parse().unwrap();
assert_eq!(
usable_source(addr),
None,
"{rejected} must never be advertised to peers as a locator"
);
}
}
#[test]
fn wildcard_zenoh_listen_is_rejected() {
for wildcard in ["0.0.0.0", "::"] {
let addr: IpAddr = wildcard.parse().unwrap();
let err = validate_zenoh_listen(addr)
.expect_err("wildcard listen address must be rejected, not accepted");
assert!(
err.to_string().contains("concrete"),
"error should tell the operator to pass a concrete address, got: {err}"
);
}
}
#[test]
fn zenoh_endpoint_brackets_ipv6_and_matches_the_reserved_shape() {
assert_eq!(
zenoh_endpoint("10.0.2.100".parse().unwrap(), 5456),
"tcp/10.0.2.100:5456"
);
assert_eq!(
zenoh_endpoint("::1".parse().unwrap(), 5456),
"tcp/[::1]:5456"
);
let reserved = reserve_zenoh_endpoint(LOCALHOST).expect("loopback reservation");
let port = reserved
.rsplit_once(':')
.expect("reserved endpoint carries a port")
.1
.parse()
.expect("reserved port is numeric");
assert_eq!(zenoh_endpoint(LOCALHOST, port), reserved);
}
#[test]
fn zenoh_listen_parses_address_with_and_without_port() {
let addr_only: ZenohListen = "10.0.2.100".parse().unwrap();
assert_eq!(addr_only.addr, "10.0.2.100".parse::<IpAddr>().unwrap());
assert_eq!(addr_only.port, None);
let with_port: ZenohListen = "10.0.2.100:5456".parse().unwrap();
assert_eq!(with_port.port, Some(5456));
assert_eq!(with_port.endpoint().unwrap(), "tcp/10.0.2.100:5456");
let v6_port: ZenohListen = "[fd7a:1::2]:5456".parse().unwrap();
assert_eq!(v6_port.port, Some(5456));
assert_eq!(v6_port.endpoint().unwrap(), "tcp/[fd7a:1::2]:5456");
let v6_bare: ZenohListen = "fd7a:1::2:5456".parse().unwrap();
assert_eq!(v6_bare.port, None);
}
#[test]
fn zenoh_listen_rejects_unadvertisable_forms() {
for bad in ["0.0.0.0", "0.0.0.0:5456", "::"] {
let err = bad
.parse::<ZenohListen>()
.expect_err("wildcard must be rejected");
assert!(
err.contains("concrete"),
"unexpected error for {bad}: {err}"
);
}
let err = "10.0.2.100:0"
.parse::<ZenohListen>()
.expect_err("port 0 must be rejected");
assert!(err.contains("port 0"), "unexpected error: {err}");
let err = "not-an-address"
.parse::<ZenohListen>()
.expect_err("garbage must be rejected");
assert!(err.contains("neither"), "unexpected error: {err}");
}
#[cfg(feature = "zenoh")]
#[test]
fn an_overlay_splits_off_the_endpoint_lists() {
let overlay = ZenohOverlay::parse(
r#"{
// a router of our own, plus a second listen address
connect: { endpoints: ["tcp/10.0.0.1:7447"], timeout_ms: 5000 },
listen: { endpoints: ["tcp/10.0.0.2:7448"] },
scouting: { multicast: { enabled: true } },
}"#,
)
.expect("valid overlay");
assert_eq!(overlay.connect_endpoints, ["tcp/10.0.0.1:7447"]);
assert_eq!(overlay.listen_endpoints, ["tcp/10.0.0.2:7448"]);
assert!(overlay.rest.contains_key("connect"));
assert!(!overlay.rest.contains_key("listen"));
assert!(overlay.rest.contains_key("scouting"));
}
#[cfg(feature = "zenoh")]
#[test]
fn an_overlay_applies_onto_a_zenoh_config() {
let mut config = zenoh::Config::default();
config
.insert_json5("connect/endpoints", r#"["tcp/127.0.0.1:1"]"#)
.expect("dora's own endpoints");
ZenohOverlay::parse(r#"{ connect: { timeout_ms: 5000 }, mode: "peer" }"#)
.expect("valid overlay")
.apply(&mut config)
.expect("overlay applies");
let rendered = config.to_string();
assert!(
rendered.contains("5000"),
"overlay value missing: {rendered}"
);
assert!(
rendered.contains("tcp/127.0.0.1:1"),
"writing a sibling subkey must not wipe dora's endpoints: {rendered}"
);
}
#[cfg(feature = "zenoh")]
#[test]
fn a_malformed_overlay_is_rejected() {
for (bad, expected) in [
("not an object", "object"),
(
r#"{ connect: { endpoints: "tcp/10.0.0.1:7447" } }"#,
"array",
),
(r#"{ listen: { endpoints: [7447] } }"#, "endpoint strings"),
] {
let err = ZenohOverlay::parse(bad)
.expect_err("must be rejected")
.to_string();
assert!(
err.contains(expected),
"unexpected error for `{bad}`: {err}"
);
}
}
#[test]
fn concrete_zenoh_listen_addresses_are_accepted() {
for ok in ["127.0.0.1", "10.0.2.100", "::1"] {
let addr: IpAddr = ok.parse().unwrap();
assert!(
validate_zenoh_listen(addr).is_ok(),
"{ok} is concrete and must pass validation"
);
}
}
#[cfg(feature = "zenoh")]
#[test]
fn output_and_control_topics_are_distinct() {
use dora_message::id::{DataId, NodeId};
let dataflow_id = uuid::Uuid::nil();
let node = NodeId::from("node".to_string());
let output = DataId::from("out".to_string());
let output_topic = zenoh_output_publish_topic(dataflow_id, &node, &output);
let control_topic = zenoh_daemon_control_topic(dataflow_id, &node, &output);
assert!(
output_topic.contains("/output/"),
"node output key must contain `/output/`, got {output_topic}"
);
assert!(
control_topic.contains("/control/"),
"daemon control key must contain `/control/`, got {control_topic}"
);
assert_ne!(
output_topic, control_topic,
"node output and daemon control must not share a Zenoh key (dora #1992/#2008)"
);
}
#[cfg(feature = "zenoh")]
#[test]
fn per_output_topics_are_distinct() {
use dora_message::id::{DataId, NodeId};
let dataflow_id = uuid::Uuid::nil();
let node = NodeId::from("node".to_string());
let output = DataId::from("out".to_string());
let topics = [
zenoh_output_publish_topic(dataflow_id, &node, &output),
zenoh_output_schema_topic(dataflow_id, &node, &output),
zenoh_output_ack_topic(dataflow_id, &node, &output),
zenoh_daemon_control_topic(dataflow_id, &node, &output),
];
for (i, a) in topics.iter().enumerate() {
for b in &topics[i + 1..] {
assert_ne!(a, b, "per-output zenoh keys must not overlap");
}
}
}
#[cfg(feature = "zenoh")]
#[test]
fn schema_topic_has_no_verbatim_chunks() {
use dora_message::id::{DataId, NodeId};
let dataflow_id = uuid::Uuid::nil();
let node = NodeId::from("node".to_string());
let output = DataId::from("out".to_string());
let topic = zenoh_output_schema_topic(dataflow_id, &node, &output);
assert!(
topic.contains("/schema/"),
"schema side-channel must live under the dedicated `/schema/` plane, got {topic}"
);
assert!(
!topic.contains("/output/"),
"schema side-channel must not nest under the data-topic `/output/` path, got {topic}"
);
for chunk in topic.split('/') {
assert!(
!chunk.starts_with('@'),
"schema topic chunk `{chunk}` must not be verbatim (`@…`); \
otherwise zenoh_ext liveliness tokens fail to parse (#2923)"
);
}
}
#[cfg(feature = "zenoh")]
#[test]
fn schema_topic_does_not_collide_with_nested_data_id() {
use dora_message::id::{DataId, NodeId};
let dataflow_id = uuid::Uuid::nil();
let node = NodeId::from("node".to_string());
let parent = DataId::from("cmd");
let nested = DataId::from("cmd/_schema");
assert_ne!(
zenoh_output_schema_topic(dataflow_id, &node, &parent),
zenoh_output_publish_topic(dataflow_id, &node, &nested),
);
}
#[cfg(feature = "zenoh")]
#[test]
fn ack_topics_of_prefix_outputs_are_distinct() {
use dora_message::id::{DataId, NodeId};
let dataflow_id = uuid::Uuid::nil();
let node = NodeId::from("node".to_string());
let cmd = DataId::from("cmd".to_string());
let cmd_vel = DataId::from("cmd/vel".to_string());
let cmd_ack = zenoh_output_ack_topic(dataflow_id, &node, &cmd);
let cmd_vel_ack = zenoh_output_ack_topic(dataflow_id, &node, &cmd_vel);
assert_ne!(cmd_ack, cmd_vel_ack);
assert!(cmd_ack.ends_with("/cmd/@ack"));
assert!(cmd_vel_ack.ends_with("/cmd/vel/@ack"));
assert!(!cmd_vel_ack.starts_with(&cmd_ack));
assert!(!cmd_ack.starts_with(&cmd_vel_ack));
}
}