DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
IfaceWrapper.hpp
Go to the documentation of this file.
1
9#ifndef DPDKLIBS_SRC_IFACEWRAPPER_HPP_
10#define DPDKLIBS_SRC_IFACEWRAPPER_HPP_
11
12// #include "dpdklibs/nicreader/Structs.hpp"
14
15#include "SourceConcept.hpp"
16#include "dpdklibs/EALSetup.hpp"
18#include "dpdklibs/arp/ARP.hpp"
22
23#include <confmodel/Session.hpp>
24// #include <confmodel/NetworkDevice.hpp>
27
28#include <nlohmann/json.hpp>
29
30#include "logging/Logging.hpp" // NOTE: if ISSUES ARE DECLARED BEFORE include logging/Logging.hpp, TLOG_DEBUG<<issue wont work.
31#include <ers/ers.hpp>
32
33#include <memory>
34#include <set>
35#include <sstream>
36#include <string>
37
38#include <folly/ProducerConsumerQueue.h>
39
40namespace dunedaq {
41
42ERS_DECLARE_ISSUE(dpdklibs, MetricPublishFailed, "Field " << field << " was not reported", ((std::string)field))
43
46 "Unexpected stream ID " << src_id << " in UDP payoad. Total counter: " << counter,
47 ((int)src_id)((size_t)counter))
48
49namespace dpdklibs {
50
51class IfaceWrapper : public opmonlib::MonitorableObject
52{
53public:
54 using sid_to_source_map_t = std::map<int, std::shared_ptr<SourceConcept>>;
55
56 IfaceWrapper(uint iface_id,
57 const appmodel::DPDKReceiver* receiver,
58 const std::vector<const appmodel::NWDetDataSender*>& senders,
59 const std::vector<const confmodel::DetectorStream*>& active_streams,
60 sid_to_source_map_t& sources,
61 std::atomic<bool>& run_marker);
62 ~IfaceWrapper();
63
64 IfaceWrapper(const IfaceWrapper&) = delete;
65 IfaceWrapper& operator=(const IfaceWrapper&) = delete;
66 IfaceWrapper(IfaceWrapper&&) = delete;
67 IfaceWrapper& operator=(IfaceWrapper&&) = delete;
68
69 // void init();
70 void start();
71 void stop();
72
73 void generate_opmon_data() override;
74
75 void allocate_mbufs();
76 void setup_interface();
77 void setup_flow_steering();
78 void setup_xstats();
79 void stop_xstats();
80
81 void enable_flow() { m_lcore_enable_flow.store(true); }
82 void disable_flow() { m_lcore_enable_flow.store(false); }
83
84 const std::vector<uint16_t>& get_rte_cores() const { return m_rte_cores; }
85
86protected:
87 // iface_conf_t m_cfg;
88 int m_iface_id;
89 std::string m_iface_id_str;
90 bool m_configured;
91
92 bool m_with_flow;
93 bool m_prom_mode;
94 std::vector<std::string> m_ip_addr;
95 std::vector<rte_be32_t> m_ip_addr_bin;
96 std::string m_mac_addr;
97 int m_socket_id;
98 int m_mtu;
99 unsigned m_max_block_words;
100 uint16_t m_rx_ring_size;
101 uint16_t m_tx_ring_size;
102 int m_num_mbufs;
103 int m_burst_size;
104 uint32_t m_lcore_sleep_ns;
105 int m_mbuf_cache_size;
106
107private:
108 int m_num_ip_sources;
109 int m_num_rx_cores;
110 std::set<std::string> m_ips;
111 std::set<int> m_rx_qs;
112 std::set<int> m_tx_qs;
113 std::vector<uint16_t> m_rte_cores;
114
115 // CPU core ID -> [queue -> ip]
116 std::map<int, std::map<int, std::string>> m_rx_core_map;
117 unsigned m_arp_rx_queue = 0; // RS TODO: make it configurable, and queue use conf check for exclusiveness!
118
119 // Lcore stop signal
120 std::atomic<bool> m_lcore_quit_signal{ false };
121
122 std::atomic<bool> m_lcore_enable_flow{ false };
123
124 // Mbufs and pools
125 std::map<int, std::unique_ptr<rte_mempool>> m_mbuf_pools;
126 std::map<int, struct rte_mbuf**> m_bufs; // by queue
127
128 // Stats by queues
129 std::map<int, std::atomic<std::size_t>> m_num_frames_rxq;
130 std::map<int, std::atomic<std::size_t>> m_num_bytes_rxq;
131 std::map<int, std::atomic<std::size_t>> m_num_full_bursts;
132 std::map<int, std::atomic<uint16_t>> m_max_burst_size;
133
134 // Stats by rte_workers
135 std::map<int, std::atomic<std::size_t>> m_num_unhandled_non_ipv4;
136 std::map<int, std::atomic<std::size_t>> m_num_unhandled_non_udp;
137 std::map<int, std::atomic<std::size_t>> m_num_unhandled_non_jumbo_udp;
138
139 // Unexpected stream ID count
140 std::map<int, std::atomic<std::size_t>> m_num_unexid_frames;
141
142 // DPDK HW stats
143 dpdklibs::IfaceXstats m_iface_xstats;
144
145 // stream -> source id map indexed by queue id
146 // queue -> [stream_id -> sid]
147 std::map<int, std::map<uint, uint>> m_stream_id_to_source_id;
148 sid_to_source_map_t& m_sources;
149 bool m_strict_parsing{ true };
150
151 // Run marker
152 std::atomic<bool>& m_run_marker;
153
154 // GARP
155 std::unique_ptr<rte_mempool> m_garp_mbuf_pool;
156 std::map<int, struct rte_mbuf**> m_garp_bufs;
157 std::thread m_garp_thread;
158 void garp_func();
159 std::atomic<uint64_t> m_garps_sent{ 0 };
160
161 // ARP
162 std::unique_ptr<rte_mempool> m_arp_mbuf_pool;
163 std::map<int, struct rte_mbuf**> m_arp_bufs;
164 std::thread m_arp_thread;
165 void arp_func();
166 std::atomic<uint64_t> m_arps_sent{ 0 };
167
168 // Lcore processor
169 int rx_runner(void* arg __rte_unused);
170 int arp_response_runner(void* arg __rte_unused);
171
172 // Parse UDP payloads as DAQ frames
173 void parse_udp_payload(int src_rx_q, char* payload, std::size_t size);
174
175 // Pass through UDP payloads as is
176 void passthrough_udp_payload(int src_rx_q, char* payload, std::size_t size);
177};
178
179} // namespace dpdklibs
180} // namespace dunedaq
181
182// #include "detail/IfaceWrapper.hxx"
183
184#endif // DPDKLIBS_SRC_IFACEWRAPPER_HPP_
std::atomic< bool > run_marker
Global atomic for process lifetime.
The DUNE-DAQ namespace.
ERS_DECLARE_ISSUE(cibmodules, CIBCommunicationError, " CIB Hardware Communication Error: "<< descriptor,((std::string) descriptor)) ERS_DECLARE_ISSUE(cibmodules
FELIX Initialization std::string initerror FELIX queue timed std::string queuename Unexpected chunk size