DUNE-DAQ
DUNE Trigger and Data Acquisition software
Loading...
Searching...
No Matches
IfaceWrapper.hxx
Go to the documentation of this file.
2#include <rte_arp.h>
3#include <rte_ethdev.h>
4#include <time.h>
5
6namespace dunedaq {
7namespace dpdklibs {
8
9int
10IfaceWrapper::rx_runner(void* arg __rte_unused)
11{
12
13 // Timespec for opportunistic sleep. Nanoseconds configured in conf.
14 struct timespec sleep_request = { 0, (long)m_lcore_sleep_ns };
15
16 // bool once = true; // One shot action variable.
17 uint16_t iface = m_iface_id;
18
19 const uint16_t lid = rte_lcore_id();
20 auto queues = m_rx_core_map[lid];
21
22 if (rte_eth_dev_socket_id(iface) >= 0 && rte_eth_dev_socket_id(iface) != (int)rte_socket_id()) {
23 TLOG() << "WARNING, iface " << iface << " is on remote NUMA node to polling thread! "
24 << "Performance will not be optimal.";
25 }
26
27 TLOG() << "LCore RX runner on CPU[" << lid << "]: Main loop starts for iface " << iface << " !";
28
29 std::map<int, int> nb_rx_map;
30 // While loop of quit atomic member in IfaceWrapper
31 while (!this->m_lcore_quit_signal.load()) {
32
33 // Loop over assigned queues to process
34 uint8_t fb_count(0);
35 for (const auto& q : queues) {
36 auto src_rx_q = q.first;
37 auto* q_bufs = m_bufs[src_rx_q];
38
39 // Get burst from queue
40 const uint16_t nb_rx = rte_eth_rx_burst(iface, src_rx_q, q_bufs, m_burst_size);
41 nb_rx_map[src_rx_q] = nb_rx;
42 }
43
44 for (const auto& q : queues) {
45
46 auto src_rx_q = q.first;
47 auto* q_bufs = m_bufs[src_rx_q];
48 const uint16_t nb_rx = nb_rx_map[src_rx_q];
49
50 // We got packets from burst on this queue
51 if (nb_rx != 0) [[likely]] {
52
53 // Update max burst size counter of this queue
54 m_max_burst_size[src_rx_q] = std::max(nb_rx, m_max_burst_size[src_rx_q].load());
55
56 // -------
57 // Iterate on burst packets
58 for (int i_b = 0; i_b < nb_rx; ++i_b) {
59
60 // Check if packet is segmented. Implement support for it if needed.
61 // if (q_bufs[i_b]->nb_segs > 1) [[unlikely]] {
62 // TLOG_DEBUG(10) << "It appears a packet is spread across more than one receiving buffer;"
63 // << " there's currently no logic in this program to handle this";
64 //}
65
66 // Check packet type, decide their fate: ignore unexpected ones, FIXME: monitor occurrences
67 auto pkt_type = q_bufs[i_b]->packet_type;
68 // Handle non IPV4 frames.
69 if (not RTE_ETH_IS_IPV4_HDR(pkt_type)) [[unlikely]] {
70 // TLOG_DEBUG(10) << "Non-Ethernet packet type: " << (unsigned)pkt_type << " original: " << pkt_type;
71 if (pkt_type == RTE_PTYPE_L2_ETHER_ARP) {
72 // TLOG() << "Unexpected: Should handle an ARP request from lcore=" << lid << " rx_q=" << src_rx_q << "!
73 // Flow should be steered to dedicated RX Queue.";
74 } else if (pkt_type == RTE_PTYPE_L2_ETHER_LLDP) {
75 // TLOG_DEBUG(10) << "TODO: Handle LLDP packet!";
76 } else {
77 // TLOG_DEBUG(10) << "Unidentified! Dumping...";
78 // rte_pktmbuf_dump(stdout, q_bufs[i_b], m_bufs[src_rx_q][i_b]->pkt_len);
79 }
80 ++m_num_unhandled_non_ipv4[lid];
81 continue;
82 }
83
84 // Check if frame is non UDP: in that case, ignore it.
85 if ((pkt_type & RTE_PTYPE_L4_MASK) != RTE_PTYPE_L4_UDP) [[unlikely]] {
86 ++m_num_unhandled_non_udp[lid];
87 continue; // ommit it
88 }
89
90 // Check for JUMBO frames (bigger than 1500 Bytes)
91 if (q_bufs[i_b]->pkt_len > 1500) [[likely]] { // RS FIXME: do proper check on data length later
92
93 // If flow enabled, handle the payload.
94 if (m_lcore_enable_flow.load()) [[likely]] {
95 // Get length of user payload. (Ethernet headers excluded.)
96 struct udp::ipv4_udp_packet_hdr* udp_packet =
97 rte_pktmbuf_mtod(q_bufs[i_b], struct udp::ipv4_udp_packet_hdr*);
98 char* message = udp::get_udp_payload(q_bufs[i_b]);
99 std::size_t udp_payload_len = udp::get_payload_size_udp_hdr(&udp_packet->udp_hdr);
100
101 if (m_strict_parsing) { // all sources maintain DAQ protocol
102 parse_udp_payload(src_rx_q, message, udp_payload_len);
103 } else { // avoid size checks and scattering
104 passthrough_udp_payload(src_rx_q, message, udp_payload_len);
105 }
106 }
107
108 // Update metrics of queue: frame and Byte counters
109 ++m_num_frames_rxq[src_rx_q];
110 std::size_t data_len = q_bufs[i_b]->data_len;
111 m_num_bytes_rxq[src_rx_q] += data_len;
112 } else {
113 ++m_num_unhandled_non_jumbo_udp[lid];
114 }
115 }
116
117 // Bulk free of mbufs
118 rte_pktmbuf_free_bulk(q_bufs, nb_rx);
119
120 } // per burst
121 // -----------
122
123 // Full burst counter
124 if (nb_rx == m_burst_size) {
125 ++fb_count;
126 ++m_num_full_bursts[src_rx_q];
127 }
128 } // per queue
129
130 // If no full buffers in burst...
131 if (!fb_count) {
132 if (m_lcore_sleep_ns) {
133 // Sleep n nanoseconds... (value from config, timespec initialized in lcore first lines)
134 /*int response =*/nanosleep(&sleep_request, nullptr);
135 }
136 }
137
138 } // main while(quit) loop
139
140 TLOG() << "LCore RX runner on CPU[" << lid << "] returned.";
141 return 0;
142}
143
144int
145IfaceWrapper::arp_response_runner(void* arg __rte_unused)
146{
147
148 // Timespec for opportunistic sleep. Nanoseconds configured in conf.
149 struct timespec sleep_request = { 0, (long)900000 };
150
151 // bool once = true; // One shot action variable.
152 uint16_t iface = m_iface_id;
153
154 const uint16_t lid = rte_lcore_id();
155 unsigned arp_rx_queue = m_arp_rx_queue;
156
157 TLOG() << "LCore ARP responder on CPU[" << lid << "]: Main loop starts for iface " << iface
158 << " rx queue: " << arp_rx_queue;
159
160 // While loop of quit atomic member in IfaceWrapper
161 while (!this->m_lcore_quit_signal.load()) {
162
163 const uint16_t nb_rx = rte_eth_rx_burst(iface, arp_rx_queue, m_arp_bufs[arp_rx_queue], m_burst_size);
164
165 // We got packets from burst on this queue
166 if (nb_rx != 0) {
167 // Iterate on burst packets
168 for (int i_b = 0; i_b < nb_rx; ++i_b) {
169
170 // Check packet type, ommit/drop unexpected ones.
171 auto pkt_type = m_arp_bufs[arp_rx_queue][i_b]->packet_type;
173 if (not RTE_ETH_IS_IPV4_HDR(pkt_type)) {
174 // TLOG_DEBUG(10) << "Non-Ethernet packet type: " << (unsigned)pkt_type << " original: " << pkt_type;
175 if (pkt_type == RTE_PTYPE_L2_ETHER_ARP) {
176 TLOG_DEBUG(10) << "Handling ARP request";
177 struct rte_ether_hdr* eth_hdr = rte_pktmbuf_mtod(m_arp_bufs[arp_rx_queue][i_b], struct rte_ether_hdr*);
178 struct rte_arp_hdr* arp_hdr = (struct rte_arp_hdr*)(eth_hdr + 1);
179
181 dunedaq::dpdklibs::udp::ip_address_binary_to_dotdecimal(rte_be_to_cpu_32(arp_hdr->arp_data.arp_sip)));
182 TLOG_DEBUG(10) << "SRC IP: " << srcaddr;
184 dunedaq::dpdklibs::udp::ip_address_binary_to_dotdecimal(rte_be_to_cpu_32(arp_hdr->arp_data.arp_tip)));
185 TLOG_DEBUG(10) << "DEST IP: " << dstaddr;
186
187 for (const auto& ip_addr_bin : m_ip_addr_bin) {
189 dunedaq::dpdklibs::udp::ip_address_binary_to_dotdecimal(rte_be_to_cpu_32(ip_addr_bin)));
190 TLOG_DEBUG(10) << "LOCAL IP: " << localaddr;
191 }
192
193 if (std::find(m_ip_addr_bin.begin(), m_ip_addr_bin.end(), arp_hdr->arp_data.arp_tip) !=
194 m_ip_addr_bin.end()) {
195 arp::pktgen_process_arp(m_arp_bufs[arp_rx_queue][i_b], m_iface_id, arp_hdr->arp_data.arp_tip);
196 } else {
197 TLOG_DEBUG(10) << "I'm not the ARP target";
198 }
199 } else if (pkt_type == RTE_PTYPE_L2_ETHER_LLDP) {
200 // TLOG_DEBUG(10) << "TODO: Handle LLDP packet!";
201 } else {
202 // TLOG_DEBUG(10) << "Unidentified! Dumping...";
203 // rte_pktmbuf_dump(stdout, m_arp_bufs[arp_rx_queue][i_b], m_bufs[src_rx_q][i_b]->pkt_len);
204 }
205 continue;
206 }
207 }
208
209 // Bulk free of mbufs
210 rte_pktmbuf_free_bulk(m_arp_bufs[arp_rx_queue], nb_rx);
211
212 } // per burst
213
214 // If no full buffers in burst...
215 if (m_lcore_sleep_ns) {
216 // Sleep n nanoseconds... (value from config, timespec initialized in lcore first lines)
217 /*int response =*/nanosleep(&sleep_request, nullptr);
218 }
219
220 } // main while(quit) loop
221
222 TLOG() << "LCore ARP responder on CPU[" << lid << "] returned.";
223 return 0;
224}
225
226} // namespace dpdklibs
227} // namespace dunedaq
#define TLOG_DEBUG(lvl,...)
Definition Logging.hpp:116
#define TLOG(...)
Definition macro.hpp:21
void pktgen_process_arp(struct rte_mbuf *m, uint32_t pid, rte_be32_t binary_ip_address)
Definition ARP.cpp:78
char * get_udp_payload(const rte_mbuf *mbuf)
Definition Utils.cpp:71
std::string get_ipv4_decimal_addr_str(struct ipaddr ipv4_address)
Definition Utils.cpp:56
struct ipaddr ip_address_binary_to_dotdecimal(rte_le32_t binary_ipv4_address)
Definition Utils.cpp:48
std::uint16_t get_payload_size_udp_hdr(struct rte_udp_hdr *udp_hdr)
Definition Utils.cpp:29
The DUNE-DAQ namespace.
message(message)
Definition __init__.py:84