otsdaq-mu2e-stm  5.02.01
udp.hh
1 #ifndef UDP_hh_
2 #define UDP_hh_
3 
4 #include <iostream>
5 #include <thread>
6 #include <mutex>
7 #include <condition_variable>
8 #include <vector>
9 #include <queue>
10 #include <atomic>
11 #include <memory>
12 #include <cstring>
13 #include <array>
14 #include <functional>
15 #include <chrono>
16 #include <iomanip>
17 #include <sched.h>
18 #include <unistd.h>
19 #include <sys/socket.h>
20 #include <netinet/in.h>
21 #include <arpa/inet.h>
22 #include <linux/udp.h>
23 #include <errno.h>
24 #include <sys/types.h>
25 #include <sys/mman.h>
26 #include <fcntl.h>
27 #include <poll.h>
28 #include <sys/sysinfo.h>
29 #include <cmath>
30 #include <signal.h>
31 #include <sys/select.h>
32 #include <stdio.h>
33 #include <string.h>
34 
35 // Async Logger code
36 #include "Mu2e-STMDAQ/utils/async_logger.hh"
37 // Signal handler header
38 #include "Mu2e-STMDAQ/utils/signal_handler.hh"
39 // STM data header
40 #include "Mu2e-STMDAQ/config/stm_data.hh"
41 // Operations base header
42 #include "Mu2e-STMDAQ/processing/operations_base.hh"
43 // UDP debug code
44 #include "Mu2e-STMDAQ/debug/udp_debug.hh"
45 
46 
47 // UDP Class for socket setup and data reception
48 class UDP : public OperationMap {
49 
50 private:
51 
52  // Async Logger
53  const std::shared_ptr<AsyncLogger> logger;
54 
55  // STM data info
56  const std::shared_ptr<STMdata>& stm;
57 
58  // Signal Handler
59  const std::shared_ptr<SignalHandler>& signal;
60 
61  // Is this socket a client (if not, then a server)
62  const bool is_client;
63 
64  // Socket IP address
65  const std::string ip;
66  // Socket port
67  const uint16_t port;
68 
69  // File descriptor for the UDP socket
70  int socket_fd;
71  // Structure to hold the socket's address and port
72  sockaddr_in address;
73 
74  // System rmem_max
75  size_t RMEM_MAX;
76  // System wmem_max
77  size_t WMEM_MAX;
78 
79  // The time value of the UDP timeout
80  struct timeval read_timeout;
81  struct timespec timeout;
82 
83  // Server polling structure
84  struct pollfd pfd;
85 
86  // Pointer to an array of mmsghdr structures for recvmmsg
87  struct mmsghdr* messages;
88  // Pointer to an array of iovec structures describing the buffers
89  struct iovec* iovecs;
90 
91  // Maximum number of packets in buffer
92  const size_t max_packet_num;
93 
94  // Have we received the first packet?
95  bool has_received_data = false;
96 
97  // Are we waiting for data from the socket
98  bool waiting_for_data = false;
99 
100  // The total number of packets received in this run
101  size_t total_packets_received = 0;
102 
103  // Timing
104  std::chrono::duration<double> wait_time;
105  std::chrono::time_point<std::chrono::high_resolution_clock> start_wait =
106  std::chrono::high_resolution_clock::now();
107  std::chrono::time_point<std::chrono::high_resolution_clock> end_wait = start_wait;
108 
109  const size_t idle_timeout_ms;
110  std::chrono::steady_clock::time_point last_recv_time;
111  bool idle_timeout_logged = false;
112 
113 
114  std::atomic<double> wait_secs{0.0};
115  void add_wait(std::chrono::duration<double> delta){
116  double old_val = wait_secs.load(std::memory_order_relaxed);
117  double new_val;
118  do {
119  new_val = old_val + delta.count();
120  } while (!wait_secs.compare_exchange_weak(
121  old_val, new_val,
122  std::memory_order_release,
123  std::memory_order_relaxed));
124  }
125 
126 public:
127 
128  // Constructor
129  UDP(const std::shared_ptr<AsyncLogger>& logger_,
130  const std::shared_ptr<STMdata>& stm_,
131  const std::shared_ptr<SignalHandler>& signal_,
132  const bool is_client_);
133 
134  // Destructor
135  ~UDP() {
136  // Close the socket when the object is destroyed
137  close(socket_fd);
138  // If server
139  if (!is_client){
140  delete[] messages; // Free the allocated mmsghdr structures
141  delete[] iovecs; // Free the allocated iovec structures
142  }
143  if (logger) logger->log("UDP server received a total of " + std::to_string(total_packets_received) + " packets.",1);
144  if (logger) logger->log("Total UDP server dead time = " + std::to_string(wait_time.count()) + " seconds.",1);
145  std::cout << "UDP destructor called.\n";
146  }
147 
148  // Get the systems kernel buffer values
149  void get_mem_max();
150 
151  // Configure server
152  void configure_server();
153 
154  // Set up the UDP server buffers
155  void setup_buffers();
156 
157  // Configure client
158  void configure_client();
159 
160  // Receive data from UDP server
161  void receive_data(std::shared_ptr<DataStruct>& buffer);
162 
163  // Function to send single UDP packet
164  void send_packet(const std::vector<int16_t>& data_vec);
165 
166  // Return the address of the total wait time
167  std::chrono::duration<double>* get_wait_time(){
168  return &wait_time;
169  }
170 
171  std::chrono::duration<double> get_wait() {
172  return std::chrono::duration<double>(wait_secs.load(std::memory_order_acquire));
173  }
174 
175 };
176 
177 #endif //UDP_hh_
Definition: udp.hh:48