6 #include <condition_variable>
20 std::vector<int16_t>
data;
26 std::fill(std::begin(metadata), std::end(metadata),
'\0');
35 std::vector<std::shared_ptr<T>> buffer;
39 std::atomic<size_t> size = 0;
41 std::condition_variable cv_push;
42 std::condition_variable cv_pop;
46 explicit RingBuffer(
size_t max_size) : capacity(max_size), buffer(max_size) {}
50 std::cout <<
"RingBuffer destructor called.\n";
53 void push(
const std::shared_ptr<T>& item)
55 std::unique_lock<std::mutex> lock(mutex);
56 cv_push.wait(lock, [
this]() {
return stop || size < capacity; });
60 head = (head + 1) % capacity;
65 std::shared_ptr<T> pop()
67 std::unique_lock<std::mutex> lock(mutex);
68 cv_pop.wait(lock, [
this]() {
return stop || size > 0; });
71 auto item = buffer[tail];
72 tail = (tail + 1) % capacity;
81 std::lock_guard<std::mutex> lock(mutex);
90 std::lock_guard<std::mutex> lock(mutex);
99 std::vector<std::shared_ptr<DataStruct>> pool;
103 BufferPool(
size_t pool_size,
size_t buffer_size)
105 for(
size_t i = 0; i < pool_size; ++i)
107 pool.emplace_back(std::make_shared<DataStruct>(buffer_size));
113 std::cout <<
"BufferPool destructor called.\n";
116 std::shared_ptr<DataStruct> acquire()
118 std::lock_guard<std::mutex> lock(mutex);
121 return std::make_shared<DataStruct>(60 * 1024 * 1024 /
sizeof(int16_t));
125 auto buffer = pool.back();
131 void release(std::shared_ptr<DataStruct> buffer)
133 std::lock_guard<std::mutex> lock(mutex);
134 pool.push_back(buffer);
139 std::atomic<uint64_t> globalTotalBytes{0};
140 const uint64_t targetTotalBytes = 10L * 1024 * 1024 * 1024;
141 std::atomic<bool> kernelDone{
false};
144 void pin_thread_to_core(
size_t core_id)
148 CPU_SET(core_id, &cpuset);
149 pthread_setaffinity_np(pthread_self(),
sizeof(cpu_set_t), &cpuset);
157 std::atomic<uint64_t>& totalBytesProcessed,
159 std::function<
void(std::shared_ptr<DataStruct>&)> operation)
161 pin_thread_to_core(core_id);
162 auto start_time = std::chrono::high_resolution_clock::now();
163 while(!kernelDone || !inputQueue.is_empty())
165 auto buffer = inputQueue.pop();
172 strncat(buffer->metadata,
174 sizeof(buffer->metadata) - strlen(buffer->metadata) - 1);
177 uint64_t bytesProcessed = buffer->size *
sizeof(int16_t);
178 totalBytesProcessed += bytesProcessed;
181 outputQueue.push(buffer);
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;
187 std::cout <<
"[" << stage <<
"] Processing speed: " << std::fixed
188 << std::setprecision(2) << (gbits / elapsed.count()) <<
" Gbit/s\n";
196 pin_thread_to_core(core_id);
197 while(globalTotalBytes < targetTotalBytes)
199 auto buffer = pool.acquire();
200 buffer->size = (rand() % (60 * 1024 * 1024 /
sizeof(int16_t))) +
202 memset(buffer->data.data(), rand() % 100, buffer->size *
sizeof(int16_t));
203 strncpy(buffer->metadata,
204 "Generated by kernelBufferWorker",
205 sizeof(buffer->metadata));
207 uint64_t bytesGenerated = buffer->size *
sizeof(int16_t);
208 globalTotalBytes += bytesGenerated;
210 outputQueue.push(buffer);
218 pin_thread_to_core(core_id);
219 while(!kernelDone || !inputQueue.is_empty())
221 auto buffer = inputQueue.pop();
226 pool.release(buffer);
228 std::cout <<
"[Consumer] Final buffer emptied.\n";
233 const size_t max_queue_size = 20;
234 const size_t buffer_pool_size = 30;
235 const size_t buffer_size =
236 60 * 1024 * 1024 /
sizeof(int16_t);
239 BufferPool pool(buffer_pool_size, buffer_size);
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};
258 auto operationExample = [](std::shared_ptr<DataStruct>& buffer) {
259 for(
size_t i = 0; i < buffer->size; ++i)
261 buffer->data[i] += 1;
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),
273 std::ref(totalBytesCheck),
276 std::thread formEventsThread(workerThread,
277 std::ref(checkDataQueue),
278 std::ref(formEventsQueue),
281 std::ref(totalBytesForm),
284 std::thread rawProcessThread(workerThread,
285 std::ref(formEventsQueue),
286 std::ref(rawBufferQueue),
288 " | Algorithm Processed",
289 std::ref(totalBytesRaw),
292 std::thread zsProcessThread(workerThread,
293 std::ref(rawBufferQueue),
294 std::ref(zsBufferQueue),
297 std::ref(totalBytesZS),
300 std::thread mwdProcessThread(workerThread,
301 std::ref(zsBufferQueue),
302 std::ref(mwdBufferQueue),
305 std::ref(totalBytesMWD),
308 std::thread finalConsumerThread(
309 consumerThread, std::ref(mwdBufferQueue), std::ref(pool), 6);
313 checkDataThread.join();
314 formEventsThread.join();
315 rawProcessThread.join();
316 zsProcessThread.join();
317 mwdProcessThread.join();
318 finalConsumerThread.join();
321 kernelBufferQueue.shutdown();
322 checkDataQueue.shutdown();
323 formEventsQueue.shutdown();
324 rawBufferQueue.shutdown();
325 zsBufferQueue.shutdown();
326 mwdBufferQueue.shutdown();