otsdaq-mu2e-stm  5.02.01
udp.cc
1 // UDP header
2 #include "Mu2e-STMDAQ/processing/udp.hh"
3 
4 // Constructor
5 UDP::UDP(const std::shared_ptr<AsyncLogger>& logger_,
6  const std::shared_ptr<STMdata>& stm_,
7  const std::shared_ptr<SignalHandler>& signal_,
8  const bool is_client_)
9  : logger(logger_)
10  , stm(stm_)
11  , signal(signal_)
12  , is_client(is_client_)
13  , // Server or client
14  ip((stm->master_config.ch_num) ? stm->udp_config.rcv_ip : stm->udp_config.snd_ip)
15  , // IP address
16  port((stm->master_config.ch_num) ? stm->udp_config.rcv_port
17  : stm->udp_config.snd_port)
18  , // Port
19  max_packet_num(stm->buffer_config.max_packet_num)
20  , wait_time(0.0)
21  , // Time waiting with no packets received
22  idle_timeout_ms(stm->udp_config.idle_timeout_ms)
23  , last_recv_time(std::chrono::steady_clock::now())
24 
25 {
26  // Create a UDP socket
27  socket_fd = socket(AF_INET, SOCK_DGRAM, 0);
28  // Handle socket creation failure
29  if(socket_fd < 0)
30  {
31  if(logger)
32  logger->log("UDP: socket creation failed. Exiting...", 0);
33  return;
34  }
35 
36  // Notify user
37  std::string type = (is_client) ? "client" : "server";
38  if(logger)
39  logger->log("UDP: Configuring UDP " + type + " for " +
40  stm->master_config.ch_name + " channel...",
41  1);
42 
43  // Setup socket
44  memset((char*)&address, 0, sizeof(address)); // zero out the structure
45  address.sin_family = AF_INET; // Set address family to IPv4
46  address.sin_addr.s_addr = inet_addr(ip.c_str()); // Set the IP address
47  address.sin_port = htons(port); // Convert and set the port number
48 
49  // Set socket to allow port re-use / to reuse port
50  int optval = 1;
51  setsockopt(socket_fd, SOL_SOCKET, SO_REUSEPORT, &optval, sizeof(optval));
52  // Get rmem_max and wmem_max
53  get_mem_max();
54 
55  // If creating a client
56  if(is_client)
57  {
58  configure_client();
59  }
60  // If creating a server
61  else
62  {
63  configure_server();
64  }
65 
66  // Register operations for OperationManager
67  register_operation("receive_data", [this](auto& b) { receive_data(b); });
68 }
69 
70 // Get kernel buffer memory values
71 void UDP::get_mem_max()
72 {
73  // Get system set RCVBUF size from system rmem_max
74  std::string size;
75  uint64_t size_val;
76  std::ifstream rmem_max("/proc/sys/net/core/rmem_max");
77  getline(rmem_max, size);
78  std::istringstream iss1(size);
79  iss1 >> size_val;
80  rmem_max.close(); // Close the file;
81  RMEM_MAX = size_val;
82  if(logger)
83  logger->log(
84  "UDP: /proc/sys/net/core/rmem_max = " + std::to_string(RMEM_MAX) + " bytes.",
85  1);
86 
87  // Get system set SNDBUF size from system wmem_max
88  std::ifstream wmem_max("/proc/sys/net/core/wmem_max");
89  getline(wmem_max, size);
90  iss1 >> size_val;
91  wmem_max.close(); // Close the file
92  WMEM_MAX = size_val;
93  if(logger)
94  logger->log(
95  "UDP: /proc/sys/net/core/wmem_max = " + std::to_string(WMEM_MAX) + " bytes.",
96  1);
97 }
98 
99 // Configure the server
100 void UDP::configure_server()
101 {
102  // Set the receive buffer size based on the system vale for rmem_max
103  socklen_t optlen = sizeof(RMEM_MAX);
104  setsockopt(socket_fd, SOL_SOCKET, SO_RCVBUF, &RMEM_MAX, optlen);
105  // Retrieve the actual buffer size set
106  int actual_rmem_max;
107  getsockopt(socket_fd, SOL_SOCKET, SO_RCVBUF, &actual_rmem_max, &optlen);
108  // Print the actual buffer size
109  if(logger)
110  logger->log("UDP: SO_RCVBUF set to: " + std::to_string(actual_rmem_max) + " (" +
111  std::to_string(actual_rmem_max / 2) +
112  ") bytes [RMEM_MAX = " + std::to_string(RMEM_MAX) + " byes].",
113  1);
114 
115  // Bind the socket to the server address
116  if(bind(socket_fd, (struct sockaddr*)&address, sizeof(address)) < 0)
117  {
118  if(logger)
119  logger->log("UDP: Unable to bind socket to IP: " +
120  std::string(inet_ntoa(address.sin_addr)) + ", Port: " +
121  std::to_string(ntohs(address.sin_port)) + ". Exiting....",
122  0); // Handle binding failure
123  return;
124  }
125 
126  // Get the current socket flags
127  int flags = fcntl(socket_fd, F_GETFL, 0);
128  // Set the socket to non-blocking mode
129  fcntl(socket_fd, F_SETFL, flags | O_NONBLOCK);
130 
131  // Allocate memory for mmsghdr structures
132  messages = new mmsghdr[max_packet_num];
133  // Allocate memory for iovec structures
134  iovecs = new iovec[max_packet_num];
135  // Set up the UDP server buffers
136  setup_buffers();
137 
138  // Set the server's polling structure
139  pfd.fd = socket_fd;
140  pfd.events = POLLIN;
141 
142  // Set recvmmsg timeout
143  read_timeout.tv_sec = stm->udp_config.recv_timeout_s;
144  read_timeout.tv_usec = stm->udp_config.recv_timeout_us;
145  timeout = {read_timeout.tv_sec, read_timeout.tv_usec};
146  setsockopt(socket_fd, SOL_SOCKET, SO_RCVTIMEO, &timeout, sizeof(timeout));
147 
148  // Print to user
149  if(logger)
150  logger->log("UDP: CREATED UDP SERVER. IP = " + ip + ", PORT = " +
151  std::to_string(port) + ", socket = " + std::to_string(socket_fd) +
152  ", timeout = " + std::to_string(read_timeout.tv_sec) + " s, " +
153  std::to_string(read_timeout.tv_usec) + " us.",
154  1);
155 }
156 
157 // Set up the UDP server buffers
158 // void UDP::setup_buffers() {
159 
160 // // Loop over number of packets per call
161 // for (size_t i = 0; i < max_packet_num; ++i) {
162 // iovecs[i].iov_len = MAX_PACKET_SIZE; // Set the size of each packet buffer
163 // messages[i].msg_hdr.msg_iov = &iovecs[i]; // Associate the iovec with the corresponding mmsghdr structure
164 // messages[i].msg_hdr.msg_iovlen = 1; // Indicate that each mmsghdr describes one buffer
165 // }
166 
167 // }
168 
169 // #include <cstring> // for std::memset
170 void UDP::setup_buffers()
171 {
172  for(size_t i = 0; i < max_packet_num; ++i)
173  {
174  // Clear everything so we don't have garbage in msg_name, msg_control, etc.
175  std::memset(&messages[i], 0, sizeof(mmsghdr));
176  std::memset(&iovecs[i], 0, sizeof(iovec));
177 
178  // Each iovec describes one packet buffer; length is fixed here
179  iovecs[i].iov_len = MAX_PACKET_SIZE;
180 
181  // Hook the iovec into the msghdr inside mmsghdr
182  messages[i].msg_hdr.msg_iov = &iovecs[i];
183  messages[i].msg_hdr.msg_iovlen = 1;
184 
185  // We don't care about the sender address → disable it explicitly
186  messages[i].msg_hdr.msg_name = nullptr;
187  messages[i].msg_hdr.msg_namelen = 0;
188 
189  // No ancillary data
190  messages[i].msg_hdr.msg_control = nullptr;
191  messages[i].msg_hdr.msg_controllen = 0;
192 
193  // msg_len will be filled in by recvmmsg
194  messages[i].msg_len = 0;
195  }
196 }
197 
198 // Configure the client
199 void UDP::configure_client()
200 {
201  // Set the receive buffer size based on the system vale for wmem_max
202  socklen_t optlen = sizeof(WMEM_MAX);
203  setsockopt(socket_fd, SOL_SOCKET, SO_SNDBUF, &WMEM_MAX, optlen);
204  // Retrieve the actual buffer size set
205  int actual_wmem_max;
206  getsockopt(socket_fd, SOL_SOCKET, SO_SNDBUF, &actual_wmem_max, &optlen);
207  // Print the actual buffer size
208  if(logger)
209  logger->log("UDP: SO_SNDBUF set to: " + std::to_string(actual_wmem_max) + " (" +
210  std::to_string(actual_wmem_max / 2) +
211  ") bytes [WMEM_MAX = " + std::to_string(WMEM_MAX) + " byes].",
212  1);
213 
214  // Print to user
215  if(logger)
216  logger->log("UDP: CREATED UDP CLIENT. IP = " + ip +
217  ", PORT = " + std::to_string(port) +
218  ", socket = " + std::to_string(socket_fd) + ".",
219  1);
220 }
221 
222 // Receive data from UDP server
223 void UDP::receive_data(std::shared_ptr<DataStruct>& buffer)
224 {
225  // Check correct buffer sizing
226  size_t required = max_packet_num * MAX_PACKET_SIZE;
227  size_t available = buffer->zs.size() * sizeof(int16_t);
228  if(available < required)
229  {
230  std::cerr << "ERROR: raw buffer too small: " << available << " bytes available, "
231  << required << " required.\n";
232  std::abort();
233  }
234 
235  // Set up iovec structures to point directly into DataStruct
236  for(size_t i = 0; i < max_packet_num; ++i)
237  {
238  // Point each iovec into the correct slice of the shared data buffer
239  iovecs[i].iov_base = buffer->zs.data() + (i * MAX_PACKET_LEN);
240  messages[i].msg_hdr.msg_iov = &iovecs[i];
241  // Reset per-call output fields
242  messages[i].msg_len = 0;
243  messages[i].msg_hdr.msg_flags = 0;
244  }
245 
246  // Ensure buffer is zero before call
247  buffer->zs_len = 0;
248  buffer->orig_len = 0;
249 
250  // Start and end wait times
251  if(!waiting_for_data)
252  start_wait = std::chrono::high_resolution_clock::now();
253 
254  // Packet counter
255  int packets_this_call = 0;
256 
257  // Loop until stop signal is triggered
258  while(!stop::should_stop())
259  {
260  // Poll for available data
261  int poll_ret = poll(&pfd, 1, stm->udp_config.poll_timeout_ms);
262 
263  // If poll failed
264  if(poll_ret < 0)
265  {
266  if(logger)
267  logger->log("UDP: UDP::receive_data - poll failed!", 0);
268  break;
269  }
270 
271  // If data is available
272  else if(poll_ret > 0)
273  {
274  // If we haven't yet received data
275  if(!has_received_data)
276  {
277  // Store waiting end time
278  end_wait = std::chrono::high_resolution_clock::now();
279  // ... and add to wait time counter
280  if(packets_this_call > 0)
281  {
282  wait_time += end_wait - start_wait;
283  add_wait(end_wait - start_wait);
284  }
285  }
286 
287  // Receive data from socket
288  int retval = recvmmsg(socket_fd,
289  &messages[packets_this_call],
290  max_packet_num - packets_this_call,
291  MSG_DONTWAIT,
292  &timeout);
293  // int retval = udp_debug::debug_recvmmsg(socket_fd, &messages[packets_this_call],
294  // max_packet_num - packets_this_call, MSG_DONTWAIT, &timeout);
295 
296  // If recvmmsg error
297  if(retval < 0)
298  {
299  char buf[256];
300  strerror_r(errno, buf, sizeof(buf));
301  std::string errorStr(buf);
302  if(logger)
303  logger->log(
304  "UDP: UDP::receive_data - recvmmsg failed with error code: " +
305  (std::to_string(errno)) + " " + errorStr,
306  0);
307  std::this_thread::sleep_for(std::chrono::seconds(3));
308  return;
309  }
310 
311  // Update number of packets received
312  packets_this_call += retval;
313 
314  // If packets update last receive time for idle timeout
315  if(retval > 0)
316  {
317  last_recv_time = std::chrono::steady_clock::now();
318  idle_timeout_logged = false;
319  }
320 
321  // Update size of received data in DataStruct
322  buffer->zs_len += retval * MAX_PACKET_LEN;
323  buffer->orig_len = buffer->zs_len;
324 
325  // Update total number of packets received
326  total_packets_received += retval;
327 
328  // If first data of call
329  if(!has_received_data)
330  has_received_data = true;
331 
332  // Tell user we're receiving data
333  if(logger && waiting_for_data)
334  logger->log("UDP: Receiving data...", 1);
335 
336  // Signal we're not waiting for data
337  waiting_for_data = false;
338 
339  // Reached maximum number of packets per call, so return
340  if(packets_this_call == max_packet_num)
341  return;
342  }
343 
344  // Else, if no data available
345  else
346  {
347  // If we were previously receiving data
348  if(!waiting_for_data)
349  {
350  // If we've not still received the first data yet
351  if(!has_received_data)
352  {
353  if(logger)
354  logger->log("UDP: Waiting for data from STM firmware...", 1);
355  }
356  else
357  {
358  // Start new waiting time
359  start_wait = std::chrono::high_resolution_clock::now();
360  // Warn user
361  if(logger)
362  logger->log(
363  "UDP: Stopped receiving from STM firmware. Waiting for "
364  "data...",
365  2);
366  }
367  }
368 
369  // Signal we're waiting for data
370  waiting_for_data = true;
371 
372  // If we were previously receiving data and have now gone idle
373  if(has_received_data && !idle_timeout_logged)
374  {
375  auto now = std::chrono::steady_clock::now();
376  auto idle_ms = std::chrono::duration_cast<std::chrono::milliseconds>(
377  now - last_recv_time)
378  .count();
379 
380  if(idle_ms > idle_timeout_ms)
381  {
382  idle_timeout_logged = true;
383  logger->log("UDP: Idle timeout after " + std::to_string(idle_ms) +
384  "ms, releasing buffer.",
385  2);
386  buffer->has_idle_timeout = true;
387  return;
388  }
389  }
390 
391  // Push what has already been pulled to the queue
392  if(packets_this_call > 0)
393  return;
394  }
395 
396  } // End while (!stop::should_stop())
397 
398  // If a signal has been caught and we've been waiting for data
399  if(waiting_for_data)
400  {
401  // Get waiting end time
402  end_wait = std::chrono::high_resolution_clock::now();
403  // ... and add to wait time counter
404  wait_time += end_wait - start_wait;
405  add_wait(end_wait - start_wait);
406  }
407 
408  return;
409 }
410 
411 // Function to send a single UDP packet
412 void UDP::send_packet(const std::vector<int16_t>& data_vec)
413 {
414  // Prepare the message header
415  struct iovec iov;
416  iov.iov_base = const_cast<int16_t*>(data_vec.data());
417  iov.iov_len = data_vec.size() * sizeof(int16_t);
418 
419  struct msghdr msg;
420  memset(&msg, 0, sizeof(msg));
421  msg.msg_name = &address;
422  msg.msg_namelen = sizeof(address);
423  msg.msg_iov = &iov;
424  msg.msg_iovlen = 1;
425 
426  ssize_t sent_bytes = sendmsg(socket_fd, &msg, 0);
427 
428  if(__builtin_expect(sent_bytes == -1, 0))
429  {
430  if(logger)
431  logger->log(std::string("send_packet error: ") + strerror(errno), 0);
432  }
433 }