use crate::error::{NetError, Result};
use crate::source_admission::SourceAdmissionEngine;
use zenith_foundation::FramePool;
const FRAME_SIZE: usize = 4096;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum WorkerState {
Created,
Running,
Stopped,
Error,
}
#[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,
pub cpu_usage: f64,
pub queue_depth: u64,
pub latency_p99: u64,
pub throughput: u64,
pub unexpected_dropped: u64,
}
pub struct Worker {
id: u32,
frame_pool: FramePool,
admission: SourceAdmissionEngine,
state: WorkerState,
stats: WorkerStats,
quarantine_on_error: bool,
quarantine_on_deny: bool,
}
impl Worker {
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,
})
}
#[inline]
pub fn id(&self) -> u32 {
self.id
}
#[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
}
pub fn start(&mut self) {
self.state = WorkerState::Running;
}
pub fn stop(&mut self) {
self.state = WorkerState::Stopped;
}
pub fn process_cycle(&mut self) -> Result<u32> {
if self.state != WorkerState::Running {
return Err(NetError::WorkerState {
reason: "worker not running",
});
}
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 构造应成功");
assert!(worker.process_cycle().is_err());
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);
}
}