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