Skip to main content

retina_core/lcore/
rx_core.rs

1use 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
16/// A RxCore polls from `rxqueues` and reduces the stream of packets into
17/// a stream of higher-level network events to be processed by the user.
18pub(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        // TODO: need check to enforce that each core only has same queue types
72        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                    // log::debug!("{:#?}", mbuf);
99                    // log::debug!("Mark: {}", mbuf.mark());
100                    // log::debug!("RSS Hash: 0x{:x}", mbuf.rss_hash());
101                    // log::debug!(
102                    //     "Queue ID: {}, Port ID: {}, Core ID: {}",
103                    //     rxqueue.qid,
104                    //     rxqueue.pid,
105                    //     self.id,
106                    // );
107                    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        // // Deliver remaining data in table from unfinished connections
116        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}