otsdaq-mu2e-stm  5.02.01
best_so_far.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  size_t size; // Actual size of the data used
22  char metadata[128]; // Fixed-size metadata array
23 
24  DataStruct(size_t buffer_size) : data(buffer_size), size(0)
25  {
26  std::fill(std::begin(metadata), std::end(metadata), '\0');
27  }
28 };
29 
30 // Ring buffer for thread-safe buffer management
31 template<typename T>
33 {
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  void push(const std::shared_ptr<T>& item)
49  {
50  std::unique_lock<std::mutex> lock(mutex);
51  cv_push.wait(lock, [this]() { return stop || size < capacity; });
52  if(stop)
53  return;
54  buffer[head] = item;
55  head = (head + 1) % capacity;
56  ++size;
57  cv_pop.notify_one();
58  }
59 
60  std::shared_ptr<T> pop()
61  {
62  std::unique_lock<std::mutex> lock(mutex);
63  cv_pop.wait(lock, [this]() { return stop || size > 0; });
64  if(size == 0)
65  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  {
76  std::lock_guard<std::mutex> lock(mutex);
77  stop = true;
78  }
79  cv_push.notify_all();
80  cv_pop.notify_all();
81  }
82 
83  bool is_empty()
84  {
85  std::lock_guard<std::mutex> lock(mutex);
86  return size == 0;
87  }
88 };
89 
90 // Pre-allocated buffer pool for memory reuse
92 {
93  private:
94  std::vector<std::shared_ptr<DataStruct>> pool;
95  std::mutex mutex;
96 
97  public:
98  BufferPool(size_t pool_size, size_t buffer_size)
99  {
100  for(size_t i = 0; i < pool_size; ++i)
101  {
102  pool.emplace_back(std::make_shared<DataStruct>(buffer_size));
103  }
104  }
105 
106  std::shared_ptr<DataStruct> acquire()
107  {
108  std::lock_guard<std::mutex> lock(mutex);
109  if(pool.empty())
110  {
111  return std::make_shared<DataStruct>(60 * 1024 * 1024 / sizeof(int16_t));
112  }
113  else
114  {
115  auto buffer = pool.back();
116  pool.pop_back();
117  return buffer;
118  }
119  }
120 
121  void release(std::shared_ptr<DataStruct> buffer)
122  {
123  std::lock_guard<std::mutex> lock(mutex);
124  pool.push_back(buffer);
125  }
126 };
127 
128 // Global total data tracker
129 std::atomic<uint64_t> globalTotalBytes{0};
130 const uint64_t targetTotalBytes = 10L * 1024 * 1024 * 1024; // 10 GB
131 std::atomic<bool> kernelDone{false}; // Signal when kernelBufferWorker is finished
132 
133 // Pin thread to specific CPU core
134 void pin_thread_to_core(size_t core_id)
135 {
136  cpu_set_t cpuset;
137  CPU_ZERO(&cpuset);
138  CPU_SET(core_id, &cpuset);
139  pthread_setaffinity_np(pthread_self(), sizeof(cpu_set_t), &cpuset);
140 }
141 
142 // General worker function for processing stages
143 void workerThread(RingBuffer<DataStruct>& inputQueue,
144  RingBuffer<DataStruct>& outputQueue,
145  BufferPool& pool,
146  const char* stage,
147  std::atomic<uint64_t>& totalBytesProcessed,
148  size_t core_id)
149 {
150  pin_thread_to_core(core_id);
151  auto start_time = std::chrono::high_resolution_clock::now();
152  while(!kernelDone || !inputQueue.is_empty())
153  {
154  std::vector<std::shared_ptr<DataStruct>> batch;
155  {
156  auto buffer = inputQueue.pop();
157  if(buffer)
158  {
159  batch.push_back(buffer);
160  }
161  }
162  for(auto& buffer : batch)
163  {
164  // Process the data and append metadata
165  strncat(buffer->metadata,
166  stage,
167  sizeof(buffer->metadata) - strlen(buffer->metadata) - 1);
168 
169  // Update the total bytes processed
170  uint64_t bytesProcessed = buffer->size * sizeof(int16_t);
171  totalBytesProcessed += bytesProcessed;
172 
173  // Push the processed buffer to the next stage
174  outputQueue.push(buffer);
175  }
176  }
177  auto end_time = std::chrono::high_resolution_clock::now();
178  std::chrono::duration<double> elapsed = end_time - start_time;
179  double gbits = totalBytesProcessed.load() * 8.0 / 1e9; // Convert to Gbits
180  std::cout << "[" << stage << "] Processing speed: " << std::fixed
181  << std::setprecision(2) << (gbits / elapsed.count()) << " Gbit/s\n";
182 }
183 
184 // Kernel buffer simulation function
185 void kernelBufferWorker(RingBuffer<DataStruct>& outputQueue,
186  BufferPool& pool,
187  size_t core_id)
188 {
189  pin_thread_to_core(core_id);
190  while(globalTotalBytes < targetTotalBytes)
191  {
192  auto buffer = pool.acquire();
193  buffer->size = (rand() % (60 * 1024 * 1024 / sizeof(int16_t))) +
194  1000; // Simulate variable size
195  memset(buffer->data.data(), rand() % 100, buffer->size * sizeof(int16_t));
196  strncpy(buffer->metadata,
197  "Generated by kernelBufferWorker",
198  sizeof(buffer->metadata));
199 
200  uint64_t bytesGenerated = buffer->size * sizeof(int16_t);
201  globalTotalBytes += bytesGenerated;
202 
203  outputQueue.push(buffer);
204  }
205  kernelDone = true; // Signal that kernel is done producing data
206 }
207 
208 // Consumer thread to ensure the final buffer is emptied
209 void consumerThread(RingBuffer<DataStruct>& inputQueue, BufferPool& pool, size_t core_id)
210 {
211  pin_thread_to_core(core_id);
212  while(!kernelDone || !inputQueue.is_empty())
213  {
214  auto buffer = inputQueue.pop();
215  if(!buffer)
216  continue; // Wait for remaining data
217 
218  // Simulate consumption of the final data
219  pool.release(buffer);
220  }
221  std::cout << "[Consumer] Final buffer emptied.\n";
222 }
223 
224 int main()
225 {
226  const size_t max_queue_size = 20; // Increased maximum size for each ring buffer
227  const size_t buffer_pool_size = 30; // Increased pre-allocated buffer pool size
228  const size_t buffer_size =
229  60 * 1024 * 1024 / sizeof(int16_t); // Buffer size in elements
230 
231  // Create buffer pool
232  BufferPool pool(buffer_pool_size, buffer_size);
233 
234  // Define ring buffers for each step
235  RingBuffer<DataStruct> kernelBufferQueue(max_queue_size);
236  RingBuffer<DataStruct> checkDataQueue(max_queue_size);
237  RingBuffer<DataStruct> formEventsQueue(max_queue_size);
238  RingBuffer<DataStruct> rawBufferQueue(max_queue_size);
239  RingBuffer<DataStruct> zsBufferQueue(max_queue_size);
240  RingBuffer<DataStruct> mwdBufferQueue(max_queue_size);
241 
242  // Track total bytes processed for each thread
243  std::atomic<uint64_t> totalBytesKernel{0};
244  std::atomic<uint64_t> totalBytesCheck{0};
245  std::atomic<uint64_t> totalBytesForm{0};
246  std::atomic<uint64_t> totalBytesRaw{0};
247  std::atomic<uint64_t> totalBytesZS{0};
248  std::atomic<uint64_t> totalBytesMWD{0};
249 
250  // Launch worker threads for each stage
251  std::thread kernelThread(
252  kernelBufferWorker, std::ref(kernelBufferQueue), std::ref(pool), 0);
253  std::thread checkDataThread(workerThread,
254  std::ref(kernelBufferQueue),
255  std::ref(checkDataQueue),
256  std::ref(pool),
257  " | Checked",
258  std::ref(totalBytesCheck),
259  1);
260  std::thread formEventsThread(workerThread,
261  std::ref(checkDataQueue),
262  std::ref(formEventsQueue),
263  std::ref(pool),
264  " | Events Formed",
265  std::ref(totalBytesForm),
266  2);
267  std::thread rawProcessThread(workerThread,
268  std::ref(formEventsQueue),
269  std::ref(rawBufferQueue),
270  std::ref(pool),
271  " | Algorithm Processed",
272  std::ref(totalBytesRaw),
273  3);
274  std::thread zsProcessThread(workerThread,
275  std::ref(rawBufferQueue),
276  std::ref(zsBufferQueue),
277  std::ref(pool),
278  " | ZS Processed",
279  std::ref(totalBytesZS),
280  4);
281  std::thread mwdProcessThread(workerThread,
282  std::ref(zsBufferQueue),
283  std::ref(mwdBufferQueue),
284  std::ref(pool),
285  " | MWD Processed",
286  std::ref(totalBytesMWD),
287  5);
288  std::thread finalConsumerThread(
289  consumerThread, std::ref(mwdBufferQueue), std::ref(pool), 6);
290 
291  // Wait for threads to complete
292  kernelThread.join();
293  checkDataThread.join();
294  formEventsThread.join();
295  rawProcessThread.join();
296  zsProcessThread.join();
297  mwdProcessThread.join();
298  finalConsumerThread.join();
299 
300  // Shutdown all queues
301  kernelBufferQueue.shutdown();
302  checkDataQueue.shutdown();
303  formEventsQueue.shutdown();
304  rawBufferQueue.shutdown();
305  zsBufferQueue.shutdown();
306  mwdBufferQueue.shutdown();
307 
308  return 0;
309 }
Definition: data.hh:4