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