Skip to main content

retina_core/runtime/
offline.rs

1use crate::config::{ConnTrackConfig, OfflineConfig};
2use crate::conntrack::{ConnTracker, TrackerConfig};
3use crate::dpdk;
4use crate::filter::Filter;
5use crate::lcore::{CoreId, SocketId};
6use crate::memory::mbuf::Mbuf;
7use crate::memory::mempool::Mempool;
8use crate::protocols::stream::ParserRegistry;
9use crate::subscription::*;
10
11use std::collections::BTreeMap;
12use std::ffi::CString;
13use std::sync::Arc;
14
15use cpu_time::ProcessTime;
16use pcap::Capture;
17
18pub(crate) struct OfflineRuntime<'a, S>
19where
20    S: Subscribable,
21{
22    pub(crate) mempool_name: String,
23    pub(crate) filter: Filter,
24    pub(crate) subscription: Arc<Subscription<'a, S>>,
25    pub(crate) options: OfflineOptions,
26}
27
28impl<'a, S> OfflineRuntime<'a, S>
29where
30    S: Subscribable,
31{
32    pub(crate) fn new(
33        options: OfflineOptions,
34        mempools: &BTreeMap<SocketId, Mempool>,
35        filter: Filter,
36        subscription: Arc<Subscription<'a, S>>,
37    ) -> Self {
38        let core_id = CoreId(unsafe { dpdk::rte_lcore_id() } as u32);
39        let mempool_name = mempools
40            .get(&core_id.socket_id())
41            .expect("Get offline mempool")
42            .name()
43            .to_string();
44        OfflineRuntime {
45            mempool_name,
46            filter,
47            subscription,
48            options,
49        }
50    }
51
52    pub(crate) fn run(&self) {
53        log::info!(
54            "Launched offline analysis. Processing pcap: {}",
55            self.options.offline.pcap,
56        );
57
58        let mut nb_pkts = 0;
59        let mut nb_bytes = 0;
60
61        let config = TrackerConfig::from(&self.options.conntrack);
62        let registry = ParserRegistry::build::<S>(&self.filter).expect("Unable to build registry");
63        log::debug!("{:#?}", registry);
64        let mut stream_table = ConnTracker::<S::Tracked>::new(config, registry);
65
66        let mempool_raw = self.get_mempool_raw();
67        let pcap = self.options.offline.pcap.as_str();
68        let mut cap = Capture::from_file(pcap).expect("Error opening pcap. Aborting.");
69        let start = ProcessTime::try_now().expect("Getting process time failed");
70        while let Ok(frame) = cap.next() {
71            if frame.header.len as usize > self.options.offline.mtu {
72                continue;
73            }
74            let mbuf = Mbuf::from_bytes(frame.data, mempool_raw)
75                .expect("Unable to allocate mbuf. Try increasing mempool size.");
76            nb_pkts += 1;
77            nb_bytes += mbuf.data_len() as u64;
78
79            S::process_packet(mbuf, &self.subscription, &mut stream_table);
80        }
81
82        // // Deliver remaining data in table
83        stream_table.drain(&self.subscription);
84        let cpu_time = start.elapsed();
85        println!("Processed: {} pkts, {} bytes", nb_pkts, nb_bytes);
86        println!("CPU time: {:?}ms", cpu_time.as_millis());
87    }
88
89    fn get_mempool_raw(&self) -> *mut dpdk::rte_mempool {
90        let cname = CString::new(self.mempool_name.clone()).expect("Invalid CString conversion");
91        unsafe { dpdk::rte_mempool_lookup(cname.as_ptr()) }
92    }
93}
94
95/// Read-only runtime options for the offline core
96#[derive(Debug)]
97pub(crate) struct OfflineOptions {
98    pub(crate) offline: OfflineConfig,
99    pub(crate) conntrack: ConnTrackConfig,
100}