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,
79 std::set<int> src_in_d2d;
81 for (
auto& det_stream : active_streams) {
82 src_in_d2d.insert(det_stream->get_source_id());
86 std::set<int> src_models;
87 for (
const auto& [
src_id, _] : m_sources) {
92 if (!std::includes(src_models.begin(), src_models.end(), src_in_d2d.begin(), src_in_d2d.end())) {
95 for (
auto src : src_models)
97 for (
auto src : src_in_d2d)
103 std::vector<int> src_missing;
105 src_models.begin(), src_models.end(), src_in_d2d.begin(), src_in_d2d.end(), std::back_inserter(src_missing));
107 std::stringstream ss;
108 for (
int src : src_missing) {
114 throw MissingSourceIDOutputs(
ERS_HERE, m_iface_id, ss.str());
117 auto net_device = receiver->get_uses();
119 m_iface_id = iface_id;
120 m_mac_addr = net_device->get_mac_address();
121 m_ip_addr = net_device->get_ip_address();
123 TLOG() <<
"Building IfaceWrapper " << m_iface_id;
125 s <<
'IfaceWrapper (port ' << m_iface_id <<
") responding to : ";
126 for (
const std::string& ip_addr : m_ip_addr) {
132 for (
const std::string& ip_addr : m_ip_addr) {
133 IpAddr ip_addr_struct(ip_addr);
135 ip_addr_struct.addr_bytes[1],
136 ip_addr_struct.addr_bytes[2],
137 ip_addr_struct.addr_bytes[3]));
140 auto iface_cfg = receiver->get_configuration();
142 m_with_flow = iface_cfg->get_flow_control();
143 m_prom_mode = iface_cfg->get_promiscuous_mode();
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();
153 m_lcore_sleep_ns = iface_cfg->get_lcore_sleep_us() * 1000;
154 m_socket_id = rte_eth_dev_socket_id(m_iface_id);
156 m_iface_id_str = iface_cfg->UID();
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());
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!"));
171 std::map<std::string, std::map<uint, uint>> ip_to_stream_src_groups;
173 for (
auto nw_sender : nw_senders) {
174 auto sender_ni = nw_sender->get_uses();
176 std::string
tx_ip = sender_ni->get_ip_address().at(0);
179 for (
auto det_stream : nw_sender->get_streams()) {
182 if (std::find(active_streams.begin(), active_streams.end(), det_stream) == active_streams.end())
185 uint32_t tx_geo_stream_id = det_stream->get_geo_id()->get_stream_id();
187 ip_to_stream_src_groups[
tx_ip][tx_geo_stream_id] = det_stream->get_source_id();
192 uint32_t core_idx(0), rx_q(0);
194 m_rx_qs.insert(rx_q);
195 m_arp_rx_queue = rx_q;
199 for (
const auto& [
tx_ip, strm_src] : ip_to_stream_src_groups) {
201 m_rx_qs.insert(rx_q);
202 m_num_frames_rxq[rx_q] = { 0 };
203 m_num_bytes_rxq[rx_q] = { 0 };
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;
210 if (++core_idx == m_rte_cores.size()) {
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;
224 TLOG() <<
"Append TX_Q=0 for ARP responses.";
228 for (
auto const& [sid, src_concept] : m_sources) {
229 if (!src_concept->m_daq_protocol_ensured) {
230 m_strict_parsing =
false;
236IfaceWrapper::~IfaceWrapper()
240 struct rte_flow_error
error;
241 rte_flow_flush(m_iface_id, &
error);
249IfaceWrapper::allocate_mbufs()
251 TLOG() <<
"Allocating pools and mbufs for UDP, GARP, and ARP.";
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);
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;
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);
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;
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);
283IfaceWrapper::setup_interface()
285 TLOG() <<
"Initialize interface " << m_iface_id;
286 bool with_reset =
false, with_mq_mode =
true;
287 bool check_link_status =
false;
299 throw FailedToSetupInterface(
ERS_HERE, m_iface_id, retval);
307IfaceWrapper::setup_flow_steering()
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!
317 TLOG() <<
"Create control flow rules (ARP) assinged to rxq=" << m_arp_rx_queue;
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");
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) {
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;
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");
354IfaceWrapper::setup_xstats()
357 m_iface_xstats.setup(m_iface_id);
358 m_iface_xstats.reset_counters();
363IfaceWrapper::stop_xstats()
366 m_iface_xstats.reset_counters();
367 m_iface_xstats.stop();
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 };
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 };
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);
394 TLOG() <<
"Interface id=" << m_iface_id <<
" starting ARP LCore processor:";
395 m_arp_thread = std::thread(&IfaceWrapper::IfaceWrapper::arp_response_runner,
this,
nullptr);
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) :
"");
409 m_lcore_enable_flow.store(
false);
410 m_lcore_quit_signal.store(
true);
412 if (m_garp_thread.joinable()) {
413 m_garp_thread.join();
415 TLOG() <<
"GARP thread is not joinable!";
418 if (m_arp_thread.joinable()) {
421 TLOG() <<
"ARP thread is not joinable!";
435IfaceWrapper::generate_opmon_data()
438 if (m_iface_xstats.m_enabled) {
440 m_iface_xstats.poll();
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));
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));
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));
463 std::map<std::string, opmon::QueueEthXStats> xq;
465 for (
int i = 0; i < m_iface_xstats.m_len; ++i) {
467 std::string name(m_iface_xstats.m_xstats_names[i].name);
470 static std::regex queue_regex(R
"((rx|tx)_q(\d+)_([^_]+))");
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];
478 }
catch (
const ers::Issue& e) {
484 google::protobuf::Message* metric_p =
nullptr;
485 static std::regex err_regex(R
"(.+error.*)");
486 if (std::regex_match(name, err_regex))
493 }
catch (
const ers::Issue& e) {
500 m_iface_xstats.reset_counters();
502 publish(std::move(xinfos));
503 publish(std::move(xerrs));
504 for (
auto [
id, stat] : xq) {
505 publish(std::move(stat), { {
"queue",
id } });
509 for (
const auto& [src_rx_q, _] : m_num_frames_rxq) {
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));
516 publish(std::move(i), { {
"queue", std::to_string(src_rx_q) } });
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) } });
528 for (
auto& [
id, counter] : m_num_unexid_frames) {
529 auto val = counter.exchange(0);
538IfaceWrapper::garp_func()
540 TLOG() <<
"Launching GARP sender...";
541 while (m_run_marker.load()) {
542 for (
const auto& ip_addr_bin : m_ip_addr_bin) {
546 std::this_thread::sleep_for(std::chrono::seconds(1));
548 TLOG() <<
"GARP function joins.";
553IfaceWrapper::parse_udp_payload(
int src_rx_q,
char*
payload, std::size_t
size)
560 while (plptr +
sizeof(dunedaq::detdataformats::DAQEthHeader) < plendptr) {
563 auto daqhdrptr =
reinterpret_cast<dunedaq::detdataformats::DAQEthHeader*
>(plptr);
566 unsigned block_words = unsigned(daqhdrptr->block_length) - 1;
569 if (block_words == 0 || block_words > m_max_block_words) {
578 char* daqframe_endptr = plptr +
sizeof(dunedaq::detdataformats::DAQEthHeader) + data_bytes;
581 if (daqframe_endptr > plendptr) {
587 std::size_t daq_frame_size =
sizeof(dunedaq::detdataformats::DAQEthHeader) + data_bytes;
590 uint
strm_id = daqhdrptr->stream_id;
593 auto& strm_to_src = m_stream_id_to_source_id[src_rx_q];
595 if (
auto strm_it = strm_to_src.find(
strm_id); strm_it != strm_to_src.end()) {
597 m_sources[strm_it->second]->handle_daq_frame((
char*)daqhdrptr, daq_frame_size);
603 if (m_num_unexid_frames.count(
strm_id) == 0) {
604 m_num_unexid_frames[
strm_id] = 0;
606 m_num_unexid_frames[
strm_id]++;
610 plptr += daq_frame_size;
617IfaceWrapper::passthrough_udp_payload(
int src_rx_q,
char*
payload, std::size_t
size)
620 auto* daqhdrptr =
reinterpret_cast<dunedaq::detdataformats::DAQEthHeader*
>(
payload);
623 uint
strm_id = daqhdrptr->stream_id;
626 auto& strm_to_src = m_stream_id_to_source_id[src_rx_q];
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);
634 if (m_num_unexid_frames.count(
strm_id) == 0) {
635 m_num_unexid_frames[
strm_id] = 0;
637 m_num_unexid_frames[
strm_id]++;
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,...)
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)
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)
int iface_promiscuous_mode(std::uint16_t iface, bool mode=false)
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)
rte_be32_t ip_address_dotdecimal_to_binary(std::uint8_t byte1, std::uint8_t byte2, std::uint8_t byte3, std::uint8_t byte4)
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)
FELIX Initialization std::string initerror FELIX queue timed std::string queuename Unexpected chunk size
CIB Buffer std::string descriptor Message from std::string descriptor CIB process error
void warning(const Issue &issue)
void fatal(const Issue &issue)