6 #include <condition_variable>
31 std::vector<int16_t>
data;
41 :
data(buffer_size), size(0)
43 std::fill(std::begin(metadata), std::end(metadata),
'\0');
52 std::vector<std::shared_ptr<T>> buffer;
56 std::atomic<size_t> size = 0;
58 std::condition_variable cv_push;
59 std::condition_variable cv_pop;
63 explicit RingBuffer(
size_t max_size) : capacity(max_size), buffer(max_size) {}
67 std::cout <<
"RingBuffer destructor called.\n";
70 void push(
const std::shared_ptr<T>& item)
72 std::unique_lock<std::mutex> lock(mutex);
73 cv_push.wait(lock, [
this]() {
return stop || size < capacity; });
77 head = (head + 1) % capacity;
82 std::shared_ptr<T> pop()
84 std::unique_lock<std::mutex> lock(mutex);
85 cv_pop.wait(lock, [
this]() {
return stop || size > 0; });
88 auto item = buffer[tail];
89 tail = (tail + 1) % capacity;
98 std::lock_guard<std::mutex> lock(mutex);
101 cv_push.notify_all();
107 std::lock_guard<std::mutex> lock(mutex);
116 std::vector<std::shared_ptr<DataStruct>> pool;
120 BufferPool(
size_t pool_size,
size_t buffer_size)
122 for(
size_t i = 0; i < pool_size; ++i)
124 pool.emplace_back(std::make_shared<DataStruct>(buffer_size));
130 std::cout <<
"BufferPool destructor called.\n";
133 std::shared_ptr<DataStruct> acquire()
135 std::lock_guard<std::mutex> lock(mutex);
138 return std::make_shared<DataStruct>(60 * 1024 * 1024 /
sizeof(int16_t));
142 auto buffer = pool.back();
148 void release(std::shared_ptr<DataStruct> buffer)
150 std::lock_guard<std::mutex> lock(mutex);
151 pool.push_back(buffer);
156 std::atomic<uint64_t> globalTotalBytes{0};
157 const uint64_t targetTotalBytes = 1L * 1024 * 1024 * 1024;
158 std::atomic<bool> kernelDone{
false};
159 std::vector<std::atomic<bool>> threadDone(
161 std::vector<int16_t> dummy_data;
162 const size_t buffer_size = 60 * 1024 * 1024 /
sizeof(int16_t);
163 auto dummy_struct = std::make_shared<DataStruct>(buffer_size);
166 void pin_thread_to_core(
size_t core_id)
170 CPU_SET(core_id, &cpuset);
171 pthread_setaffinity_np(pthread_self(),
sizeof(cpu_set_t), &cpuset);
179 std::atomic<uint64_t>& totalBytesProcessed,
181 std::function<
void(std::shared_ptr<DataStruct>&)> operation,
184 pin_thread_to_core(core_id);
185 auto start_time = std::chrono::high_resolution_clock::now();
187 while(!kernelDone || !inputQueue.is_empty() || !threadDone[threadIndex - 1].load())
189 auto buffer = inputQueue.pop();
195 strncat(buffer->metadata,
197 sizeof(buffer->metadata) - strlen(buffer->metadata) - 1);
199 uint64_t bytesProcessed = buffer->size *
sizeof(int16_t);
200 totalBytesProcessed += bytesProcessed;
202 outputQueue.push(buffer);
207 threadDone[threadIndex].store(
true);
209 auto end_time = std::chrono::high_resolution_clock::now();
210 std::chrono::duration<double> elapsed = end_time - start_time;
211 double gbits = totalBytesProcessed.load() * 8.0 / 1e9;
212 std::cout <<
"[" << stage <<
"] Processing speed: " << std::fixed
213 << std::setprecision(2) << (gbits / elapsed.count()) <<
" Gbit/s\n";
221 pin_thread_to_core(core_id);
222 while(globalTotalBytes < targetTotalBytes)
224 auto buffer = pool.acquire();
227 buffer->size = buffer_size;
228 buffer->data = dummy_data;
230 strncpy(buffer->metadata,
231 "Generated by kernelBufferWorker",
232 sizeof(buffer->metadata));
234 uint64_t bytesGenerated = buffer->size *
sizeof(int16_t);
235 globalTotalBytes += bytesGenerated;
237 outputQueue.push(buffer);
242 threadDone[0].store(
true);
251 pin_thread_to_core(core_id);
252 while(!kernelDone || !inputQueue.is_empty() || !threadDone[threadIndex - 1].load())
254 auto buffer = inputQueue.pop();
257 pool.release(buffer);
261 threadDone[threadIndex].store(
true);
263 std::cout <<
"[Consumer] Final buffer emptied.\n";
268 const size_t max_queue_size = 20;
269 const size_t buffer_pool_size = 30;
272 BufferPool pool(buffer_pool_size, buffer_size);
274 std::vector<int16_t> dummy_data;
275 dummy_data.resize(buffer_size);
276 std::fill(dummy_data.begin(), dummy_data.end(), rand() % 100);
287 std::atomic<uint64_t> totalBytesKernel{0};
288 std::atomic<uint64_t> totalBytesCheck{0};
289 std::atomic<uint64_t> totalBytesForm{0};
290 std::atomic<uint64_t> totalBytesRaw{0};
291 std::atomic<uint64_t> totalBytesZS{0};
292 std::atomic<uint64_t> totalBytesMWD{0};
297 for(
auto& flag : threadDone)
303 auto operationExample = [](std::shared_ptr<DataStruct>& buffer) {
304 for(
size_t i = 0; i < buffer->size; ++i)
306 buffer->data[i] += 1;
330 std::thread kernelThread(
331 kernelBufferWorker, std::ref(kernelBufferQueue), std::ref(pool), 0);
332 std::thread checkDataThread(workerThread,
333 std::ref(kernelBufferQueue),
334 std::ref(checkDataQueue),
337 std::ref(totalBytesCheck),
341 std::thread formEventsThread(workerThread,
342 std::ref(checkDataQueue),
343 std::ref(formEventsQueue),
346 std::ref(totalBytesForm),
350 std::thread rawProcessThread(workerThread,
351 std::ref(formEventsQueue),
352 std::ref(rawBufferQueue),
354 " | Algorithm Processed",
355 std::ref(totalBytesRaw),
359 std::thread zsProcessThread(workerThread,
360 std::ref(rawBufferQueue),
361 std::ref(zsBufferQueue),
364 std::ref(totalBytesZS),
368 std::thread mwdProcessThread(workerThread,
369 std::ref(zsBufferQueue),
370 std::ref(mwdBufferQueue),
373 std::ref(totalBytesMWD),
377 std::thread finalConsumerThread(
378 consumerThread, std::ref(mwdBufferQueue), std::ref(pool), 6, 6);
382 checkDataThread.join();
383 formEventsThread.join();
384 rawProcessThread.join();
385 zsProcessThread.join();
386 mwdProcessThread.join();
387 finalConsumerThread.join();
390 kernelBufferQueue.shutdown();
391 checkDataQueue.shutdown();
392 formEventsQueue.shutdown();
393 rawBufferQueue.shutdown();
394 zsBufferQueue.shutdown();
395 mwdBufferQueue.shutdown();