zenith-net 0.1.0

Zenith 网络地址与传输层抽象:L2-L4 协议解析、TCP/UDP/QUIC 状态机、来源准入引擎、单队列 Worker 数据面循环
//! 单队列 Worker 数据面循环(非 Linux 平台桩)
//!
//! AF_XDP/eBPF 数据面仅 Linux 可用。本桩模块在非 Linux 目标(或未启用 `linux`
//! feature)下提供与 [`crate::worker`] 一致的公开 API 表面,保证可移植子集可编译:
//! - 帧池、准入引擎、状态机、统计、隔离配置均为真实可用(平台无关)
//! - `process_cycle` 诚实返回 0 帧:本平台无 AF_XDP 套接字,不存在可处理的报文
//!   (等价于 Linux 模拟模式无 RX 数据的空周期);内核数据面的显式 fail-closed
//!   门禁在 `WorkerRuntime::init_ebpf`(非 Linux 返回错误)

use crate::error::{NetError, Result};
use crate::source_admission::SourceAdmissionEngine;

use zenith_foundation::FramePool;

/// 帧大小(与真实 Worker 的 UMEM 帧大小一致)
const FRAME_SIZE: usize = 4096;

/// Worker 状态
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum WorkerState {
    /// 已创建
    Created,
    /// 正在运行
    Running,
    /// 已停止
    Stopped,
    /// 已出错
    Error,
}

/// Worker 统计信息
#[derive(Debug, Clone, Copy, Default)]
pub struct WorkerStats {
    /// 总接收包数
    pub rx_packets: u64,
    /// 总发送包数
    pub tx_packets: u64,
    /// 准入拒绝包数
    pub rejected_packets: u64,
    /// 解析失败包数
    pub parse_errors: u64,
    /// 帧池分配次数
    pub frame_allocs: u64,
    /// 帧池释放次数
    pub frame_frees: u64,
    /// 隔离区帧数
    pub quarantined_frames: u64,
    /// 处理周期数
    pub cycles_completed: u64,
    /// CPU 使用率(0.0 - 1.0)
    pub cpu_usage: f64,
    /// 当前队列深度
    pub queue_depth: u64,
    /// P99 延迟(微秒)
    pub latency_p99: u64,
    /// 吞吐量(包/秒)
    pub throughput: u64,
    /// 非预期包丢弃数(协议预期检查未通过)
    pub unexpected_dropped: u64,
}

/// 单队列 Worker(非 Linux 桩)
///
/// 帧池/准入/状态机为真实实现;数据面循环 fail-closed。
pub struct Worker {
    /// Worker ID
    id: u32,
    /// 帧池(真实,平台无关)
    frame_pool: FramePool,
    /// 准入引擎(真实,平台无关)
    admission: SourceAdmissionEngine,
    /// Worker 状态
    state: WorkerState,
    /// 统计信息
    stats: WorkerStats,
    /// 解析错误时是否隔离帧
    quarantine_on_error: bool,
    /// 准入拒绝时是否隔离帧
    quarantine_on_deny: bool,
}

impl Worker {
    /// 创建新的 Worker(桩:无 XSK 配置,数据面不可用)
    ///
    /// # 参数
    /// * `id` - Worker 唯一 ID
    /// * `frame_pool_capacity` - 帧池容量
    /// * `admission` - 预配置的准入引擎
    pub fn new(
        id: u32,
        frame_pool_capacity: u32,
        admission: SourceAdmissionEngine,
    ) -> Result<Self> {
        let frame_pool = FramePool::new(
            format!("worker-{}", id),
            frame_pool_capacity,
            FRAME_SIZE as u32,
        );
        Ok(Self {
            id,
            frame_pool,
            admission,
            state: WorkerState::Created,
            stats: WorkerStats::default(),
            quarantine_on_error: false,
            quarantine_on_deny: false,
        })
    }

    /// 获取 Worker ID
    #[inline]
    pub fn id(&self) -> u32 {
        self.id
    }

    /// 获取 Worker 状态
    #[inline]
    pub fn state(&self) -> WorkerState {
        self.state
    }

    /// 获取统计信息
    #[inline]
    pub fn stats(&self) -> WorkerStats {
        self.stats
    }

    /// 获取准入引擎的可变引用
    #[inline]
    pub fn admission_mut(&mut self) -> &mut SourceAdmissionEngine {
        &mut self.admission
    }

    /// 获取准入引擎的不可变引用
    #[inline]
    pub fn admission(&self) -> &SourceAdmissionEngine {
        &self.admission
    }

    /// 设置解析错误时是否隔离帧
    #[inline]
    pub fn set_quarantine_on_error(&mut self, enabled: bool) {
        self.quarantine_on_error = enabled;
    }

    /// 设置准入拒绝时是否隔离帧
    #[inline]
    pub fn set_quarantine_on_deny(&mut self, enabled: bool) {
        self.quarantine_on_deny = enabled;
    }

    /// 获取解析错误隔离配置
    #[inline]
    pub fn quarantine_on_error(&self) -> bool {
        self.quarantine_on_error
    }

    /// 获取准入拒绝隔离配置
    #[inline]
    pub fn quarantine_on_deny(&self) -> bool {
        self.quarantine_on_deny
    }

    /// 获取帧池引用
    #[inline]
    pub fn frame_pool(&self) -> &FramePool {
        &self.frame_pool
    }

    /// 获取帧池可变引用
    #[inline]
    pub fn frame_pool_mut(&mut self) -> &mut FramePool {
        &mut self.frame_pool
    }

    /// 启动 Worker
    pub fn start(&mut self) {
        self.state = WorkerState::Running;
    }

    /// 停止 Worker
    pub fn stop(&mut self) {
        self.state = WorkerState::Stopped;
    }

    /// 执行一个数据面处理周期
    ///
    /// 非 Linux 平台无 AF_XDP 套接字,不存在可处理的报文,诚实返回 0 帧
    /// (等价于 Linux 模拟模式的空周期);帧池守恒校验每周期真实执行。
    pub fn process_cycle(&mut self) -> Result<u32> {
        if self.state != WorkerState::Running {
            return Err(NetError::WorkerState {
                reason: "worker not running",
            });
        }
        // 无内核数据面:零帧周期。所有权守恒校验真实执行(与真实 Worker 一致)
        if !self.verify_conservation() {
            self.state = WorkerState::Error;
            return Err(NetError::Internal("frame ownership conservation violated"));
        }
        self.stats.cycles_completed = self.stats.cycles_completed.saturating_add(1);
        Ok(0)
    }

    /// 验证帧所有权守恒
    #[inline]
    pub fn verify_conservation(&self) -> bool {
        self.frame_pool.verify_conservation().is_ok()
    }

    /// 获取帧池信息(容量、已分配、已隔离)
    pub fn frame_pool_info(&self) -> (u32, u32, u32) {
        (
            self.frame_pool.capacity(),
            self.frame_pool.allocated_count(),
            self.frame_pool.quarantined_count(),
        )
    }
}

impl std::fmt::Debug for Worker {
    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        f.debug_struct("Worker")
            .field("id", &self.id)
            .field("state", &self.state)
            .field("stats", &self.stats)
            .finish()
    }
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn test_stub_process_cycle_zero_frame() {
        let mut worker = Worker::new(0, 16, SourceAdmissionEngine::allow_all())
            .expect("桩 Worker 构造应成功");
        // 未启动:状态错误(fail-closed)
        assert!(worker.process_cycle().is_err());
        // 启动后:无 AF_XDP 套接字,诚实返回 0 帧且守恒校验通过
        worker.start();
        assert_eq!(worker.process_cycle().expect("零帧周期应成功"), 0);
        assert_eq!(worker.stats().cycles_completed, 1);
    }

    #[test]
    fn test_stub_real_subsystems() {
        let mut worker = Worker::new(1, 16, SourceAdmissionEngine::allow_all())
            .expect("桩 Worker 构造应成功");
        assert_eq!(worker.id(), 1);
        assert_eq!(worker.state(), WorkerState::Created);
        assert!(worker.verify_conservation());
        worker.set_quarantine_on_error(true);
        assert!(worker.quarantine_on_error());
        assert!(!worker.quarantine_on_deny());
        let (capacity, allocated, quarantined) = worker.frame_pool_info();
        assert_eq!(capacity, 16);
        assert_eq!(allocated, 0);
        assert_eq!(quarantined, 0);
    }
}