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