otsdaq-mu2e-stm  5.02.01
thread_buffer_management_bpool.cpp
1 #include <algorithm>
2 #include <array>
3 #include <atomic>
4 #include <condition_variable>
5 #include <cstring>
6 #include <iostream>
7 #include <memory>
8 #include <mutex>
9 #include <queue>
10 #include <thread>
11 #include <vector>
12 
13 // Data structure representing a buffer with metadata
14 struct DataStruct
15 {
16  std::vector<int16_t> data; // Dynamically sized data buffer
17  size_t size; // Actual size of the data used
18  char metadata[128]; // Fixed-size metadata array
19 
20  DataStruct(size_t buffer_size) : data(buffer_size), size(0)
21  {
22  std::fill(std::begin(metadata), std::end(metadata), '\0');
23  }
24 };
25 
26 // Ring buffer for thread-safe buffer management
27 template<typename T>
28 class RingBuffer
29 {
30  private:
31  std::vector<std::shared_ptr<T>> buffer;
32  size_t head = 0;
33  size_t tail = 0;
34  size_t capacity;
35  std::atomic<size_t> size = 0;
36  std::mutex mutex;
37  std::condition_variable cv_push;
38  std::condition_variable cv_pop;
39  bool stop = false;
40 
41  public:
42  explicit RingBuffer(size_t max_size) : capacity(max_size), buffer(max_size) {}
43 
44  void push(const std::shared_ptr<T>& item)
45  {
46  std::unique_lock<std::mutex> lock(mutex);
47  cv_push.wait(lock, [this]() { return stop || size < capacity; });
48  if(stop)
49  return;
50  buffer[head] = item;
51  head = (head + 1) % capacity;
52  ++size;
53  cv_pop.notify_one();
54  }
55 
56  std::shared_ptr<T> pop()
57  {
58  std::unique_lock<std::mutex> lock(mutex);
59  cv_pop.wait(lock, [this]() { return stop || size > 0; });
60  if(size == 0)
61  return nullptr;
62  auto item = buffer[tail];
63  tail = (tail + 1) % capacity;
64  --size;
65  cv_push.notify_one();
66  return item;
67  }
68 
69  void shutdown()
70  {
71  {
72  std::lock_guard<std::mutex> lock(mutex);
73  stop = true;
74  }
75  cv_push.notify_all();
76  cv_pop.notify_all();
77  }
78 };
79 
80 // Pre-allocated buffer pool for memory reuse
81 class BufferPool
82 {
83  private:
84  std::vector<std::shared_ptr<DataStruct>> pool;
85  std::mutex mutex;
86 
87  public:
88  BufferPool(size_t pool_size, size_t buffer_size)
89  {
90  for(size_t i = 0; i < pool_size; ++i)
91  {
92  pool.emplace_back(std::make_shared<DataStruct>(buffer_size));
93  }
94  }
95 
96  std::shared_ptr<DataStruct> acquire()
97  {
98  std::lock_guard<std::mutex> lock(mutex);
99  if(pool.empty())
100  {
101  return std::make_shared<DataStruct>(60 * 1024 * 1024 / sizeof(int16_t));
102  }
103  else
104  {
105  auto buffer = pool.back();
106  pool.pop_back();
107  return buffer;
108  }
109  }
110 
111  void release(std::shared_ptr<DataStruct> buffer)
112  {
113  std::lock_guard<std::mutex> lock(mutex);
114  pool.push_back(buffer);
115  }
116 };
117 
118 // Function for Kernel Buffer and UDP Server
119 auto kernelBufferThread = [](RingBuffer<DataStruct>& bufferQueue, BufferPool& pool) {
120  while(true)
121  {
122  auto buffer = pool.acquire();
123  buffer->size = (rand() % (60 * 1024 * 1024 / sizeof(int16_t))) +
124  1000; // Simulate variable size
125  std::fill(
126  buffer->data.begin(), buffer->data.begin() + buffer->size, rand() % 100);
127  strncpy(buffer->metadata,
128  "Generated by kernelBufferThread",
129  sizeof(buffer->metadata));
130  bufferQueue.push(buffer);
131  std::this_thread::sleep_for(std::chrono::milliseconds(5));
132  }
133 };
134 
135 // Function for Data Checking
136 auto checkDataThread = [](RingBuffer<DataStruct>& inputQueue,
137  RingBuffer<DataStruct>& outputQueue,
138  BufferPool& pool) {
139  while(true)
140  {
141  auto buffer = inputQueue.pop();
142  if(!buffer)
143  break;
144  strncpy(buffer->metadata + strlen(buffer->metadata),
145  " | Checked",
146  sizeof(buffer->metadata) - strlen(buffer->metadata));
147  outputQueue.push(buffer);
148  }
149 };
150 
151 // Function for Forming Events
152 auto formEventsThread = [](RingBuffer<DataStruct>& inputQueue,
153  RingBuffer<DataStruct>& outputQueue,
154  BufferPool& pool) {
155  while(true)
156  {
157  auto buffer = inputQueue.pop();
158  if(!buffer)
159  break;
160  strncpy(buffer->metadata + strlen(buffer->metadata),
161  " | Events Formed",
162  sizeof(buffer->metadata) - strlen(buffer->metadata));
163  outputQueue.push(buffer);
164  }
165 };
166 
167 // Function for Processing Algorithms
168 auto processAlgorithmThread = [](RingBuffer<DataStruct>& inputQueue,
169  RingBuffer<DataStruct>& outputQueue,
170  BufferPool& pool) {
171  while(true)
172  {
173  auto buffer = inputQueue.pop();
174  if(!buffer)
175  break;
176  strncpy(buffer->metadata + strlen(buffer->metadata),
177  " | Algorithm Processed",
178  sizeof(buffer->metadata) - strlen(buffer->metadata));
179  outputQueue.push(buffer);
180  }
181 };
182 
183 int main()
184 {
185  const size_t max_queue_size = 10; // Maximum size for each ring buffer
186  const size_t buffer_pool_size = 20; // Pre-allocated buffer pool size
187  const size_t buffer_size =
188  60 * 1024 * 1024 / sizeof(int16_t); // Buffer size in elements
189 
190  // Create buffer pool
191  BufferPool pool(buffer_pool_size, buffer_size);
192 
193  // Define ring buffers for each step
194  RingBuffer<DataStruct> kernelBufferQueue(max_queue_size);
195  RingBuffer<DataStruct> checkDataQueue(max_queue_size);
196  RingBuffer<DataStruct> formEventsQueue(max_queue_size);
197  RingBuffer<DataStruct> rawBufferQueue(max_queue_size);
198  RingBuffer<DataStruct> zsBufferQueue(max_queue_size);
199  RingBuffer<DataStruct> mwdBufferQueue(max_queue_size);
200 
201  // Launch threads for data processing
202  std::thread kernelThread(
203  kernelBufferThread, std::ref(kernelBufferQueue), std::ref(pool));
204  std::thread checkData(checkDataThread,
205  std::ref(kernelBufferQueue),
206  std::ref(checkDataQueue),
207  std::ref(pool));
208  std::thread formEvents(formEventsThread,
209  std::ref(checkDataQueue),
210  std::ref(formEventsQueue),
211  std::ref(pool));
212  std::thread rawProcess(processAlgorithmThread,
213  std::ref(formEventsQueue),
214  std::ref(rawBufferQueue),
215  std::ref(pool));
216  std::thread zsProcess(processAlgorithmThread,
217  std::ref(rawBufferQueue),
218  std::ref(zsBufferQueue),
219  std::ref(pool));
220  std::thread mwdProcess(processAlgorithmThread,
221  std::ref(zsBufferQueue),
222  std::ref(mwdBufferQueue),
223  std::ref(pool));
224 
225  // Simulate running for a limited time
226  std::this_thread::sleep_for(std::chrono::seconds(10));
227 
228  // Shutdown all queues
229  kernelBufferQueue.shutdown();
230  checkDataQueue.shutdown();
231  formEventsQueue.shutdown();
232  rawBufferQueue.shutdown();
233  zsBufferQueue.shutdown();
234  mwdBufferQueue.shutdown();
235 
236  // Join threads
237  kernelThread.join();
238  checkData.join();
239  formEvents.join();
240  rawProcess.join();
241  zsProcess.join();
242  mwdProcess.join();
243 
244  return 0;
245 }
Definition: data.hh:4