#![cfg(all(feature = "afxdp", target_os = "linux"))]
use std::fmt;
use zenith_net::NetError;
use zenith_net::worker::{UdpDatagram, Worker};
use crate::server::{ProtocolServer, ServerError};
const MAX_DATAGRAMS_PER_CYCLE: usize = 64;
const REPLY_FRAME_BUF_SIZE: usize = 2048;
#[derive(Debug, Clone, Copy, Default)]
pub struct BridgeStats {
pub cycles_completed: u64,
pub datagrams_delivered: u64,
pub l7_responses: u64,
pub tx_frames_queued: u64,
pub tx_dropped: u64,
pub l7_errors: u64,
}
impl BridgeStats {
fn saturating_add(&mut self, other: &BridgeStats) {
self.cycles_completed = self.cycles_completed.saturating_add(other.cycles_completed);
self.datagrams_delivered = self.datagrams_delivered.saturating_add(other.datagrams_delivered);
self.l7_responses = self.l7_responses.saturating_add(other.l7_responses);
self.tx_frames_queued = self.tx_frames_queued.saturating_add(other.tx_frames_queued);
self.tx_dropped = self.tx_dropped.saturating_add(other.tx_dropped);
self.l7_errors = self.l7_errors.saturating_add(other.l7_errors);
}
}
#[derive(Debug)]
pub enum BridgeError {
QuicNotBound,
Net(NetError),
Server(ServerError),
}
impl fmt::Display for BridgeError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
BridgeError::QuicNotBound => write!(
f,
"QUIC server not bound — call bind_quic_socketless() before bridging"
),
BridgeError::Net(e) => write!(f, "datapath error: {e}"),
BridgeError::Server(e) => write!(f, "L7 server error: {e}"),
}
}
}
impl std::error::Error for BridgeError {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
match self {
BridgeError::QuicNotBound => None,
BridgeError::Net(e) => Some(e),
BridgeError::Server(e) => Some(e),
}
}
}
impl From<NetError> for BridgeError {
fn from(e: NetError) -> Self {
BridgeError::Net(e)
}
}
#[derive(Debug)]
pub struct AfxdpH3Bridge {
worker: Worker,
server: ProtocolServer,
stats: BridgeStats,
datagrams_buf: Vec<UdpDatagram>,
}
impl AfxdpH3Bridge {
pub fn new(worker: Worker, server: ProtocolServer) -> Result<Self, BridgeError> {
if !server.is_quic_socketless() {
return Err(BridgeError::QuicNotBound);
}
Ok(Self {
worker,
server,
stats: BridgeStats::default(),
datagrams_buf: Vec::with_capacity(MAX_DATAGRAMS_PER_CYCLE),
})
}
#[cfg(target_os = "linux")]
pub fn register_worker_xsk(
&mut self,
maps: &zenith_ebpf::BpfMaps<'_>,
) -> Result<(), BridgeError> {
self.worker.register_xsk(maps).map_err(BridgeError::Net)
}
pub fn run_cycle(&mut self) -> Result<u32, BridgeError> {
self.datagrams_buf.clear();
let processed = {
let buf = &mut self.datagrams_buf;
self.worker
.process_cycle_with_udp_sink(&mut |d| buf.push(d))?
};
let mut cycle_stats = BridgeStats {
cycles_completed: 1,
datagrams_delivered: self.datagrams_buf.len() as u64,
..BridgeStats::default()
};
for dgram in &self.datagrams_buf {
let responses = match self
.server
.serve_udp_quic_datagram(&dgram.payload, dgram.src_socket_addr())
{
Ok(r) => r,
Err(e) => {
cycle_stats.l7_errors += 1;
tracing::debug!("afxdp bridge: L7 datagram error: {e}");
continue;
}
};
cycle_stats.l7_responses += responses.len() as u64;
for (payload, _to) in responses {
let mut frame_buf = [0u8; REPLY_FRAME_BUF_SIZE];
match dgram.build_reply_frame(&payload, &mut frame_buf) {
Some(len) => match self.worker.queue_tx_data(&frame_buf[..len]) {
Ok(_) => cycle_stats.tx_frames_queued += 1,
Err(e) => {
cycle_stats.tx_dropped += 1;
tracing::debug!("afxdp bridge: TX inject failed: {e}");
}
},
None => {
cycle_stats.tx_dropped += 1;
tracing::debug!("afxdp bridge: reply frame build failed (oversize/IPv6)");
}
}
}
}
self.stats.saturating_add(&cycle_stats);
Ok(processed)
}
#[inline]
pub fn stats(&self) -> BridgeStats {
self.stats
}
#[inline]
pub fn worker(&self) -> &Worker {
&self.worker
}
#[inline]
pub fn worker_mut(&mut self) -> &mut Worker {
&mut self.worker
}
#[inline]
pub fn server(&self) -> &ProtocolServer {
&self.server
}
}
#[cfg(test)]
mod tests {
use super::*;
use zenith_net::source_admission::SourceAdmissionEngine;
use zenith_net::worker::WorkerState;
use zenith_tls::CertGeneration;
use crate::app::{App, success_response};
use crate::quic_server::QuicServerConfig;
use crate::server::ProtocolServer;
const TEST_CERT_PEM: &[u8] = b"-----BEGIN CERTIFICATE-----
MIIDCTCCAfGgAwIBAgIUEUD9CUfA86J7odmW+9fgH5y3ziIwDQYJKoZIhvcNAQEL
BQAwFDESMBAGA1UEAwwJbG9jYWxob3N0MB4XDTI2MDcyOTIyMDAxMFoXDTI3MDcy
OTIyMDAxMFowFDESMBAGA1UEAwwJbG9jYWxob3N0MIIBIjANBgkqhkiG9w0BAQEF
AAOCAQ8AMIIBCgKCAQEAuAcuUu9Ajdh//C3n75jRuI0CAfI1EdX/SzDAUZoK1Vbi
+7ZQJ2y7YXvJEa03/CYMn7Qr7Cj/NW+7HZNvxD4o6gAMrNdB4qIEOEM/mSQwXryH
ELJb/it/mC66Sklm3hYjNx0naPLf/5ZGlXuBjr1Um5dT5V1H9BG/tssO2jFbcGpo
MnbZ16J0LO4Jgh8ojaRqzu408CpFnYLUueNwNeG/T+3LUxFpdWNf5uvu4ErW0R4/
rWH4dPlS6/UktQ92jxHHukikPj8hnDkv0H4TfNn+H3cIRIA76wNqif6UTPYPLtkU
4WfvUKHPsto4SQSvG2nk31wtWJZ7PFIS611KpnbXYQIDAQABo1MwUTAdBgNVHQ4E
FgQUxsaRnqWYFIFm8wcXRUMeUGWNOdYwHwYDVR0jBBgwFoAUxsaRnqWYFIFm8wcX
RUMeUGWNOdYwDwYDVR0TAQH/BAUwAwEB/zANBgkqhkiG9w0BAQsFAAOCAQEAgAzn
7GtDuMU1hLiQdXohLzZh2bu1ySmas3ZoAm/Ydbs6+wQA7CaGZ+Q0gTb0CZZUR60r
6qQvD3JkYqNuakUHZlDWZ6NboO9O3F4SybRR3F9Y08kjog7jORAFqh/QpfVgh2Ep
psDAK2TxIFKaLG4It4wE68x9FhHQZfJIekSe1qJXAd6CYmhbZRCIu/5xScijCpJ7
BO0mTj7WU1u3bj8/7jld4ct5dbtEIHpF4rq/5uSCbqPNvA8sqA4i4+beK3aEiXl1
yCPpOR133rFfOg/t0FiqKy1Vtj+kDJgVhzjPrw1w3L50h6iRmug52ozLYMdyOQOI
rEpfPnlOnBBdLz8hjA==
-----END CERTIFICATE-----
";
const TEST_KEY_PEM: &[u8] = b"-----BEGIN PRIVATE KEY-----
MIIEvgIBADANBgkqhkiG9w0BAQEFAASCBKgwggSkAgEAAoIBAQC4By5S70CN2H/8
LefvmNG4jQIB8jUR1f9LMMBRmgrVVuL7tlAnbLthe8kRrTf8JgyftCvsKP81b7sd
k2/EPijqAAys10HiogQ4Qz+ZJDBevIcQslv+K3+YLrpKSWbeFiM3HSdo8t//lkaV
e4GOvVSbl1PlXUf0Eb+2yw7aMVtwamgydtnXonQs7gmCHyiNpGrO7jTwKkWdgtS5
43A14b9P7ctTEWl1Y1/m6+7gStbRHj+tYfh0+VLr9SS1D3aPEce6SKQ+PyGcOS/Q
fhN82f4fdwhEgDvrA2qJ/pRM9g8u2RThZ+9Qoc+y2jhJBK8baeTfXC1Ylns8UhLr
XUqmdtdhAgMBAAECggEAEh9y6mvxWYa2o+kJbEkKbjhEuFhX7Ze7enYkmmSnKHdU
ByHfJuLIWUNNe9YpK0W7/IZLxQgMigCk1rbMTPEqKlEy7lqMfHskGz5UJwqvUMUU
MArAkHlMKXqAkgxEex6G/Uh7txQkBxGPhe0RxzLSADiY5H+ZNGoDDUdWARrXPGyz
ZWlNY/nwg5qKO8zQePNAtUvHUWtNq3G8achi96NsFuuxuXwLKOx60JCMcg5BpRpb
2v8F+2r0ZjngkhI+wYGE7E5E/BdbbooHb1MFyajc+C/GH+sOhSfMMthjwKiSSaIB
/Jpqp9ETJ4GsxTxZTTOszyM5yPBIiPIeuwR0LwXRWQKBgQDjwBJuUqNl0CNMj7I7
6FmPrylzY7JKQB3q1SXqVPKeuPArd2BNBjGX6Ic9k0lLtmcDcsHnKNSXZaTb6UQw
APFlsdWBkWfaNeVh7SLlhez1MWwryuWCw/V8RdMR8MiOcZ+QfjPBLpUgJi/e6CLE
kvv8806vSKwhXQf1Ho43cmhrvwKBgQDO2sOBcYHupfjqOzWewJfR+rOs9F3/YYr/
fZQpdRu4z8zPfTfeR7O8NhvTtNJuzwi/yOjas1DEriBql8zUVDRHJZibj43qvuG6
Ua/cTmbcnyiEv15V0nmv91rEPgTtaf6/iz/zVvlKYT0Zj3wIxpX3rTiw6Ba3wD7I
81EM4K4E3wKBgQCVtmgnP4mL3xOlO3y1ptpg+ossAChuaNGB0lXXQbovnnC6kgGr
AFxPeJqWXqC69Y+oE6LlWtDNKRMDQMcDK0uERy+LudLj/bPo+KKM8MnAsJlj/D99
A2X3KEtEqtybzpNOv7cz0XRUKuYjCMP6JokhUauyy/njAK2/czOXvUxpLwKBgQCX
h+haRdViBpGevPsdrYZKGzZOR8EoGMOjP9IuwIwrSYaGpPstSSdgg97EqpzQ8bc+
DyaNN3i+a7RxgXxaOskFKYRuyK20vlpLjBWg9IojqjAbdrjbc9ES18fVJH2lkdU9
afvR/e+mzi7dL6A0KY2on2t9JLenqhwURzIjld/EzwKBgEoHKFsBZQqYoTBjw46g
VP3U0WBNPxtLlFNFWTh7SVAXTDkSOi4j5YeSEzGt4EHdd+RSVSjap0Yq+kLE3ZzB
hDJ2OdlHsaaSJcZ9+Vnnj44fSAJUKZcAQhq9CfZKcxvQiDuxTP3GZRl0urUdiedO
n6Pgv5yWjDwcNLUlPWhljvBK
-----END PRIVATE KEY-----
";
fn make_worker() -> Worker {
let config = zenith_net::worker::XskConfig {
ifindex: 0,
queue_id: 0,
zero_copy: false,
fill_ring_size: 256,
rx_ring_size: 256,
tx_ring_size: 256,
completion_ring_size: 256,
shared_umem: false,
frame_size: 4096,
headroom: 0,
..Default::default()
};
Worker::new(0, config, 256, SourceAdmissionEngine::allow_all())
.expect("模拟 Worker 构造必须成功")
}
fn make_server() -> ProtocolServer {
let mut app = App::new();
app.get("/", |_req, _rm| {
Ok(success_response("hello from h3", "text/plain"))
});
let mut server = ProtocolServer::new(app);
let cert = CertGeneration::from_pem(TEST_CERT_PEM, TEST_KEY_PEM)
.expect("测试证书必须可解析");
let cfg = QuicServerConfig {
bind_addr: "10.0.0.1:443".parse().expect("合法地址"),
..QuicServerConfig::default()
};
server
.bind_quic_socketless(cfg, &cert)
.expect("socketless 绑定必须成功");
server
}
fn build_udp_frame(payload: &[u8]) -> Vec<u8> {
let ip_total = (20 + 8 + payload.len()) as u16;
let udp_len = (8 + payload.len()) as u16;
let mut frame = vec![0u8; 14 + ip_total as usize];
frame[0..6].copy_from_slice(&[0x11, 0x22, 0x33, 0x44, 0x55, 0x66]); frame[6..12].copy_from_slice(&[0xAA, 0xBB, 0xCC, 0xDD, 0xEE, 0xFF]); frame[12] = 0x08;
frame[13] = 0x00;
let ip = 14;
frame[ip] = 0x45;
frame[ip + 2..ip + 4].copy_from_slice(&ip_total.to_be_bytes());
frame[ip + 8] = 64;
frame[ip + 9] = 17; frame[ip + 12..ip + 16].copy_from_slice(&[172, 16, 0, 1]); frame[ip + 16..ip + 20].copy_from_slice(&[10, 0, 0, 1]);
let l4 = ip + 20;
frame[l4..l4 + 2].copy_from_slice(&12345u16.to_be_bytes());
frame[l4 + 2..l4 + 4].copy_from_slice(&443u16.to_be_bytes());
frame[l4 + 4..l4 + 6].copy_from_slice(&udp_len.to_be_bytes());
frame[l4 + 8..].copy_from_slice(payload);
frame
}
fn build_initial_like_payload() -> Vec<u8> {
let mut p = Vec::with_capacity(64);
p.push(0xC3); p.extend_from_slice(&1u32.to_be_bytes()); p.push(8); p.extend_from_slice(&[0xC1, 0xC2, 0xC3, 0xC4, 0xC5, 0xC6, 0xC7, 0xC8]); p.push(8); p.extend_from_slice(&[0xD1, 0xD2, 0xD3, 0xD4, 0xD5, 0xD6, 0xD7, 0xD8]); p.push(0); p.extend_from_slice(&[0x40, 0x14]); p.extend_from_slice(&[0, 0, 0, 0]); p.extend_from_slice(&[0xAB; 16]); p
}
#[test]
fn test_bridge_requires_socketless_quic() {
let worker = make_worker();
let app = App::new();
let server = ProtocolServer::new(app);
match AfxdpH3Bridge::new(worker, server) {
Err(BridgeError::QuicNotBound) => {}
other => panic!("未绑定 QUIC 时必须返回 QuicNotBound,实际: {}", match other {
Ok(_) => "Ok".to_string(),
Err(e) => e.to_string(),
}),
}
}
#[test]
fn test_socketless_bind_has_no_kernel_socket() {
let server = make_server();
assert!(server.quic_local_addr().is_none());
}
#[test]
fn test_bridge_cycle_delivers_datagram_to_l7() {
let mut worker = make_worker();
worker.start();
worker.process_cycle().expect("预填充周期必须成功");
let server = make_server();
let mut bridge = AfxdpH3Bridge::new(worker, server).expect("桥构造必须成功");
let frame = build_udp_frame(&build_initial_like_payload());
bridge.worker_mut().simulate_rx_transfer(1, frame.len(), |_idx, buf| {
buf[..frame.len()].copy_from_slice(&frame);
});
let processed = bridge.run_cycle().expect("桥周期必须成功");
assert_eq!(processed, 1, "Worker 必须处理 1 个包");
let stats = bridge.stats();
assert_eq!(stats.cycles_completed, 1);
assert_eq!(stats.datagrams_delivered, 1, "数据报必须交付 L7");
assert_eq!(stats.l7_errors, 0, "CONNECTION_CLOSE 不应计为 L7 错误");
assert!(stats.l7_responses >= 1, "CONNECTION_CLOSE 响应应生成 >=1 个响应包");
assert!(stats.tx_frames_queued >= 1, "CONNECTION_CLOSE 应注入 TX");
}
#[test]
fn test_bridge_admission_deny_blocks_delivery() {
let config = zenith_net::worker::XskConfig {
ifindex: 0,
queue_id: 0,
zero_copy: false,
fill_ring_size: 256,
rx_ring_size: 256,
tx_ring_size: 256,
completion_ring_size: 256,
shared_umem: false,
frame_size: 4096,
headroom: 0,
..Default::default()
};
let mut worker = Worker::new(
0,
config,
256,
SourceAdmissionEngine::deny_all(),
)
.expect("Worker 构造必须成功");
worker.start();
worker.process_cycle().expect("预填充周期必须成功");
let server = make_server();
let mut bridge = AfxdpH3Bridge::new(worker, server).expect("桥构造必须成功");
let frame = build_udp_frame(&build_initial_like_payload());
bridge.worker_mut().simulate_rx_transfer(1, frame.len(), |_idx, buf| {
buf[..frame.len()].copy_from_slice(&frame);
});
let processed = bridge.run_cycle().expect("桥周期必须成功");
assert_eq!(processed, 1, "拒绝的包也计入处理数");
let stats = bridge.stats();
assert_eq!(stats.datagrams_delivered, 0, "准入拒绝的数据报禁止交付 L7");
assert_eq!(bridge.worker().stats().rejected_packets, 1);
}
#[test]
fn test_bridge_tx_injection_path() {
let mut worker = make_worker();
worker.start();
let dgram = UdpDatagram {
src_ip: zenith_net::IpAddr::V4([172, 16, 0, 1]),
dst_ip: zenith_net::IpAddr::V4([10, 0, 0, 1]),
src_port: 12345,
dst_port: 443,
src_mac: [0xAA, 0xBB, 0xCC, 0xDD, 0xEE, 0xFF],
dst_mac: [0x11, 0x22, 0x33, 0x44, 0x55, 0x66],
payload: smallvec::SmallVec::new(),
};
let payload = b"quic-response-bytes";
let mut frame_buf = [0u8; 2048];
let len = dgram
.build_reply_frame(payload, &mut frame_buf)
.expect("回包帧构造必须成功");
let frame_idx = worker
.queue_tx_data(&frame_buf[..len])
.expect("TX 注入必须成功");
assert!(frame_idx < 256, "返回值是注入帧的索引,必须落在帧池容量内");
worker.simulate_tx_complete(1);
let stats = worker.stats();
assert_eq!(stats.tx_packets, 0, "queue_tx_data 不计入 rx/tx 转发统计(独立注入路径)");
}
#[test]
fn test_bridge_worker_state_passthrough() {
let worker = make_worker();
let server = make_server();
let mut bridge = AfxdpH3Bridge::new(worker, server).expect("桥构造必须成功");
assert!(bridge.run_cycle().is_err());
bridge.worker_mut().start();
assert_eq!(bridge.worker().state(), WorkerState::Running);
assert!(bridge.run_cycle().is_ok(), "空周期必须成功");
assert_eq!(bridge.stats().cycles_completed, 1);
}
}