6 #include <condition_variable>
20 std::vector<int16_t>
data;
21 std::vector<int16_t> data2;
22 std::vector<int16_t> data3;
30 :
data(buffer_size), data2(), data3(), size(0), size2(0), size3(0)
32 std::fill(std::begin(metadata), std::end(metadata),
'\0');
41 std::vector<std::shared_ptr<T>> buffer;
45 std::atomic<size_t> size = 0;
47 std::condition_variable cv_push;
48 std::condition_variable cv_pop;
52 explicit RingBuffer(
size_t max_size) : capacity(max_size), buffer(max_size) {}
56 std::cout <<
"RingBuffer destructor called.\n";
59 void push(
const std::shared_ptr<T>& item)
61 std::unique_lock<std::mutex> lock(mutex);
62 cv_push.wait(lock, [
this]() {
return stop || size < capacity; });
66 head = (head + 1) % capacity;
71 std::shared_ptr<T> pop()
73 std::unique_lock<std::mutex> lock(mutex);
74 cv_pop.wait(lock, [
this]() {
return stop || size > 0; });
77 auto item = buffer[tail];
78 tail = (tail + 1) % capacity;
87 std::lock_guard<std::mutex> lock(mutex);
96 std::lock_guard<std::mutex> lock(mutex);
105 std::vector<std::shared_ptr<DataStruct>> pool;
109 BufferPool(
size_t pool_size,
size_t buffer_size)
111 for(
size_t i = 0; i < pool_size; ++i)
113 pool.emplace_back(std::make_shared<DataStruct>(buffer_size));
119 std::cout <<
"BufferPool destructor called.\n";
133 std::shared_ptr<DataStruct> acquire()
135 std::lock_guard<std::mutex> lock(mutex);
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 = 10L * 1024 * 1024 * 1024;
158 std::atomic<bool> kernelDone{
false};
159 std::vector<std::atomic<bool>> threadDone(
163 void pin_thread_to_core(
size_t core_id)
167 CPU_SET(core_id, &cpuset);
168 pthread_setaffinity_np(pthread_self(),
sizeof(cpu_set_t), &cpuset);
176 std::atomic<uint64_t>& totalBytesProcessed,
178 std::function<
void(std::shared_ptr<DataStruct>&)> operation,
181 pin_thread_to_core(core_id);
182 auto start_time = std::chrono::high_resolution_clock::now();
184 while(!kernelDone || !inputQueue.is_empty() || !threadDone[threadIndex - 1].load())
186 auto buffer = inputQueue.pop();
192 strncat(buffer->metadata,
194 sizeof(buffer->metadata) - strlen(buffer->metadata) - 1);
196 uint64_t bytesProcessed = buffer->size *
sizeof(int16_t);
197 totalBytesProcessed += bytesProcessed;
199 outputQueue.push(buffer);
204 threadDone[threadIndex].store(
true);
206 auto end_time = std::chrono::high_resolution_clock::now();
207 std::chrono::duration<double> elapsed = end_time - start_time;
208 double gbits = totalBytesProcessed.load() * 8.0 / 1e9;
209 std::cout <<
"[" << stage <<
"] Processing speed: " << std::fixed
210 << std::setprecision(2) << (gbits / elapsed.count()) <<
" Gbit/s\n";
218 pin_thread_to_core(core_id);
220 while(globalTotalBytes < targetTotalBytes)
222 auto buffer = pool.acquire();
225 std::this_thread::yield();
228 buffer->size = (rand() % (60 * 1024 * 1024 /
sizeof(int16_t))) +
230 memset(buffer->data.data(), rand() % 100, buffer->size *
sizeof(int16_t));
231 strncpy(buffer->metadata,
232 "Generated by kernelBufferWorker",
233 sizeof(buffer->metadata));
235 uint64_t bytesGenerated = buffer->size *
sizeof(int16_t);
236 globalTotalBytes += bytesGenerated;
238 outputQueue.push(buffer);
243 threadDone[0].store(
true);
252 pin_thread_to_core(core_id);
253 while(!kernelDone || !inputQueue.is_empty() || !threadDone[threadIndex - 1].load())
255 auto buffer = inputQueue.pop();
258 pool.release(buffer);
262 threadDone[threadIndex].store(
true);
264 std::cout <<
"[Consumer] Final buffer emptied.\n";
269 const size_t max_queue_size = 20;
270 const size_t buffer_pool_size = 30;
271 const size_t buffer_size =
272 60 * 1024 * 1024 /
sizeof(int16_t);
275 BufferPool pool(buffer_pool_size, buffer_size);
286 std::atomic<uint64_t> totalBytesKernel{0};
287 std::atomic<uint64_t> totalBytesCheck{0};
288 std::atomic<uint64_t> totalBytesForm{0};
289 std::atomic<uint64_t> totalBytesRaw{0};
290 std::atomic<uint64_t> totalBytesZS{0};
291 std::atomic<uint64_t> totalBytesMWD{0};
296 for(
auto& flag : threadDone)
302 auto operationExample = [](std::shared_ptr<DataStruct>& buffer) {
303 for(
size_t i = 0; i < buffer->size; ++i)
305 buffer->data[i] += 1;
310 auto zsOperation = [](std::shared_ptr<DataStruct>& buffer) {
311 buffer->data2.resize(buffer->data.size());
312 buffer->size2 = buffer->size;
313 for(
size_t i = 0; i < buffer->size; ++i)
315 buffer->data2[i] = buffer->data[i] * 2;
320 auto mwdOperation = [](std::shared_ptr<DataStruct>& buffer) {
321 buffer->data3.resize(buffer->data2.size());
322 buffer->size3 = buffer->size2;
323 for(
size_t i = 0; i < buffer->size2; ++i)
325 double temp =
static_cast<double>(buffer->data2[i]) / 2.0;
326 buffer->data3[i] =
static_cast<int16_t
>(temp);
331 std::thread kernelThread(
332 kernelBufferWorker, std::ref(kernelBufferQueue), std::ref(pool), 0);
333 std::thread checkDataThread(workerThread,
334 std::ref(kernelBufferQueue),
335 std::ref(checkDataQueue),
338 std::ref(totalBytesCheck),
342 std::thread formEventsThread(workerThread,
343 std::ref(checkDataQueue),
344 std::ref(formEventsQueue),
347 std::ref(totalBytesForm),
351 std::thread rawProcessThread(workerThread,
352 std::ref(formEventsQueue),
353 std::ref(rawBufferQueue),
355 " | Algorithm Processed",
356 std::ref(totalBytesRaw),
360 std::thread zsProcessThread(workerThread,
361 std::ref(rawBufferQueue),
362 std::ref(zsBufferQueue),
365 std::ref(totalBytesZS),
369 std::thread mwdProcessThread(workerThread,
370 std::ref(zsBufferQueue),
371 std::ref(mwdBufferQueue),
374 std::ref(totalBytesMWD),
378 std::thread finalConsumerThread(
379 consumerThread, std::ref(mwdBufferQueue), std::ref(pool), 6, 6);
383 checkDataThread.join();
384 formEventsThread.join();
385 rawProcessThread.join();
386 zsProcessThread.join();
387 mwdProcessThread.join();
388 finalConsumerThread.join();
391 kernelBufferQueue.shutdown();
392 checkDataQueue.shutdown();
393 formEventsQueue.shutdown();
394 rawBufferQueue.shutdown();
395 zsBufferQueue.shutdown();
396 mwdBufferQueue.shutdown();