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