DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
IfaceWrapper.cpp
Go to the documentation of this file.
1
9#include "logging/Logging.hpp"
10
11#include "opmonlib/Utils.hpp"
12
13#include "dpdklibs/Issues.hpp"
14
16
17#include "IfaceWrapper.hpp"
18#include "dpdklibs/EALSetup.hpp"
20#include "dpdklibs/arp/ARP.hpp"
24
26// #include "confmodel/DROStreamConf.hpp"
27// #include "confmodel/StreamParameters.hpp"
30#include "confmodel/GeoId.hpp"
32// #include "confmodel/NetworkDevice.hpp"
33// #include "appmodel/NICInterfaceConfiguration.hpp"
34// #include "appmodel/NICStatsConf.hpp"
35// #include "appmodel/EthStreamParameters.hpp"
36
38
39#include <chrono>
40#include <format>
41#include <memory>
42#include <regex>
43#include <stdexcept>
44#include <string>
45
49enum
50{
54};
55
56namespace dunedaq {
57namespace dpdklibs {
58
59//-----------------------------------------------------------------------------
60// TODO: the constructor signature shall be reviewd.
61// The current constructor takes a set of largely correlated arguments
62// - a receiver object
63// - the list of active senders
64// - the list of active streams
65//
66// These arguments are created by applying the enable mask to detector2daq connection object.
67// They are preferred to the d2d object not to expose the IfaceWrapper code to the System class
68IfaceWrapper::IfaceWrapper(uint iface_id,
69 const appmodel::DPDKReceiver* receiver,
70 const std::vector<const appmodel::NWDetDataSender*>& nw_senders,
71 const std::vector<const confmodel::DetectorStream*>& active_streams,
72 sid_to_source_map_t& sources,
73 std::atomic<bool>& run_marker)
74 : m_sources(sources)
75 , m_run_marker(run_marker)
76{
77
78 // Arguments consistency check: collect source ids in senders
79 std::set<int> src_in_d2d;
80
81 for (auto& det_stream : active_streams) {
82 src_in_d2d.insert(det_stream->get_source_id());
83 }
84
85 // Arguments consistency check: collect source ids in source model map
86 std::set<int> src_models;
87 for (const auto& [src_id, _] : m_sources) {
88 src_models.insert(src_id);
89 }
90
91 // check that the 2 sets are identical.
92 if (!std::includes(src_models.begin(), src_models.end(), src_in_d2d.begin(), src_in_d2d.end())) {
93
94 // TODO: remove, possibly
95 for (auto src : src_models)
96 TLOG_DEBUG(TLVL_BOOKKEEPING) << "model srcid " << src;
97 for (auto src : src_in_d2d)
98 TLOG_DEBUG(TLVL_BOOKKEEPING) << "d2d srcid " << src;
99
100 // D2D sources are not included in the source model list
101 // Extract the differences: src_in_d2d - src_models
102
103 std::vector<int> src_missing;
104 std::set_difference(
105 src_models.begin(), src_models.end(), src_in_d2d.begin(), src_in_d2d.end(), std::back_inserter(src_missing));
106
107 std::stringstream ss;
108 for (int src : src_missing) {
109 ss << src << " ";
110 }
111
112 // TLOG() << std::format("WARNING : these source ids are present in the d2d connection but no corresponding source
113 // objects are found {}", ss.str());
114 throw MissingSourceIDOutputs(ERS_HERE, m_iface_id, ss.str());
115 }
116
117 auto net_device = receiver->get_uses();
118
119 m_iface_id = iface_id;
120 m_mac_addr = net_device->get_mac_address();
121 m_ip_addr = net_device->get_ip_address();
122
123 TLOG() << "Building IfaceWrapper " << m_iface_id;
124 std::stringstream s;
125 s << 'IfaceWrapper (port ' << m_iface_id << ") responding to : ";
126 for (const std::string& ip_addr : m_ip_addr) {
127 s << ip_addr << " ";
128 }
129
130 TLOG() << s.str();
131
132 for (const std::string& ip_addr : m_ip_addr) {
133 IpAddr ip_addr_struct(ip_addr);
134 m_ip_addr_bin.push_back(udp::ip_address_dotdecimal_to_binary(ip_addr_struct.addr_bytes[0],
135 ip_addr_struct.addr_bytes[1],
136 ip_addr_struct.addr_bytes[2],
137 ip_addr_struct.addr_bytes[3]));
138 }
139
140 auto iface_cfg = receiver->get_configuration();
141
142 m_with_flow = iface_cfg->get_flow_control();
143 m_prom_mode = iface_cfg->get_promiscuous_mode();
144 ;
145 m_mtu = iface_cfg->get_mtu();
146 m_max_block_words = unsigned(m_mtu) / sizeof(uint64_t);
147 m_rx_ring_size = iface_cfg->get_rx_ring_size();
148 m_tx_ring_size = iface_cfg->get_tx_ring_size();
149 m_num_mbufs = iface_cfg->get_num_bufs();
150 m_burst_size = iface_cfg->get_burst_size();
151 m_mbuf_cache_size = iface_cfg->get_mbuf_cache_size();
152
153 m_lcore_sleep_ns = iface_cfg->get_lcore_sleep_us() * 1000;
154 m_socket_id = rte_eth_dev_socket_id(m_iface_id);
155
156 m_iface_id_str = iface_cfg->UID();
157
158 // Here is my list of cores
159 for (const auto* proc_res : iface_cfg->get_used_lcores()) {
160 m_rte_cores.insert(m_rte_cores.end(), proc_res->get_cpu_cores().begin(), proc_res->get_cpu_cores().end());
161 }
162 if (std::find(m_rte_cores.begin(), m_rte_cores.end(), rte_get_main_lcore()) != m_rte_cores.end()) {
163 TLOG() << "ERROR! Throw ERS error here that LCore=0 should not be used, as it's a control RTE core!";
164 throw std::runtime_error(
165 std::string("ERROR! Throw ERS here that LCore=0 should not be used, as it's a control RTE core!"));
166 }
167
168 // iterate through active streams
169
170 // Create a map of sender ni (ip) to streams from the d2d connection object
171 std::map<std::string, std::map<uint, uint>> ip_to_stream_src_groups;
172
173 for (auto nw_sender : nw_senders) {
174 auto sender_ni = nw_sender->get_uses();
175
176 std::string tx_ip = sender_ni->get_ip_address().at(0);
177
178 // Loop over streams
179 for (auto det_stream : nw_sender->get_streams()) {
180
181 // Only include active streams
182 if (std::find(active_streams.begin(), active_streams.end(), det_stream) == active_streams.end())
183 continue;
184
185 uint32_t tx_geo_stream_id = det_stream->get_geo_id()->get_stream_id();
186 // (tx, geo_stream) -> source_id
187 ip_to_stream_src_groups[tx_ip][tx_geo_stream_id] = det_stream->get_source_id();
188 }
189 }
190
191 // RS FIXME: Is this RX_Q bump is enough??? I don't remember how the RX_Qs are assigned...
192 uint32_t core_idx(0), rx_q(0); // RS FIXME: Ensure that no RX_Q=0 is used for UDP RX, ever.
193
194 m_rx_qs.insert(rx_q);
195 m_arp_rx_queue = rx_q;
196 ++rx_q;
197
198 // Build additional helper maps
199 for (const auto& [tx_ip, strm_src] : ip_to_stream_src_groups) {
200 m_ips.insert(tx_ip);
201 m_rx_qs.insert(rx_q);
202 m_num_frames_rxq[rx_q] = { 0 };
203 m_num_bytes_rxq[rx_q] = { 0 };
204
205 m_rx_core_map[m_rte_cores[core_idx]][rx_q] = tx_ip;
206 m_stream_id_to_source_id[rx_q] = strm_src;
207 // TLOG() << "+++ ip, rx_q : (" << tx_ip << ", " << rx_q << ") -> " << strm_src;
208
209 ++rx_q;
210 if (++core_idx == m_rte_cores.size()) {
211 core_idx = 0;
212 }
213 }
214
215 // Log mapping
216 for (auto const& [lcore, rx_qs] : m_rx_core_map) {
217 TLOG() << "Lcore=" << lcore << " handles: ";
218 for (auto const& [rx_q, src_ip] : rx_qs) {
219 TLOG() << " rx_q=" << rx_q << " src_ip=" << src_ip;
220 }
221 }
222
223 // Adding single TX queue for ARP responses
224 TLOG() << "Append TX_Q=0 for ARP responses.";
225 m_tx_qs.insert(0);
226
227 // Strict parsing (DAQ protocol) or pass through of UDP payloads to SourceModels
228 for (auto const& [sid, src_concept] : m_sources) {
229 if (!src_concept->m_daq_protocol_ensured) {
230 m_strict_parsing = false;
231 }
232 }
233}
234
235//-----------------------------------------------------------------------------
236IfaceWrapper::~IfaceWrapper()
237{
238 TLOG_DEBUG(TLVL_ENTER_EXIT_METHODS) << "IfaceWrapper destructor called. First stop check, then closing iface.";
239
240 struct rte_flow_error error;
241 rte_flow_flush(m_iface_id, &error);
242 // graceful_stop();
243 // close_iface();
244 TLOG_DEBUG(TLVL_ENTER_EXIT_METHODS) << "IfaceWrapper destroyed.";
245}
246
247//-----------------------------------------------------------------------------
248void
249IfaceWrapper::allocate_mbufs()
250{
251 TLOG() << "Allocating pools and mbufs for UDP, GARP, and ARP.";
252
253 // Pools for UDP RX messages
254 for (size_t i = 0; i < m_rx_qs.size(); ++i) {
255 std::stringstream bufss;
256 bufss << "MBP-" << m_iface_id << '-' << i;
257 TLOG() << "Acquire pool with name=" << bufss.str() << " for iface_id=" << m_iface_id << " rxq=" << i;
258 m_mbuf_pools[i] = ealutils::get_mempool(bufss.str(), m_num_mbufs, m_mbuf_cache_size, 16384, m_socket_id);
259 m_bufs[i] = (rte_mbuf**)malloc(sizeof(struct rte_mbuf*) * m_burst_size);
260 // No need to alloc?
261 // rte_pktmbuf_alloc_bulk(m_mbuf_pools[i].get(), m_bufs[i], m_burst_size);
262 }
263
264 // Pools for GARP messages
265 std::stringstream garpss;
266 garpss << "GARPMBP-" << m_iface_id;
267 TLOG() << "Acquire GARP pool with name=" << garpss.str() << " for iface_id=" << m_iface_id;
268 m_garp_mbuf_pool = ealutils::get_mempool(garpss.str());
269 m_garp_bufs[0] = (rte_mbuf**)malloc(sizeof(struct rte_mbuf*) * m_burst_size);
270 rte_pktmbuf_alloc_bulk(m_garp_mbuf_pool.get(), m_garp_bufs[0], m_burst_size);
271
272 // Pools for ARP request/responses
273 std::stringstream arpss;
274 arpss << "ARPMBP-" << m_iface_id;
275 TLOG() << "Acquire ARP pool with name=" << arpss.str() << " for iface_id=" << m_iface_id;
276 m_arp_mbuf_pool = ealutils::get_mempool(arpss.str());
277 m_arp_bufs[0] = (rte_mbuf**)malloc(sizeof(struct rte_mbuf*) * m_burst_size);
278 rte_pktmbuf_alloc_bulk(m_arp_mbuf_pool.get(), m_arp_bufs[0], m_burst_size);
279}
280
281//-----------------------------------------------------------------------------
282void
283IfaceWrapper::setup_interface()
284{
285 TLOG() << "Initialize interface " << m_iface_id;
286 bool with_reset = false, with_mq_mode = true; // go to config
287 bool check_link_status = false;
288
289 int retval = ealutils::iface_init(m_iface_id,
290 m_rx_qs.size(),
291 m_tx_qs.size(),
292 m_rx_ring_size,
293 m_tx_ring_size,
294 m_mbuf_pools,
295 with_reset,
296 with_mq_mode,
297 check_link_status);
298 if (retval != 0) {
299 throw FailedToSetupInterface(ERS_HERE, m_iface_id, retval);
300 }
301 // Promiscuous mode
302 ealutils::iface_promiscuous_mode(m_iface_id, m_prom_mode); // should come from config
303}
304
305//-----------------------------------------------------------------------------
306void
307IfaceWrapper::setup_flow_steering()
308{
309 // Flow steering setup
310 TLOG() << "Configuring Flow steering rules for iface=" << m_iface_id;
311 struct rte_flow_error error;
312 struct rte_flow* flow;
313 TLOG() << "Attempt to flush previous flow rules...";
314 rte_flow_flush(m_iface_id, &error);
315#warning RS: FIXME -> Check for flow flush return!
316
317 TLOG() << "Create control flow rules (ARP) assinged to rxq=" << m_arp_rx_queue;
318 flow = generate_arp_flow(m_iface_id, m_arp_rx_queue, &error);
319 if (not flow) { // ers::fatal
320 TLOG() << "ARP flow can't be created for " << m_arp_rx_queue << " Error type: " << (unsigned)error.type
321 << " Message: '" << error.message << "'";
322 ers::fatal(dunedaq::datahandlinglibs::InitializationError(ERS_HERE, "Couldn't create ARP flow API rules!"));
323 rte_exit(EXIT_FAILURE, "error in creating ARP flow");
324 }
325
326 TLOG() << "Create flow rules for UDP RX.";
327 for (auto const& [lcoreid, rxqs] : m_rx_core_map) {
328 for (auto const& [rxqid, srcip] : rxqs) {
329 // Put the IP numbers temporarily in a vector, so they can be converted easily to uint32_t
330 TLOG() << "Creating flow rule for src_ip=" << srcip << " assigned to rxq=" << rxqid;
331 size_t ind = 0, current_ind = 0;
332 std::vector<uint8_t> v;
333 for (int i = 0; i < 4; ++i) {
334 v.push_back(std::stoi(srcip.substr(current_ind, srcip.size() - current_ind), &ind));
335 current_ind += ind + 1;
336 }
337
338 flow = generate_ipv4_flow(m_iface_id, rxqid, RTE_IPV4(v[0], v[1], v[2], v[3]), 0xffffffff, 0, 0, &error);
339
340 if (not flow) { // ers::fatal
341 TLOG() << "Flow can't be created for " << rxqid << " Error type: " << (unsigned)error.type << " Message: '"
342 << error.message << "'";
343 ers::fatal(dunedaq::datahandlinglibs::InitializationError(ERS_HERE, "Couldn't create Flow API rules!"));
344 rte_exit(EXIT_FAILURE, "error in creating flow");
345 }
346 }
347 }
348
349 return;
350}
351
352//-----------------------------------------------------------------------------
353void
354IfaceWrapper::setup_xstats()
355{
356 // Stats setup
357 m_iface_xstats.setup(m_iface_id);
358 m_iface_xstats.reset_counters();
359}
360
361//-----------------------------------------------------------------------------
362void
363IfaceWrapper::stop_xstats()
364{
365 // Stopping stats
366 m_iface_xstats.reset_counters();
367 m_iface_xstats.stop();
368}
369
370//-----------------------------------------------------------------------------
371void
372IfaceWrapper::start()
373{
374 // Reset counters for RX queues
375 for (auto const& [rx_q, _] : m_num_frames_rxq) {
376 m_num_frames_rxq[rx_q] = { 0 };
377 m_num_bytes_rxq[rx_q] = { 0 };
378 m_num_full_bursts[rx_q] = { 0 };
379 m_max_burst_size[rx_q] = { 0 };
380 }
381
382 // Reset counters for rte_workers
383 for (auto const& [lcore, _] : m_rx_core_map) {
384 m_num_unhandled_non_ipv4[lcore] = { 0 };
385 m_num_unhandled_non_udp[lcore] = { 0 };
386 m_num_unhandled_non_jumbo_udp[lcore] = { 0 };
387 }
388
389 m_lcore_enable_flow.store(false);
390 m_lcore_quit_signal.store(false);
391 TLOG() << "Interface id=" << m_iface_id << " Launching GARP thread with garp_func...";
392 m_garp_thread = std::thread(&IfaceWrapper::garp_func, this);
393
394 TLOG() << "Interface id=" << m_iface_id << " starting ARP LCore processor:";
395 m_arp_thread = std::thread(&IfaceWrapper::IfaceWrapper::arp_response_runner, this, nullptr);
396
397 TLOG() << "Interface id=" << m_iface_id << " starting LCore processors:";
398 for (auto const& [lcoreid, _] : m_rx_core_map) {
399 int ret = rte_eal_remote_launch((int (*)(void*))(&IfaceWrapper::rx_runner), this, lcoreid);
400 TLOG() << " -> LCore[" << lcoreid << "] launched with return code=" << ret << " "
401 << (ret < 0 ? rte_strerror(-ret) : "");
402 }
403}
404
405//-----------------------------------------------------------------------------
406void
407IfaceWrapper::stop()
408{
409 m_lcore_enable_flow.store(false);
410 m_lcore_quit_signal.store(true);
411 // Stop GARP sender thread
412 if (m_garp_thread.joinable()) {
413 m_garp_thread.join();
414 } else {
415 TLOG() << "GARP thread is not joinable!";
416 }
417
418 if (m_arp_thread.joinable()) {
419 m_arp_thread.join();
420 } else {
421 TLOG() << "ARP thread is not joinable!";
422 }
423}
424/*
425void
426IfaceWrapper::scrap()
427{
428 struct rte_flow_error error;
429 rte_flow_flush(m_iface_id, &error);
430}
431*/
432
433//-----------------------------------------------------------------------------
434void
435IfaceWrapper::generate_opmon_data()
436{
437
438 if (m_iface_xstats.m_enabled) {
439 // Poll stats from HW
440 m_iface_xstats.poll();
441
443 s.set_ipackets(m_iface_xstats.m_eth_stats.ipackets);
444 s.set_opackets(m_iface_xstats.m_eth_stats.opackets);
445 s.set_ibytes(m_iface_xstats.m_eth_stats.ibytes);
446 s.set_obytes(m_iface_xstats.m_eth_stats.obytes);
447 s.set_imissed(m_iface_xstats.m_eth_stats.imissed);
448 s.set_ierrors(m_iface_xstats.m_eth_stats.ierrors);
449 s.set_oerrors(m_iface_xstats.m_eth_stats.oerrors);
450 s.set_rx_nombuf(m_iface_xstats.m_eth_stats.rx_nombuf);
451 publish(std::move(s));
452
453 if (m_iface_xstats.m_eth_stats.imissed > 0) {
454 ers::warning(PacketErrors(ERS_HERE, m_iface_id_str, "missed", m_iface_xstats.m_eth_stats.imissed));
455 }
456 if (m_iface_xstats.m_eth_stats.ierrors > 0) {
457 ers::warning(PacketErrors(ERS_HERE, m_iface_id_str, "dropped", m_iface_xstats.m_eth_stats.ierrors));
458 }
459
460 // loop over all the xstats information
463 std::map<std::string, opmon::QueueEthXStats> xq;
464
465 for (int i = 0; i < m_iface_xstats.m_len; ++i) {
466
467 std::string name(m_iface_xstats.m_xstats_names[i].name);
468
469 // first we select the info from the queue
470 static std::regex queue_regex(R"((rx|tx)_q(\d+)_([^_]+))");
471 std::smatch match;
472
473 if (std::regex_match(name, match, queue_regex)) {
474 auto queue_name = match[1].str() + '-' + match[2].str();
475 auto& entry = xq[queue_name];
476 try {
477 opmonlib::set_value(entry, match[3], m_iface_xstats.m_xstats_values[i]);
478 } catch (const ers::Issue& e) {
479 ers::warning(MetricPublishFailed(ERS_HERE, name, e));
480 }
481 continue;
482 }
483
484 google::protobuf::Message* metric_p = nullptr;
485 static std::regex err_regex(R"(.+error.*)");
486 if (std::regex_match(name, err_regex))
487 metric_p = &xerrs;
488 else
489 metric_p = &xinfos;
490
491 try {
492 opmonlib::set_value(*metric_p, name, m_iface_xstats.m_xstats_values[i]);
493 } catch (const ers::Issue& e) {
494 ers::warning(MetricPublishFailed(ERS_HERE, name, e));
495 }
496
497 } // loop over xstats
498
499 // Reset HW counters
500 m_iface_xstats.reset_counters();
501 // finally we publish the information
502 publish(std::move(xinfos));
503 publish(std::move(xerrs));
504 for (auto [id, stat] : xq) {
505 publish(std::move(stat), { { "queue", id } });
506 }
507 }
508
509 for (const auto& [src_rx_q, _] : m_num_frames_rxq) {
511 i.set_packets_received(m_num_frames_rxq[src_rx_q].load());
512 i.set_bytes_received(m_num_bytes_rxq[src_rx_q].load());
513 i.set_full_rx_burst(m_num_full_bursts[src_rx_q].load());
514 i.set_max_burst_size(m_max_burst_size[src_rx_q].exchange(0));
515
516 publish(std::move(i), { { "queue", std::to_string(src_rx_q) } });
517 }
518
519 // RTE Workers
520 for (auto const& [lcore, _] : m_rx_core_map) {
522 info.set_num_unhandled_non_ipv4(m_num_unhandled_non_ipv4[lcore].exchange(0));
523 info.set_num_unhandled_non_udp(m_num_unhandled_non_udp[lcore].exchange(0));
524 info.set_num_unhandled_non_jumbo_udp(m_num_unhandled_non_jumbo_udp[lcore].exchange(0));
525 publish(std::move(info), { { "rte_worker_id", std::to_string(lcore) } });
526 }
527
528 for (auto& [id, counter] : m_num_unexid_frames) {
529 auto val = counter.exchange(0);
530 if (val > 0) {
532 }
533 }
534}
535
536//-----------------------------------------------------------------------------
537void
538IfaceWrapper::garp_func()
539{
540 TLOG() << "Launching GARP sender...";
541 while (m_run_marker.load()) {
542 for (const auto& ip_addr_bin : m_ip_addr_bin) {
543 arp::pktgen_send_garp(m_garp_bufs[0][0], m_iface_id, ip_addr_bin);
544 }
545 ++m_garps_sent;
546 std::this_thread::sleep_for(std::chrono::seconds(1));
547 }
548 TLOG() << "GARP function joins.";
549}
550
551//-----------------------------------------------------------------------------
552void
553IfaceWrapper::parse_udp_payload(int src_rx_q, char* payload, std::size_t size)
554{
555 // Pointers for parsing and to the end of the UDP payload.
556 char* plptr = payload;
557 const char* plendptr = payload + size;
558
559 // Process every DAQEth frame within UDP payload
560 while (plptr + sizeof(dunedaq::detdataformats::DAQEthHeader) < plendptr) { // Scatter loop start
561
562 // Reinterpret directly to DAQEthHeader
563 auto daqhdrptr = reinterpret_cast<dunedaq::detdataformats::DAQEthHeader*>(plptr);
564
565 // Check number of DAQEth block_words
566 unsigned block_words = unsigned(daqhdrptr->block_length) - 1; // removing timestamp word from the block length.
567
568 // Check for corrupted DAQEth frame length
569 if (block_words == 0 || block_words > m_max_block_words) {
570 // RS FIXME: corrupted length -> stop, add opmon counter or warning
571 return;
572 }
573
574 // Calculate data bytes after DAQEthHeader based on block_words
575 std::size_t data_bytes = std::size_t(block_words) * sizeof(dunedaq::detdataformats::DAQEthHeader::word_t);
576
577 // Grab end pointer of DAQEth frame
578 char* daqframe_endptr = plptr + sizeof(dunedaq::detdataformats::DAQEthHeader) + data_bytes;
579
580 // Check if full DAQEth frame fits
581 if (daqframe_endptr > plendptr) {
582 // RS FIXME: truncated payload -> stop, add opmon counter or warning
583 return;
584 }
585
586 // Calculate DAQEth frame size (used both for handling and advancing)
587 std::size_t daq_frame_size = sizeof(dunedaq::detdataformats::DAQEthHeader) + data_bytes;
588
589 // Sadly, cannot take a reference to a bitfield
590 uint strm_id = daqhdrptr->stream_id;
591
592 // Check that stream id is corresponds to a registered source
593 auto& strm_to_src = m_stream_id_to_source_id[src_rx_q];
594
595 if (auto strm_it = strm_to_src.find(strm_id); strm_it != strm_to_src.end()) {
596
597 m_sources[strm_it->second]->handle_daq_frame((char*)daqhdrptr, daq_frame_size);
598
599 } else {
600 // Really bad -> unexpeced StreamID in UDP Payload.
601 // This check is needed in order to avoid dynamically add thousands
602 // of Sources on the fly, in case the data corruption is extremely severe.
603 if (m_num_unexid_frames.count(strm_id) == 0) {
604 m_num_unexid_frames[strm_id] = 0;
605 }
606 m_num_unexid_frames[strm_id]++;
607 }
608
609 // Advance to next payload
610 plptr += daq_frame_size;
611
612 } // Scatter loop end
613}
614
615//-----------------------------------------------------------------------------
616void
617IfaceWrapper::passthrough_udp_payload(int src_rx_q, char* payload, std::size_t size)
618{
619 // Get DAQ Header and its StreamID
620 auto* daqhdrptr = reinterpret_cast<dunedaq::detdataformats::DAQEthHeader*>(payload);
621
622 // Sadly, cannot take a reference to a bitfield
623 uint strm_id = daqhdrptr->stream_id;
624
625 // Check that stream id is corresponds to a registered source
626 auto& strm_to_src = m_stream_id_to_source_id[src_rx_q];
627
628 if (auto strm_it = strm_to_src.find(strm_id); strm_it != strm_to_src.end()) {
629 m_sources[strm_it->second]->handle_daq_frame(payload, size);
630 } else {
631 // Really bad -> unexpeced StreamID in UDP Payload.
632 // This check is needed in order to avoid dynamically add thousands
633 // of Sources on the fly, in case the data corruption is extremely severe.
634 if (m_num_unexid_frames.count(strm_id) == 0) {
635 m_num_unexid_frames[strm_id] = 0;
636 }
637 m_num_unexid_frames[strm_id]++;
638 }
639}
640
641} // namespace dpdklibs
642} // namespace dunedaq
643
644//
#define ERS_HERE
void set_packets_received(::uint64_t value)
std::atomic< bool > run_marker
Global atomic for process lifetime.
#define TLVL_ENTER_EXIT_METHODS
#define TLOG_DEBUG(lvl,...)
Definition Logging.hpp:116
#define TLOG(...)
Definition macro.hpp:21
dunedaq::conffwk::relationship_t match(T const &, T const &)
void pktgen_send_garp(struct rte_mbuf *m, uint32_t port_id, rte_be32_t binary_ip_address)
Definition ARP.cpp:23
std::unique_ptr< rte_mempool > get_mempool(const std::string &pool_name, int num_mbufs=NUM_MBUFS, int mbuf_cache_size=MBUF_CACHE_SIZE, int data_room_size=9800, int socket_id=0)
Definition EALSetup.cpp:263
int iface_promiscuous_mode(std::uint16_t iface, bool mode=false)
Definition EALSetup.cpp:76
int iface_init(uint16_t iface, uint16_t rx_rings, uint16_t tx_rings, uint16_t rx_ring_size, uint16_t tx_ring_size, std::map< int, std::unique_ptr< rte_mempool > > &mbuf_pool, bool with_reset=false, bool with_mq_rss=false, bool check_link_status=false)
Definition EALSetup.cpp:95
rte_be32_t ip_address_dotdecimal_to_binary(std::uint8_t byte1, std::uint8_t byte2, std::uint8_t byte3, std::uint8_t byte4)
Definition Utils.cpp:41
struct rte_flow * generate_ipv4_flow(uint16_t port_id, uint16_t rx_q, uint32_t src_ip, uint32_t src_mask, uint32_t dest_ip, uint32_t dest_mask, struct rte_flow_error *error)
struct rte_flow * generate_arp_flow(uint16_t port_id, uint16_t rx_q, struct rte_flow_error *error)
void set_value(google::protobuf::Message &m, const std::string &name, T value)
Definition Utils.hxx:19
The DUNE-DAQ namespace.
FELIX Initialization std::string initerror FELIX queue timed std::string queuename Unexpected chunk size
default char v[0]
CIB Buffer std::string descriptor Message from std::string descriptor CIB process error
void warning(const Issue &issue)
Definition ers.hpp:150
void fatal(const Issue &issue)
Definition ers.hpp:111