7 #include <condition_variable>
18 #include "ring_buffer.hh"
92 std::vector<std::shared_ptr<DataStruct>> pool;
96 BufferPool(
size_t pool_size,
size_t buffer_size)
98 for(
size_t i = 0; i < pool_size; ++i)
100 pool.emplace_back(std::make_shared<DataStruct>(buffer_size));
106 std::cout <<
"BufferPool destructor called.\n";
109 std::shared_ptr<DataStruct> acquire()
111 std::lock_guard<std::mutex> lock(mutex);
114 return std::make_shared<DataStruct>(60 * 1024 * 1024 /
sizeof(int16_t));
118 auto buffer = pool.back();
124 void release(std::shared_ptr<DataStruct> buffer)
126 std::lock_guard<std::mutex> lock(mutex);
127 pool.push_back(buffer);
132 std::atomic<uint64_t> globalTotalBytes{0};
133 const uint64_t targetTotalBytes = 10L * 1024 * 1024 * 1024;
134 std::atomic<bool> kernelDone{
false};
135 std::vector<std::atomic<bool>> threadDone(
139 void pin_thread_to_core(
size_t core_id)
143 CPU_SET(core_id, &cpuset);
144 pthread_setaffinity_np(pthread_self(),
sizeof(cpu_set_t), &cpuset);
152 std::atomic<uint64_t>& totalBytesProcessed,
154 std::function<
void(std::shared_ptr<DataStruct>&)> operation,
157 pin_thread_to_core(core_id);
158 auto start_time = std::chrono::high_resolution_clock::now();
160 while(!kernelDone || !inputQueue.is_empty() || !threadDone[threadIndex - 1].load())
162 auto buffer = inputQueue.pop();
168 strncat(buffer->metadata,
170 sizeof(buffer->metadata) - strlen(buffer->metadata) - 1);
172 uint64_t bytesProcessed = buffer->size *
sizeof(int16_t);
173 totalBytesProcessed += bytesProcessed;
175 outputQueue.push(buffer);
180 threadDone[threadIndex].store(
true);
182 auto end_time = std::chrono::high_resolution_clock::now();
183 std::chrono::duration<double> elapsed = end_time - start_time;
184 double gbits = totalBytesProcessed.load() * 8.0 / 1e9;
185 std::cout <<
"[" << stage <<
"] Processing speed: " << std::fixed
186 << std::setprecision(2) << (gbits / elapsed.count()) <<
" Gbit/s\n";
194 pin_thread_to_core(core_id);
195 while(globalTotalBytes < targetTotalBytes)
197 auto buffer = pool.acquire();
198 buffer->size = (rand() % (60 * 1024 * 1024 /
sizeof(int16_t))) +
200 memset(buffer->data.data(), rand() % 100, buffer->size *
sizeof(int16_t));
201 strncpy(buffer->metadata,
202 "Generated by kernelBufferWorker",
203 sizeof(buffer->metadata));
205 uint64_t bytesGenerated = buffer->size *
sizeof(int16_t);
206 globalTotalBytes += bytesGenerated;
208 outputQueue.push(buffer);
213 threadDone[0].store(
true);
222 pin_thread_to_core(core_id);
223 while(!kernelDone || !inputQueue.is_empty() || !threadDone[threadIndex - 1].load())
225 auto buffer = inputQueue.pop();
228 pool.release(buffer);
232 threadDone[threadIndex].store(
true);
234 std::cout <<
"[Consumer] Final buffer emptied.\n";
239 const size_t max_queue_size = 20;
240 const size_t buffer_pool_size = 30;
241 const size_t buffer_size =
242 60 * 1024 * 1024 /
sizeof(int16_t);
245 BufferPool pool(buffer_pool_size, buffer_size);
256 std::atomic<uint64_t> totalBytesKernel{0};
257 std::atomic<uint64_t> totalBytesCheck{0};
258 std::atomic<uint64_t> totalBytesForm{0};
259 std::atomic<uint64_t> totalBytesRaw{0};
260 std::atomic<uint64_t> totalBytesZS{0};
261 std::atomic<uint64_t> totalBytesMWD{0};
266 for(
auto& flag : threadDone)
272 auto operationExample = [](std::shared_ptr<DataStruct>& buffer) {
273 for(
size_t i = 0; i < buffer->size; ++i)
275 buffer->data[i] += 1;
280 std::thread kernelThread(
281 kernelBufferWorker, std::ref(kernelBufferQueue), std::ref(pool), 0);
282 std::thread checkDataThread(workerThread,
283 std::ref(kernelBufferQueue),
284 std::ref(checkDataQueue),
287 std::ref(totalBytesCheck),
291 std::thread formEventsThread(workerThread,
292 std::ref(checkDataQueue),
293 std::ref(formEventsQueue),
296 std::ref(totalBytesForm),
300 std::thread rawProcessThread(workerThread,
301 std::ref(formEventsQueue),
302 std::ref(rawBufferQueue),
304 " | Algorithm Processed",
305 std::ref(totalBytesRaw),
309 std::thread zsProcessThread(workerThread,
310 std::ref(rawBufferQueue),
311 std::ref(zsBufferQueue),
314 std::ref(totalBytesZS),
318 std::thread mwdProcessThread(workerThread,
319 std::ref(zsBufferQueue),
320 std::ref(mwdBufferQueue),
323 std::ref(totalBytesMWD),
327 std::thread finalConsumerThread(
328 consumerThread, std::ref(mwdBufferQueue), std::ref(pool), 6, 6);
332 checkDataThread.join();
333 formEventsThread.join();
334 rawProcessThread.join();
335 zsProcessThread.join();
336 mwdProcessThread.join();
337 finalConsumerThread.join();
340 kernelBufferQueue.shutdown();
341 checkDataQueue.shutdown();
342 formEventsQueue.shutdown();
343 rawBufferQueue.shutdown();
344 zsBufferQueue.shutdown();
345 mwdBufferQueue.shutdown();