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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
270 changes: 99 additions & 171 deletions Cargo.lock

Large diffs are not rendered by default.

14 changes: 4 additions & 10 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -65,19 +65,13 @@ panic = "abort"
inherits = "release"
lto = "thin"

[profile.profiling]
inherits = "release"
# We want to crash during simulation
[profile.sim]
inherits = "dev"
opt-level = 2
debug = true
split-debuginfo = "unpacked"
strip = "none"
lto = "off"
codegen-units = 256
incremental = true

# We want to crash during simulation
[profile.sim]
inherits = "dev"
opt-level = 1

[profile.profiling.package."*"]
opt-level = 3
23 changes: 21 additions & 2 deletions pulsebeam-runtime/src/net/udp.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,9 +5,22 @@ use super::{BATCH_SIZE, CHUNK_SIZE, RecvPacketBatch, SendPacketBatch, fmt_bytes}
use quinn_udp::RecvMeta;
use std::{
io::{self, ErrorKind, IoSliceMut},
net::SocketAddr,
net::{IpAddr, SocketAddr},
};

fn normalize_v4_mapped(addr: SocketAddr) -> SocketAddr {
match addr.ip() {
IpAddr::V6(v6) => {
if let Some(v4) = v6.to_ipv4_mapped() {
SocketAddr::new(IpAddr::V4(v4), addr.port())
} else {
addr
}
}
IpAddr::V4(_) => addr,
}
}

pub const SOCKET_SEND_SIZE: usize = 2 * 1024 * 1024;
pub const SOCKET_RECV_SIZE: usize = 4 * 1024 * 1024;

Expand Down Expand Up @@ -56,6 +69,12 @@ pub async fn bind(addr: SocketAddr, external_addr: Option<SocketAddr>) -> io::Re
socket2::Type::DGRAM,
Some(socket2::Protocol::UDP),
)?;

if addr.is_ipv6() {
// Prefer dual-stack listeners so a single IPv6 socket can accept IPv4-mapped peers.
socket2_sock.set_only_v6(false)?;
}

socket2_sock.set_nonblocking(true)?;
socket2_sock.set_reuse_address(true)?;

Expand Down Expand Up @@ -152,7 +171,7 @@ impl UdpTransportReader {
let buf = &self.batch_buffer[base..tail];

out.push(RecvPacketBatch {
src: m.addr,
src: normalize_v4_mapped(m.addr),
dst: self.local_addr,
buf: buf.to_vec(), // Contains the entire GRO block
stride: m.stride, // Downstream will use this to skip through buf
Expand Down
20 changes: 20 additions & 0 deletions pulsebeam-runtime/src/net/udp_scalar.rs
Original file line number Diff line number Diff line change
Expand Up @@ -48,7 +48,27 @@ impl UdpTransport {
}

pub async fn bind(addr: SocketAddr, external_addr: Option<SocketAddr>) -> io::Result<UdpTransport> {
#[cfg(not(feature = "sim"))]
let socket = {
let socket2_sock = socket2::Socket::new(
socket2::Domain::for_address(addr),
socket2::Type::DGRAM,
Some(socket2::Protocol::UDP),
)?;

if addr.is_ipv6() {
// Prefer dual-stack listeners so a single IPv6 socket can accept IPv4-mapped peers.
socket2_sock.set_only_v6(false)?;
}

socket2_sock.set_nonblocking(true)?;
socket2_sock.bind(&addr.into())?;
UdpSocket::from_std(socket2_sock.into())?
};

#[cfg(feature = "sim")]
let socket = UdpSocket::bind(addr).await?;

let socket = Arc::new(socket);
let local_addr = external_addr.unwrap_or(socket.local_addr()?);

Expand Down
181 changes: 150 additions & 31 deletions pulsebeam-runtime/src/system.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,4 @@
use std::net::IpAddr;

use std::net::{IpAddr, Ipv4Addr, Ipv6Addr};
use systemstat::{IpAddr as SysIpAddr, Platform, System};

/// https://stackoverflow.com/questions/77585473/rust-tokio-how-to-handle-more-signals-than-just-sigint-i-e-sigquit
Expand Down Expand Up @@ -45,21 +44,66 @@ pub async fn wait_for_signal() {
wait_for_signal_impl().await
}

pub fn select_host_address() -> IpAddr {
pub fn select_host_addresses() -> Vec<IpAddr> {
#[derive(Default, Clone)]
struct InterfaceCandidates {
external_v4: Option<Ipv4Addr>,
lan_v4: Option<Ipv4Addr>,
external_v6: Option<Ipv6Addr>, // Global Unicast (e.g., 2001::)
ula_v6: Option<Ipv6Addr>, // Unique Local / LAN (fc00::/7)
link_local_v6: Option<Ipv6Addr>, // Link-Local fallback (fe80::/10)
}

impl InterfaceCandidates {
fn best_v4(&self) -> Option<Ipv4Addr> {
self.external_v4.or(self.lan_v4)
}

fn best_v6(&self) -> Option<Ipv6Addr> {
self.external_v6.or(self.ula_v6).or(self.link_local_v6)
}

// Rank 3 = Public/External, Rank 2 = Private/LAN/ULA, Rank 1 = Link-Local, 0 = Empty
fn v4_rank(&self) -> u8 {
if self.external_v4.is_some() {
3
} else if self.lan_v4.is_some() {
2
} else {
0
}
}

fn v6_rank(&self) -> u8 {
if self.external_v6.is_some() {
3
} else if self.ula_v6.is_some() {
2
} else if self.link_local_v6.is_some() {
1
} else {
0
}
}
}

let system = System::new();
let networks = match system.networks() {
Ok(n) => n,
Err(e) => {
tracing::warn!("could not get network interfaces: {e}");
return IpAddr::V4(std::net::Ipv4Addr::LOCALHOST);
return vec![
IpAddr::V4(Ipv4Addr::LOCALHOST),
IpAddr::V6(Ipv6Addr::LOCALHOST),
];
}
};

let mut external_candidates = vec![];
let mut lan_candidates = vec![];
let mut best_iface_name: Option<String> = None;
let mut best_iface: Option<InterfaceCandidates> = None;

for (name, net) in &networks {
// skip virtual / docker / bridge interfaces
// Skip virtual/container management abstractions
if name.starts_with("docker")
|| name.starts_with("veth")
|| name.starts_with("br-")
Expand All @@ -69,38 +113,113 @@ pub fn select_host_address() -> IpAddr {
continue;
}

// optionally restrict to known LAN interface patterns
// if !(name.starts_with("en") || name.starts_with("eth") || name.starts_with("wlp")) {
// tracing::debug!("skipping non-lan interface {}", name);
// continue;
// }
let mut candidates = InterfaceCandidates::default();

for n in &net.addrs {
if let SysIpAddr::V4(ipv4) = n.addr {
if ipv4.is_loopback() {
tracing::debug!("skipping loopback {}: {}", name, ipv4);
continue;
match n.addr {
SysIpAddr::V4(ipv4) => {
if ipv4.is_loopback()
|| ipv4.is_unspecified()
|| ipv4.is_multicast()
|| ipv4.is_link_local()
{
tracing::debug!("skipping loopback/unroutable v4 on {}: {}", name, ipv4);
continue;
}

if !ipv4.is_private() {
if candidates.external_v4.is_none() {
candidates.external_v4 = Some(ipv4);
}
tracing::info!("found candidate external ipv4 on {}: {}", name, ipv4);
} else {
if candidates.lan_v4.is_none() {
candidates.lan_v4 = Some(ipv4);
}
tracing::info!("found candidate lan ipv4 on {}: {}", name, ipv4);
}
}
SysIpAddr::V6(ipv6) => {
if ipv6.is_loopback() || ipv6.is_unspecified() || ipv6.is_multicast() {
tracing::debug!(
"skipping fundamental unroutable ipv6 on {}: {}",
name,
ipv6
);
continue;
}

if !ipv4.is_private() {
external_candidates.push(IpAddr::V4(ipv4));
tracing::info!("found candidate external ip on {}: {}", name, ipv4);
} else {
lan_candidates.push(IpAddr::V4(ipv4));
tracing::info!("found candidate lan ip on {}: {}", name, ipv4);
if ipv6.is_unicast_link_local() {
if candidates.link_local_v6.is_none() {
candidates.link_local_v6 = Some(ipv6);
}
tracing::info!(
"found candidate link-local fallback ipv6 on {}: {}",
name,
ipv6
);
} else if ipv6.is_unique_local() {
if candidates.ula_v6.is_none() {
candidates.ula_v6 = Some(ipv6);
}
tracing::info!("found candidate lan ula ipv6 on {}: {}", name, ipv6);
} else {
if candidates.external_v6.is_none() {
candidates.external_v6 = Some(ipv6);
}
tracing::info!(
"found candidate external global ipv6 on {}: {}",
name,
ipv6
);
}
}
SysIpAddr::Empty | SysIpAddr::Unsupported => {
tracing::debug!("skipping unsupported interface address type on {}", name);
}
}
}

// An interface is valid if it possesses AT LEAST one usable address (v4 or v6)
if candidates.best_v4().is_none() && candidates.best_v6().is_none() {
continue;
}

// Compare using tuple comparison rules: V4 rank takes priority, V6 rank breaks ties.
let replace_best = match &best_iface {
None => true,
Some(current) => {
(candidates.v4_rank(), candidates.v6_rank())
> (current.v4_rank(), current.v6_rank())
}
};

if replace_best {
best_iface_name = Some(name.clone());
best_iface = Some(candidates);
}
}

if let Some(ip) = external_candidates.first() {
tracing::info!("selecting external ip: {}", ip);
*ip
} else if let Some(ip) = lan_candidates.first() {
tracing::info!("selecting lan ip: {}", ip);
*ip
} else {
tracing::warn!("falling back to localhost");
IpAddr::V4(std::net::Ipv4Addr::LOCALHOST)
if let Some(selected) = best_iface {
let mut out = Vec::with_capacity(2);
if let Some(v4) = selected.best_v4() {
out.push(IpAddr::V4(v4));
}
if let Some(v6) = selected.best_v6() {
out.push(IpAddr::V6(v6));
}

tracing::info!(
iface = best_iface_name.unwrap_or_else(|| "<unknown>".to_string()),
?out,
"selected interface host addresses dynamically"
);
return out;
}

tracing::warn!("no active network interfaces detected; returning local fallback anchors");
vec![
IpAddr::V4(Ipv4Addr::LOCALHOST),
IpAddr::V6(Ipv6Addr::LOCALHOST),
]
}
13 changes: 10 additions & 3 deletions pulsebeam-simulator/src/tests/common/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,10 +20,17 @@ pub struct SimClientBuilder {
agent_builder: AgentBuilder,
}

fn http_base_uri(ip: IpAddr, port: u16) -> String {
match ip {
IpAddr::V4(v4) => format!("http://{}:{}", v4, port),
IpAddr::V6(v6) => format!("http://[{}]:{}", v6, port),
}
}

impl SimClientBuilder {
pub async fn bind(ip: IpAddr, server_ip: IpAddr) -> anyhow::Result<Self> {
let client = create_http_client();
let server_base_uri = format!("http://{}:7070", server_ip);
let server_base_uri = http_base_uri(server_ip, 7070);
let api = HttpApiClient::new(client, &server_base_uri)?;

let socket = UdpSocket::bind("0.0.0.0:0").await?;
Expand All @@ -38,11 +45,11 @@ impl SimClientBuilder {
/// port (3478). Use with `start_sfu_node_tcp_only` to test TCP connectivity.
pub async fn bind_tcp(ip: IpAddr, server_ip: IpAddr) -> anyhow::Result<Self> {
let client = create_http_client();
let server_base_uri = format!("http://{}:7070", server_ip);
let server_base_uri = http_base_uri(server_ip, 7070);
let api = HttpApiClient::new(client, &server_base_uri)?;

let socket = UdpSocket::bind("0.0.0.0:0").await?;
let server_tcp_addr: std::net::SocketAddr = format!("{}:3478", server_ip).parse()?;
let server_tcp_addr = std::net::SocketAddr::new(server_ip, 3478);

Ok(Self {
ip,
Expand Down
12 changes: 6 additions & 6 deletions pulsebeam-simulator/src/tests/common/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -45,14 +45,14 @@ pub fn setup_tracing() {

pub async fn start_sfu_node(ip: IpAddr, rng: pulsebeam_runtime::rand::Rng) -> anyhow::Result<()> {
let rtc_port = 3478;
let external_addr: SocketAddr = format!("{}:3478", ip).parse()?;
let external_addr = SocketAddr::new(ip, rtc_port);
let local_addr: SocketAddr = format!("0.0.0.0:{}", rtc_port).parse()?;
let http_api_addr: SocketAddr = "0.0.0.0:7070".parse()?;

pulsebeam::node::NodeBuilder::new()
.workers(1)
.local_addr(local_addr)
.external_addr(external_addr)
.external_addrs(vec![external_addr])
.rng(rng)
.with_udp_mode(UdpMode::Scalar)
.with_http_api(http_api_addr)
Expand All @@ -69,14 +69,14 @@ pub async fn start_sfu_node_tcp_only(
rng: pulsebeam_runtime::rand::Rng,
) -> anyhow::Result<()> {
let rtc_port = 3478;
let external_addr: SocketAddr = format!("{}:3478", ip).parse()?;
let external_addr = SocketAddr::new(ip, rtc_port);
let local_addr: SocketAddr = format!("0.0.0.0:{}", rtc_port).parse()?;
let http_api_addr: SocketAddr = "0.0.0.0:7070".parse()?;

pulsebeam::node::NodeBuilder::new()
.workers(1)
.local_addr(local_addr)
.external_addr(external_addr)
.external_addrs(vec![external_addr])
.rng(rng)
.with_udp_mode(UdpMode::Scalar)
.with_http_api(http_api_addr)
Expand All @@ -98,14 +98,14 @@ pub async fn start_sfu_node_tcp_only_multi_shard(
rng: pulsebeam_runtime::rand::Rng,
) -> anyhow::Result<()> {
let rtc_port = 3478;
let external_addr: SocketAddr = format!("{}:3478", ip).parse()?;
let external_addr = SocketAddr::new(ip, rtc_port);
let local_addr: SocketAddr = format!("0.0.0.0:{}", rtc_port).parse()?;
let http_api_addr: SocketAddr = "0.0.0.0:7070".parse()?;

pulsebeam::node::NodeBuilder::new()
.workers(2)
.local_addr(local_addr)
.external_addr(external_addr)
.external_addrs(vec![external_addr])
.rng(rng)
.with_udp_mode(UdpMode::Scalar)
.with_http_api(http_api_addr)
Expand Down
Loading
Loading