perf: P2 fixes — binary IPs, pre-alloc buffers, build.rs dedup

parse_packet: Replace String IP addresses with [u8; 16] binary in
UserPacket, eliminating 2M heap allocations/sec at 1Mpps. FlowKey
copies bytes directly instead of string-parse-to-bytes round-trip.

xsk_manager: Pre-allocate comp_descs (256) and rx_descs (64) once
before the main loop instead of per-iteration vec![] allocation.

build.rs: Extract duplicate build_ingress_ebpf/build_egress_ebpf into
shared build_ebpf_package(). Add cargo:rerun-if-changed for common/src
to fix stale eBPF build cache when common crate changes.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
DaLaw2 2026-03-21 14:32:48 +08:00
parent b4f49df314
commit fdcc157d59
5 changed files with 57 additions and 182 deletions

View File

@ -8,16 +8,16 @@ use std::time::SystemTime;
use cargo_metadata::{Artifact, CompilerMessage, Message, Metadata, MetadataCommand, Package, Target, TargetKind};
fn main() {
build_ingress_ebpf();
build_egress_ebpf();
build_ebpf_package("ingress-ebpf", "ingress-ebpf");
build_ebpf_package("egress-ebpf", "egress-ebpf");
build_frontend();
}
fn build_ingress_ebpf() {
fn build_ebpf_package(package_name: &str, target_subdir: &str) {
let Metadata { packages, .. } = MetadataCommand::new().no_deps().exec().unwrap();
let ebpf_package = packages
.into_iter()
.find(|Package { name, .. }| **name == "ingress-ebpf")
.find(|Package { name, .. }| **name == *package_name)
.unwrap();
let out_dir = env::var_os("OUT_DIR").unwrap();
@ -42,6 +42,7 @@ fn build_ingress_ebpf() {
let ebpf_dir = manifest_path.parent().unwrap();
println!("cargo:rerun-if-changed={}", ebpf_dir.as_str());
println!("cargo:rerun-if-changed=../common/src");
let mut cmd = Command::new("cargo");
cmd.args([
@ -62,127 +63,7 @@ fn build_ingress_ebpf() {
}
cmd.current_dir(ebpf_dir);
let ebpf_target_dir = out_dir.join("../ingress-ebpf");
cmd.arg("--target-dir").arg(&ebpf_target_dir);
let mut child = cmd
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.unwrap_or_else(|err| panic!("failed to spawn {cmd:?}: {err}"));
let Child { stdout, stderr, .. } = &mut child;
let stderr = stderr.take().unwrap();
let stderr = BufReader::new(stderr);
let stderr = std::thread::spawn(move || {
for line in stderr.lines() {
let line = line.unwrap();
println!("{line}");
}
});
let stdout = stdout.take().unwrap();
let stdout = BufReader::new(stdout);
let mut executables = Vec::new();
for message in Message::parse_stream(stdout) {
#[allow(clippy::collapsible_match)]
match message.expect("valid JSON") {
Message::CompilerArtifact(Artifact {
executable,
target: Target { name, .. },
..
}) => {
if let Some(executable) = executable {
executables.push((name, executable.into_std_path_buf()));
}
}
Message::CompilerMessage(CompilerMessage { message, .. }) => {
for line in message.rendered.unwrap_or_default().split('\n') {
println!("{line}");
}
}
Message::TextLine(line) => {
println!("{line}");
}
_ => {}
}
}
let status = child
.wait()
.unwrap_or_else(|err| panic!("failed to wait for {cmd:?}: {err}"));
assert_eq!(status.code(), Some(0), "{cmd:?} failed: {status:?}");
stderr.join().map_err(std::panic::resume_unwind).unwrap();
for (name, binary) in executables {
let dst = out_dir.join(name);
let _: u64 =
fs::copy(&binary, &dst).unwrap_or_else(|err| panic!("failed to copy {binary:?} to {dst:?}: {err}"));
}
} else {
let Package { targets, .. } = ebpf_package;
for Target { name, kind, .. } in targets {
if *kind != [TargetKind::Bin] {
continue;
}
let dst = out_dir.join(name);
fs::write(&dst, []).unwrap_or_else(|err| panic!("failed to create {dst:?}: {err}"));
}
}
}
fn build_egress_ebpf() {
let Metadata { packages, .. } = MetadataCommand::new().no_deps().exec().unwrap();
let ebpf_package = packages
.into_iter()
.find(|Package { name, .. }| **name == "egress-ebpf")
.unwrap();
let out_dir = env::var_os("OUT_DIR").unwrap();
let out_dir = PathBuf::from(out_dir);
let endian = env::var_os("CARGO_CFG_TARGET_ENDIAN").unwrap();
let target = if endian == "big" {
"bpfeb"
} else if endian == "little" {
"bpfel"
} else {
panic!("unsupported endian={:?}", endian)
};
let build_ebpf = true;
if build_ebpf {
let arch = env::var_os("CARGO_CFG_TARGET_ARCH").unwrap();
let target = format!("{target}-unknown-none");
let Package { manifest_path, .. } = ebpf_package;
let ebpf_dir = manifest_path.parent().unwrap();
println!("cargo:rerun-if-changed={}", ebpf_dir.as_str());
let mut cmd = Command::new("cargo");
cmd.args([
"build",
"-Z",
"build-std=core",
"--bins",
"--message-format=json",
"--release",
"--target",
&target,
]);
cmd.env("CARGO_CFG_BPF_TARGET_ARCH", arch);
cmd.env("CARGO_TERM_COLOR", "always");
for key in ["RUSTUP_TOOLCHAIN", "RUSTC", "RUSTC_WORKSPACE_WRAPPER"] {
cmd.env_remove(key);
}
cmd.current_dir(ebpf_dir);
let ebpf_target_dir = out_dir.join("../egress-ebpf");
let ebpf_target_dir = out_dir.join(format!("../{target_subdir}"));
cmd.arg("--target-dir").arg(&ebpf_target_dir);
let mut child = cmd

View File

@ -240,6 +240,8 @@ impl XskPair {
let mut shutdown_rx = Some(shutdown_rx);
let mut idle_count: u32 = 0;
let mut buffer_pool = BufferPool::new(self.buffer_pool_capacity, self.packet_buffer_size);
let mut comp_descs = vec![FrameDesc::default(); 256];
let mut rx_descs = vec![FrameDesc::default(); 64];
loop {
if let Some(ref mut rx) = shutdown_rx {
@ -253,17 +255,17 @@ impl XskPair {
let mut total_activity = 0;
match self.process_comp_queue() {
match self.process_comp_queue(&mut comp_descs) {
Ok(count) => total_activity += count,
Err(e) => log!(EbpfLog::CompQueueError(format!("{:?}", e))),
}
match self.process_rx_queue(&forward_tx, &mut buffer_pool) {
match self.process_rx_queue(&forward_tx, &mut buffer_pool, &mut rx_descs) {
Ok(count) => total_activity += count,
Err(e) => log!(EbpfLog::RXQueueError(format!("{:?}", e))),
}
match self.process_tx_queue(&forward_rx, &mut buffer_pool) {
match self.process_tx_queue(&forward_rx, &mut buffer_pool, &mut comp_descs) {
Ok(count) => total_activity += count,
Err(e) => log!(EbpfLog::TXQueueError(format!("{:?}", e))),
}
@ -292,10 +294,8 @@ impl XskPair {
})
}
fn process_comp_queue(&mut self) -> Result<usize, EbpfError> {
let mut comp_descs = vec![FrameDesc::default(); 256];
let nb_completed = unsafe { self.comp_queue.consume(&mut comp_descs) };
fn process_comp_queue(&mut self, comp_descs: &mut [FrameDesc]) -> Result<usize, EbpfError> {
let nb_completed = unsafe { self.comp_queue.consume(comp_descs) };
if nb_completed > 0 {
for desc in comp_descs.iter().take(nb_completed) {
@ -306,9 +306,8 @@ impl XskPair {
Ok(nb_completed)
}
fn process_rx_queue(&mut self, forward_tx: &Sender<Vec<u8>>, buffer_pool: &mut BufferPool) -> Result<usize, EbpfError> {
let mut rx_descs = vec![FrameDesc::default(); 64];
let rx_count = unsafe { self.rx.consume(&mut rx_descs) };
fn process_rx_queue(&mut self, forward_tx: &Sender<Vec<u8>>, buffer_pool: &mut BufferPool, rx_descs: &mut [FrameDesc]) -> Result<usize, EbpfError> {
let rx_count = unsafe { self.rx.consume(rx_descs) };
if rx_count > 0 {
let is_ingress = self.direction == Direction::Ingress;
@ -371,7 +370,7 @@ impl XskPair {
Ok(rx_count)
}
fn process_tx_queue(&mut self, forward_rx: &Receiver<Vec<u8>>, buffer_pool: &mut BufferPool) -> Result<usize, EbpfError> {
fn process_tx_queue(&mut self, forward_rx: &Receiver<Vec<u8>>, buffer_pool: &mut BufferPool, comp_descs: &mut [FrameDesc]) -> Result<usize, EbpfError> {
let mut packets_to_send = Vec::with_capacity(64);
while let Ok(packet) = forward_rx.try_recv() {
packets_to_send.push(packet);
@ -384,7 +383,7 @@ impl XskPair {
return Ok(0);
}
if let Err(e) = self.process_comp_queue() {
if let Err(e) = self.process_comp_queue(comp_descs) {
log!(EbpfLog::CompQueueError(format!("{:?}", e)));
}

View File

@ -1,4 +1,4 @@
use std::net::{IpAddr, Ipv4Addr, Ipv6Addr};
use std::net::{Ipv4Addr, Ipv6Addr};
use serde::{Deserialize, Serialize};
use tract_onnx::prelude::{Graph, SimplePlan, TypedFact, TypedOp};
@ -35,15 +35,13 @@ pub struct FlowKey {
impl FlowKey {
pub fn from_packet(packet: &UserPacket) -> Self {
let (src_ip, ip_version) = Self::parse_ip_to_bytes(&packet.src_ip);
let (dst_ip, _) = Self::parse_ip_to_bytes(&packet.dst_ip);
Self {
src_ip,
dst_ip,
src_ip: packet.src_ip,
dst_ip: packet.dst_ip,
src_port: packet.src_port,
dst_port: packet.dst_port,
protocol: packet.protocol,
ip_version,
ip_version: packet.ip_version,
}
}
@ -66,21 +64,6 @@ impl FlowKey {
self.ip_bytes_to_string(&self.dst_ip)
}
fn parse_ip_to_bytes(ip_str: &str) -> ([u8; 16], u8) {
if let Ok(addr) = ip_str.parse::<IpAddr>() {
match addr {
IpAddr::V4(v4) => {
let mut buf = [0u8; 16];
buf[..4].copy_from_slice(&v4.octets());
(buf, 4)
}
IpAddr::V6(v6) => (v6.octets(), 6),
}
} else {
([0u8; 16], 4)
}
}
fn ip_bytes_to_string(&self, bytes: &[u8; 16]) -> String {
if self.ip_version == 6 {
Ipv6Addr::from(*bytes).to_string()

View File

@ -1,10 +1,11 @@
use std::net::{Ipv4Addr, Ipv6Addr};
pub struct UserPacket {
#[allow(dead_code)]
pub ip_version: u8,
pub protocol: u8,
pub tcp_flags: u8,
pub src_ip: String,
pub dst_ip: String,
pub src_ip: [u8; 16],
pub dst_ip: [u8; 16],
pub src_port: u16,
pub dst_port: u16,
pub packet_length: u32,
@ -14,3 +15,23 @@ pub struct UserPacket {
pub timestamp_us: u64,
pub is_forward: bool,
}
impl UserPacket {
/// Format source IP as a human-readable string.
pub fn src_ip_string(&self) -> String {
Self::ip_bytes_to_string(self.ip_version, &self.src_ip)
}
/// Format destination IP as a human-readable string.
pub fn dst_ip_string(&self) -> String {
Self::ip_bytes_to_string(self.ip_version, &self.dst_ip)
}
fn ip_bytes_to_string(ip_version: u8, bytes: &[u8; 16]) -> String {
if ip_version == 6 {
Ipv6Addr::from(*bytes).to_string()
} else {
Ipv4Addr::from([bytes[0], bytes[1], bytes[2], bytes[3]]).to_string()
}
}
}

View File

@ -30,8 +30,10 @@ fn parse_ipv4(packet_data: &[u8], timestamp_us: u64) -> Option<(UserPacket, usiz
let protocol_byte = ip_header[9];
let src_ip_raw = u32::from_be_bytes([ip_header[12], ip_header[13], ip_header[14], ip_header[15]]);
let dst_ip_raw = u32::from_be_bytes([ip_header[16], ip_header[17], ip_header[18], ip_header[19]]);
let mut src_ip = [0u8; 16];
src_ip[..4].copy_from_slice(&ip_header[12..16]);
let mut dst_ip = [0u8; 16];
dst_ip[..4].copy_from_slice(&ip_header[16..20]);
let ihl = (ip_header[0] & 0x0F) as usize * 4;
let total_len = u16::from_be_bytes([ip_header[2], ip_header[3]]) as u32;
@ -71,8 +73,8 @@ fn parse_ipv4(packet_data: &[u8], timestamp_us: u64) -> Option<(UserPacket, usiz
ip_version: 4,
protocol: protocol_byte,
tcp_flags,
src_ip: format_ipv4(src_ip_raw),
dst_ip: format_ipv4(dst_ip_raw),
src_ip,
dst_ip,
src_port,
dst_port,
packet_length: total_len,
@ -95,13 +97,11 @@ fn parse_ipv6(packet_data: &[u8], timestamp_us: u64) -> Option<(UserPacket, usiz
let protocol_byte = ip_header[6];
let mut source_ip_bytes = [0u8; 16];
source_ip_bytes.copy_from_slice(&ip_header[8..24]);
let src_ip_raw = u128::from_be_bytes(source_ip_bytes);
let mut src_ip = [0u8; 16];
src_ip.copy_from_slice(&ip_header[8..24]);
let mut dest_ip_bytes = [0u8; 16];
dest_ip_bytes.copy_from_slice(&ip_header[24..40]);
let dst_ip_raw = u128::from_be_bytes(dest_ip_bytes);
let mut dst_ip = [0u8; 16];
dst_ip.copy_from_slice(&ip_header[24..40]);
let payload_len = u16::from_be_bytes([ip_header[4], ip_header[5]]) as u32;
let total_len = payload_len + 40;
@ -141,8 +141,8 @@ fn parse_ipv6(packet_data: &[u8], timestamp_us: u64) -> Option<(UserPacket, usiz
ip_version: 6,
protocol: protocol_byte,
tcp_flags,
src_ip: format_ipv6(src_ip_raw),
dst_ip: format_ipv6(dst_ip_raw),
src_ip,
dst_ip,
src_port,
dst_port,
packet_length: total_len,
@ -156,12 +156,3 @@ fn parse_ipv6(packet_data: &[u8], timestamp_us: u64) -> Option<(UserPacket, usiz
Some((packet, payload_start))
}
pub fn format_ipv4(addr: u32) -> String {
let bytes = addr.to_be_bytes();
format!("{}.{}.{}.{}", bytes[0], bytes[1], bytes[2], bytes[3],)
}
pub fn format_ipv6(addr: u128) -> String {
let bytes = addr.to_be_bytes();
std::net::Ipv6Addr::from(bytes).to_string()
}