retina_core/lcore/
rx_core.rs1use super::CoreId;
2use crate::config::ConnTrackConfig;
3use crate::conntrack::{ConnTracker, TrackerConfig};
4use crate::dpdk;
5use crate::filter::Filter;
6use crate::memory::mbuf::Mbuf;
7use crate::port::{RxQueue, RxQueueType};
8use crate::protocols::stream::ParserRegistry;
9use crate::subscription::*;
10
11use std::sync::atomic::{AtomicBool, Ordering};
12use std::sync::Arc;
13
14use itertools::Itertools;
15
16pub(crate) struct RxCore<'a, S>
19where
20 S: Subscribable,
21{
22 pub(crate) id: CoreId,
23 pub(crate) rxqueues: Vec<RxQueue>,
24 pub(crate) filter: Filter,
25 pub(crate) conntrack: ConnTrackConfig,
26 pub(crate) subscription: Arc<Subscription<'a, S>>,
27 pub(crate) is_running: Arc<AtomicBool>,
28}
29
30impl<'a, S> RxCore<'a, S>
31where
32 S: Subscribable,
33{
34 pub(crate) fn new(
35 core_id: CoreId,
36 rxqueues: Vec<RxQueue>,
37 filter: Filter,
38 conntrack: ConnTrackConfig,
39 subscription: Arc<Subscription<'a, S>>,
40 is_running: Arc<AtomicBool>,
41 ) -> Self {
42 RxCore {
43 id: core_id,
44 rxqueues,
45 filter,
46 conntrack,
47 subscription,
48 is_running,
49 }
50 }
51
52 pub(crate) fn rx_burst(&self, rxqueue: &RxQueue, rx_burst_size: u16) -> Vec<Mbuf> {
53 let mut ptrs = Vec::with_capacity(rx_burst_size as usize);
54 let nb_rx = unsafe {
55 dpdk::rte_eth_rx_burst(
56 rxqueue.pid.raw(),
57 rxqueue.qid.raw(),
58 ptrs.as_mut_ptr(),
59 rx_burst_size,
60 )
61 };
62 unsafe {
63 ptrs.set_len(nb_rx as usize);
64 ptrs.into_iter()
65 .map(Mbuf::new_unchecked)
66 .collect::<Vec<Mbuf>>()
67 }
68 }
69
70 pub(crate) fn rx_loop(&self) {
71 if self.rxqueues[0].ty == RxQueueType::Receive {
73 self.rx_process();
74 } else {
75 self.rx_sink();
76 }
77 }
78
79 fn rx_process(&self) {
80 log::info!(
81 "Launched RX on core {}, polling {}",
82 self.id,
83 self.rxqueues.iter().format(", "),
84 );
85
86 let mut nb_pkts = 0;
87 let mut nb_bytes = 0;
88
89 let config = TrackerConfig::from(&self.conntrack);
90 let registry = ParserRegistry::build::<S>(&self.filter).expect("Unable to build registry");
91 log::debug!("{:#?}", registry);
92 let mut conn_table = ConnTracker::<S::Tracked>::new(config, registry);
93
94 while self.is_running.load(Ordering::Relaxed) {
95 for rxqueue in self.rxqueues.iter() {
96 let mbufs: Vec<Mbuf> = self.rx_burst(rxqueue, 32);
97 for mbuf in mbufs.into_iter() {
98 nb_pkts += 1;
108 nb_bytes += mbuf.data_len() as u64;
109 S::process_packet(mbuf, &self.subscription, &mut conn_table);
110 }
111 }
112 conn_table.check_inactive(&self.subscription);
113 }
114
115 conn_table.drain(&self.subscription);
117
118 log::info!(
119 "Core {} total recv from {}: {} pkts, {} bytes",
120 self.id,
121 self.rxqueues.iter().format(", "),
122 nb_pkts,
123 nb_bytes
124 );
125 }
126
127 fn rx_sink(&self) {
128 log::info!(
129 "Launched SINK on core {}, polling {}",
130 self.id,
131 self.rxqueues.iter().format(", "),
132 );
133
134 let mut nb_pkts = 0;
135 let mut nb_bytes = 0;
136
137 while self.is_running.load(Ordering::Relaxed) {
138 for rxqueue in self.rxqueues.iter() {
139 let mbufs: Vec<Mbuf> = self.rx_burst(rxqueue, 32);
140 for mbuf in mbufs.into_iter() {
141 log::debug!("RSS Hash: 0x{:x}", mbuf.rss_hash());
142 log::debug!(
143 "Queue ID: {}, Port ID: {}, Core ID: {}",
144 rxqueue.qid,
145 rxqueue.pid,
146 self.id,
147 );
148 nb_pkts += 1;
149 nb_bytes += mbuf.data_len() as u64;
150 }
151 }
152 }
153 log::info!(
154 "Sink Core {} total recv from {}: {} pkts, {} bytes",
155 self.id,
156 self.rxqueues.iter().format(", "),
157 nb_pkts,
158 nb_bytes
159 );
160 }
161}