otsdaq-mu2e-stm  5.02.01
best_so_far_wComments.cpp
1 #define _GNU_SOURCE // Enable GNU-specific features
2 #include <sched.h> // Include for CPU affinity management
3 #include <unistd.h> // Include for POSIX system calls
4 #include <array> // Include for fixed-size arrays
5 #include <atomic> // Include for atomic variables
6 #include <chrono> // Include for timing utilities
7 #include <condition_variable> // Include for condition variables
8 #include <cstring> // Include for memory operations
9 #include <functional> // Include for function objects
10 #include <iomanip> // Include for formatting output
11 #include <iostream> // Include for standard input/output operations
12 #include <memory> // Include for smart pointers
13 #include <mutex> // Include for mutex synchronization
14 #include <queue> // Include for queue data structure
15 #include <thread> // Include for thread support
16 #include <vector> // Include for dynamic arrays
17 
18 // Data structure representing a buffer with metadata
19 struct DataStruct
20 {
21  std::vector<int16_t> data; // Dynamically sized data buffer
22  size_t size; // Actual size of the data used
23  char metadata[128]; // Fixed-size metadata array
24 
25  DataStruct(size_t buffer_size) : data(buffer_size), size(0)
26  { // Constructor to initialize data buffer
27  std::fill(std::begin(metadata),
28  std::end(metadata),
29  '\0'); // Initialize metadata to null characters
30  }
31 };
32 
33 // Ring buffer for thread-safe buffer management
34 template<typename T>
35 class RingBuffer
36 {
37  private:
38  std::vector<std::shared_ptr<T>> buffer; // Circular buffer storage
39  size_t head = 0; // Index of the next insertion point
40  size_t tail = 0; // Index of the next removal point
41  size_t capacity; // Maximum capacity of the buffer
42  std::atomic<size_t> size = 0; // Current size of the buffer
43  std::mutex mutex; // Mutex for synchronization
44  std::condition_variable cv_push; // Condition variable for producers
45  std::condition_variable cv_pop; // Condition variable for consumers
46  bool stop = false; // Flag to stop buffer operations
47 
48  public:
49  // Constructor to initialize buffer
50  explicit RingBuffer(size_t max_size) : capacity(max_size), buffer(max_size) {}
51 
52  // Destructor for logging
53  ~RingBuffer() { std::cout << "RingBuffer destructor called.\n"; }
54 
55  // Add an item to the buffer
56  void push(const std::shared_ptr<T>& item)
57  {
58  std::unique_lock<std::mutex> lock(mutex); // Lock for thread safety
59  cv_push.wait(lock, [this]() {
60  return stop || size < capacity;
61  }); // Wait if buffer is full
62  if(stop)
63  return; // Exit if stop flag is set
64  buffer[head] = item; // Add item to buffer
65  head = (head + 1) % capacity; // Update head index
66  ++size; // Increment buffer size
67  cv_pop.notify_one(); // Notify consumers
68  }
69 
70  // Remove and return an item from the buffer
71  std::shared_ptr<T> pop()
72  {
73  std::unique_lock<std::mutex> lock(mutex); // Lock for thread safety
74  cv_pop.wait(lock,
75  [this]() { return stop || size > 0; }); // Wait if buffer is empty
76  if(size == 0)
77  return nullptr; // Return null if buffer is empty
78  auto item = buffer[tail]; // Retrieve item from buffer
79  tail = (tail + 1) % capacity; // Update tail index
80  --size; // Decrement buffer size
81  cv_push.notify_one(); // Notify producers
82  return item; // Return the item
83  }
84 
85  // Stop buffer operations
86  void shutdown()
87  {
88  {
89  std::lock_guard<std::mutex> lock(mutex); // Lock for thread safety
90  stop = true; // Set stop flag
91  }
92  cv_push.notify_all(); // Notify all waiting producers
93  cv_pop.notify_all(); // Notify all waiting consumers
94  }
95 
96  bool is_empty()
97  { // Check if buffer is empty
98  std::lock_guard<std::mutex> lock(mutex); // Lock for thread safety
99  return size == 0; // Return true if buffer size is zero
100  }
101 };
102 
103 // Pre-allocated buffer pool for memory reuse
104 class BufferPool
105 {
106  private:
107  std::vector<std::shared_ptr<DataStruct>> pool; // Pool of pre-allocated buffers
108  std::mutex mutex; // Mutex for thread safety
109 
110  public:
111  // Constructor to initialize buffer pool
112  BufferPool(size_t pool_size, size_t buffer_size)
113  {
114  for(size_t i = 0; i < pool_size; ++i)
115  { // Create specified number of buffers
116  pool.emplace_back(
117  std::make_shared<DataStruct>(buffer_size)); // Add buffers to pool
118  }
119  }
120 
121  // Destructor for logging
122  ~BufferPool() { std::cout << "BufferPool destructor called.\n"; }
123 
124  // Acquire a buffer from the pool
125  std::shared_ptr<DataStruct> acquire()
126  {
127  std::lock_guard<std::mutex> lock(mutex); // Lock for thread safety
128  if(pool.empty())
129  { // If pool is empty
130  return std::make_shared<DataStruct>(60 * 1024 * 1024 /
131  sizeof(int16_t)); // Create a new buffer
132  }
133  else
134  {
135  auto buffer = pool.back(); // Retrieve buffer from the pool
136  pool.pop_back(); // Remove buffer from pool
137  return buffer; // Return the buffer
138  }
139  }
140 
141  // Return a buffer to the pool
142  void release(std::shared_ptr<DataStruct> buffer)
143  {
144  std::lock_guard<std::mutex> lock(mutex); // Lock for thread safety
145  pool.push_back(buffer); // Add buffer back to pool
146  }
147 };
148 
149 // Global total data tracker
150 std::atomic<uint64_t> globalTotalBytes{0}; // Total bytes processed
151 const uint64_t targetTotalBytes = 1000L * 1024 * 1024 * 1024; // 10 GB target
152 std::atomic<bool> kernelDone{false}; // Signal when kernelBufferWorker is finished
153 
154 // Pin thread to specific CPU core
155 void pin_thread_to_core(size_t core_id)
156 {
157  cpu_set_t cpuset; // Define CPU set
158  CPU_ZERO(&cpuset); // Clear CPU set
159  CPU_SET(core_id, &cpuset); // Add specified core to set
160  pthread_setaffinity_np(
161  pthread_self(), sizeof(cpu_set_t), &cpuset); // Set thread affinity
162 }
163 
164 // General worker function for processing stages
165 void workerThread(RingBuffer<DataStruct>& inputQueue,
166  RingBuffer<DataStruct>& outputQueue,
167  BufferPool& pool,
168  const char* stage,
169  std::atomic<uint64_t>& totalBytesProcessed,
170  size_t core_id)
171 {
172  pin_thread_to_core(core_id); // Bind thread to specific core
173  auto start_time = std::chrono::high_resolution_clock::now(); // Start timing
174  while(!kernelDone || !inputQueue.is_empty())
175  { // Continue until kernel is done and input queue is empty
176  auto buffer = inputQueue.pop(); // Retrieve a buffer from the input queue
177  if(buffer)
178  { // If buffer is valid
179  strncat(buffer->metadata,
180  stage,
181  sizeof(buffer->metadata) - strlen(buffer->metadata) -
182  1); // Append stage info to metadata
183  uint64_t bytesProcessed =
184  buffer->size * sizeof(int16_t); // Calculate bytes processed
185  totalBytesProcessed += bytesProcessed; // Update total bytes processed
186  outputQueue.push(buffer); // Push buffer to output queue
187  }
188  }
189  auto end_time = std::chrono::high_resolution_clock::now(); // End timing
190  std::chrono::duration<double> elapsed =
191  end_time - start_time; // Calculate elapsed time
192  double gbits = totalBytesProcessed.load() * 8.0 / 1e9; // Convert bytes to Gbits
193  std::cout << "[" << stage << "] Processing speed: " << std::fixed
194  << std::setprecision(2) << (gbits / elapsed.count())
195  << " Gbit/s\n"; // Print processing speed
196 }
197 
198 // Kernel buffer simulation function
199 void kernelBufferWorker(RingBuffer<DataStruct>& outputQueue,
200  BufferPool& pool,
201  size_t core_id)
202 {
203  pin_thread_to_core(core_id); // Bind thread to specific core
204  while(globalTotalBytes < targetTotalBytes)
205  { // Continue until total bytes processed reaches target
206  auto buffer = pool.acquire(); // Acquire a buffer from the pool
207  buffer->size = (rand() % (60 * 1024 * 1024 / sizeof(int16_t))) +
208  1000; // Simulate variable buffer size
209  memset(buffer->data.data(),
210  rand() % 100,
211  buffer->size * sizeof(int16_t)); // Fill buffer with simulated data
212  strncpy(buffer->metadata,
213  "Generated by kernelBufferWorker",
214  sizeof(buffer->metadata)); // Add metadata
215  uint64_t bytesGenerated =
216  buffer->size * sizeof(int16_t); // Calculate bytes generated
217  globalTotalBytes += bytesGenerated; // Update global total bytes processed
218  outputQueue.push(buffer); // Push buffer to output queue
219  }
220  kernelDone = true; // Signal that kernel is done producing data
221 }
222 
223 // Consumer thread to ensure the final buffer is emptied
224 void consumerThread(RingBuffer<DataStruct>& inputQueue, BufferPool& pool, size_t core_id)
225 {
226  pin_thread_to_core(core_id); // Bind thread to specific core
227  while(!kernelDone || !inputQueue.is_empty())
228  { // Continue until kernel is done and input queue is empty
229  auto buffer = inputQueue.pop(); // Retrieve a buffer from the input queue
230  if(buffer)
231  { // If buffer is valid
232  pool.release(buffer); // Return buffer to pool
233  }
234  }
235  std::cout << "[Consumer] Final buffer emptied.\n"; // Print completion message
236 }
237 
238 int main()
239 {
240  const size_t max_queue_size = 20; // Set maximum size for each ring buffer
241  const size_t buffer_pool_size = 30; // Set number of pre-allocated buffers
242  const size_t buffer_size =
243  60 * 1024 * 1024 / sizeof(int16_t); // Set size of each buffer in elements
244 
245  BufferPool pool(buffer_pool_size, buffer_size); // Create buffer pool
246  RingBuffer<DataStruct> kernelBufferQueue(
247  max_queue_size); // Create kernel buffer queue
248  RingBuffer<DataStruct> checkDataQueue(max_queue_size); // Create check data queue
249  RingBuffer<DataStruct> formEventsQueue(max_queue_size); // Create form events queue
250  RingBuffer<DataStruct> rawBufferQueue(max_queue_size); // Create raw buffer queue
251  RingBuffer<DataStruct> zsBufferQueue(max_queue_size); // Create ZS buffer queue
252  RingBuffer<DataStruct> mwdBufferQueue(max_queue_size); // Create MWD buffer queue
253 
254  std::atomic<uint64_t> totalBytesKernel{
255  0}; // Initialize total bytes processed for kernel
256  std::atomic<uint64_t> totalBytesCheck{
257  0}; // Initialize total bytes processed for check data
258  std::atomic<uint64_t> totalBytesForm{
259  0}; // Initialize total bytes processed for form events
260  std::atomic<uint64_t> totalBytesRaw{
261  0}; // Initialize total bytes processed for raw buffer
262  std::atomic<uint64_t> totalBytesZS{
263  0}; // Initialize total bytes processed for ZS buffer
264  std::atomic<uint64_t> totalBytesMWD{
265  0}; // Initialize total bytes processed for MWD buffer
266 
267  std::thread kernelThread(kernelBufferWorker,
268  std::ref(kernelBufferQueue),
269  std::ref(pool),
270  0); // Launch kernel buffer worker thread
271  std::thread checkDataThread(workerThread,
272  std::ref(kernelBufferQueue),
273  std::ref(checkDataQueue),
274  std::ref(pool),
275  " | Checked",
276  std::ref(totalBytesCheck),
277  1); // Launch check data worker thread
278  std::thread formEventsThread(workerThread,
279  std::ref(checkDataQueue),
280  std::ref(formEventsQueue),
281  std::ref(pool),
282  " | Events Formed",
283  std::ref(totalBytesForm),
284  2); // Launch form events worker thread
285  std::thread rawProcessThread(workerThread,
286  std::ref(formEventsQueue),
287  std::ref(rawBufferQueue),
288  std::ref(pool),
289  " | Algorithm Processed",
290  std::ref(totalBytesRaw),
291  3); // Launch raw processing worker thread
292  std::thread zsProcessThread(workerThread,
293  std::ref(rawBufferQueue),
294  std::ref(zsBufferQueue),
295  std::ref(pool),
296  " | ZS Processed",
297  std::ref(totalBytesZS),
298  4); // Launch ZS processing worker thread
299  std::thread mwdProcessThread(workerThread,
300  std::ref(zsBufferQueue),
301  std::ref(mwdBufferQueue),
302  std::ref(pool),
303  " | MWD Processed",
304  std::ref(totalBytesMWD),
305  5); // Launch MWD processing worker thread
306  std::thread finalConsumerThread(consumerThread,
307  std::ref(mwdBufferQueue),
308  std::ref(pool),
309  6); // Launch consumer thread
310 
311  // Wait for all threads to complete
312  kernelThread.join(); // Wait for kernel thread to finish
313  checkDataThread.join(); // Wait for check data thread to finish
314  formEventsThread.join(); // Wait for form events thread to finish
315  rawProcessThread.join(); // Wait for raw processing thread to finish
316  zsProcessThread.join(); // Wait for ZS processing thread to finish
317  mwdProcessThread.join(); // Wait for MWD processing thread to finish
318  finalConsumerThread.join(); // Wait for consumer thread to finish
319 
320  // Shutdown all queues
321  kernelBufferQueue.shutdown(); // Shut down kernel buffer queue
322  checkDataQueue.shutdown(); // Shut down check data queue
323  formEventsQueue.shutdown(); // Shut down form events queue
324  rawBufferQueue.shutdown(); // Shut down raw buffer queue
325  zsBufferQueue.shutdown(); // Shut down ZS buffer queue
326  mwdBufferQueue.shutdown(); // Shut down MWD buffer queue
327 
328  std::cout
329  << "[Main] All threads completed successfully.\n"; // Print final success message
330  return 0; // Exit program
331 }
Definition: data.hh:4