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