use std::net::{IpAddr, Ipv4Addr, SocketAddr};
use anyhow::{Context, Result, bail};
use serde::{Deserialize, Serialize};
use validator::Validate;
use crate::discovery::DiscoveryConfig;
#[derive(Debug, Clone, Default, Serialize, Deserialize, Validate)]
pub struct MessengerConfig {
#[validate(nested)]
pub backend: MessengerBackendConfig,
#[serde(default)]
pub discovery: Option<DiscoveryConfig>,
}
impl MessengerConfig {
pub async fn build_messenger(&self) -> Result<std::sync::Arc<velo::Messenger>> {
use std::net::TcpListener;
use std::sync::Arc;
use velo::Messenger;
use velo::backend::tcp::TcpTransportBuilder;
let bind_addr = self.backend.resolve_bind_addr()?;
let listener = TcpListener::bind(bind_addr)
.with_context(|| format!("Failed to bind TCP listener to {}", bind_addr))?;
let actual_addr = listener
.local_addr()
.context("Failed to get local address from listener")?;
tracing::info!("Built TCP transport bound to {}", actual_addr);
let tcp_transport = TcpTransportBuilder::new()
.from_listener(listener)?
.build()
.context("Failed to build TCP transport")?;
let tcp_transport = Arc::new(tcp_transport);
let mut builder = Messenger::builder().add_transport(tcp_transport);
if let Some(discovery_config) = &self.discovery {
match discovery_config {
DiscoveryConfig::Etcd(_cfg) => {
bail!("Etcd discovery not yet supported in velo");
}
DiscoveryConfig::P2p(_cfg) => {
bail!("P2P discovery not yet supported in velo");
}
DiscoveryConfig::Filesystem(cfg) => {
use velo::discovery::FilesystemPeerDiscovery;
let peer_discovery = FilesystemPeerDiscovery::new(&cfg.path)
.context("Failed to build filesystem discovery")?;
builder = builder.discovery(Arc::new(peer_discovery));
tracing::info!("Built filesystem discovery from: {:?}", cfg.path);
}
}
}
let messenger = builder.build().await.context("Failed to build Messenger")?;
Ok(messenger)
}
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, Validate)]
pub struct MessengerBackendConfig {
pub tcp_addr: Option<String>,
pub tcp_interface: Option<String>,
#[serde(default)]
pub tcp_port: u16,
}
impl MessengerBackendConfig {
pub fn resolve_bind_addr(&self) -> Result<SocketAddr> {
let ip = match (&self.tcp_addr, &self.tcp_interface) {
(Some(_), Some(_)) => {
bail!("tcp_addr and tcp_interface are mutually exclusive")
}
(Some(addr), None) => addr
.parse::<IpAddr>()
.with_context(|| format!("Invalid IP address: {}", addr))?,
(None, Some(iface)) => get_interface_ip(iface)
.with_context(|| format!("Failed to get IP for interface: {}", iface))?,
(None, None) => IpAddr::V4(Ipv4Addr::UNSPECIFIED),
};
Ok(SocketAddr::new(ip, self.tcp_port))
}
}
fn get_interface_ip(interface_name: &str) -> Result<IpAddr> {
use nix::ifaddrs::getifaddrs;
let addrs = getifaddrs().context("Failed to get interface addresses")?;
for ifaddr in addrs {
if ifaddr.interface_name == interface_name
&& let Some(addr) = ifaddr.address
{
if let Some(sockaddr) = addr.as_sockaddr_in() {
return Ok(IpAddr::V4(sockaddr.ip()));
}
if let Some(sockaddr) = addr.as_sockaddr_in6() {
return Ok(IpAddr::V6(sockaddr.ip()));
}
}
}
bail!("No IP address found for interface: {}", interface_name)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_default_backend_config() {
let config = MessengerBackendConfig::default();
assert!(config.tcp_addr.is_none());
assert!(config.tcp_interface.is_none());
assert_eq!(config.tcp_port, 0);
}
#[test]
fn test_resolve_bind_addr_default() {
let config = MessengerBackendConfig::default();
let addr = config.resolve_bind_addr().unwrap();
assert_eq!(addr.ip(), IpAddr::V4(Ipv4Addr::UNSPECIFIED));
assert_eq!(addr.port(), 0);
}
#[test]
fn test_resolve_bind_addr_explicit() {
let config = MessengerBackendConfig {
tcp_addr: Some("192.168.1.100".to_string()),
tcp_interface: None,
tcp_port: 8080,
};
let addr = config.resolve_bind_addr().unwrap();
assert_eq!(addr.ip(), IpAddr::V4(Ipv4Addr::new(192, 168, 1, 100)));
assert_eq!(addr.port(), 8080);
}
#[test]
fn test_resolve_bind_addr_mutual_exclusivity() {
let config = MessengerBackendConfig {
tcp_addr: Some("0.0.0.0".to_string()),
tcp_interface: Some("eth0".to_string()),
tcp_port: 0,
};
let result = config.resolve_bind_addr();
assert!(result.is_err());
assert!(
result
.unwrap_err()
.to_string()
.contains("mutually exclusive")
);
}
}