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) {}
48 void push(
const std::shared_ptr<T>& item)
50 std::unique_lock<std::mutex> lock(mutex);
51 cv_push.wait(lock, [
this]() {
return stop || size < capacity; });
55 head = (head + 1) % capacity;
60 std::shared_ptr<T> pop()
62 std::unique_lock<std::mutex> lock(mutex);
63 cv_pop.wait(lock, [
this]() {
return stop || size > 0; });
66 auto item = buffer[tail];
67 tail = (tail + 1) % capacity;
76 std::lock_guard<std::mutex> lock(mutex);
85 std::lock_guard<std::mutex> lock(mutex);
94 std::vector<std::shared_ptr<DataStruct>> pool;
98 BufferPool(
size_t pool_size,
size_t buffer_size)
100 for(
size_t i = 0; i < pool_size; ++i)
102 pool.emplace_back(std::make_shared<DataStruct>(buffer_size));
106 std::shared_ptr<DataStruct> acquire()
108 std::lock_guard<std::mutex> lock(mutex);
111 return std::make_shared<DataStruct>(60 * 1024 * 1024 /
sizeof(int16_t));
115 auto buffer = pool.back();
121 void release(std::shared_ptr<DataStruct> buffer)
123 std::lock_guard<std::mutex> lock(mutex);
124 pool.push_back(buffer);
129 std::atomic<uint64_t> globalTotalBytes{0};
130 const uint64_t targetTotalBytes = 10L * 1024 * 1024 * 1024;
131 std::atomic<bool> kernelDone{
false};
134 void pin_thread_to_core(
size_t core_id)
138 CPU_SET(core_id, &cpuset);
139 pthread_setaffinity_np(pthread_self(),
sizeof(cpu_set_t), &cpuset);
147 std::atomic<uint64_t>& totalBytesProcessed,
150 pin_thread_to_core(core_id);
151 auto start_time = std::chrono::high_resolution_clock::now();
152 while(!kernelDone || !inputQueue.is_empty())
154 std::vector<std::shared_ptr<DataStruct>> batch;
156 auto buffer = inputQueue.pop();
159 batch.push_back(buffer);
162 for(
auto& buffer : batch)
165 strncat(buffer->metadata,
167 sizeof(buffer->metadata) - strlen(buffer->metadata) - 1);
170 uint64_t bytesProcessed = buffer->size *
sizeof(int16_t);
171 totalBytesProcessed += bytesProcessed;
174 outputQueue.push(buffer);
177 auto end_time = std::chrono::high_resolution_clock::now();
178 std::chrono::duration<double> elapsed = end_time - start_time;
179 double gbits = totalBytesProcessed.load() * 8.0 / 1e9;
180 std::cout <<
"[" << stage <<
"] Processing speed: " << std::fixed
181 << std::setprecision(2) << (gbits / elapsed.count()) <<
" Gbit/s\n";
189 pin_thread_to_core(core_id);
190 while(globalTotalBytes < targetTotalBytes)
192 auto buffer = pool.acquire();
193 buffer->size = (rand() % (60 * 1024 * 1024 /
sizeof(int16_t))) +
195 memset(buffer->data.data(), rand() % 100, buffer->size *
sizeof(int16_t));
196 strncpy(buffer->metadata,
197 "Generated by kernelBufferWorker",
198 sizeof(buffer->metadata));
200 uint64_t bytesGenerated = buffer->size *
sizeof(int16_t);
201 globalTotalBytes += bytesGenerated;
203 outputQueue.push(buffer);
211 pin_thread_to_core(core_id);
212 while(!kernelDone || !inputQueue.is_empty())
214 auto buffer = inputQueue.pop();
219 pool.release(buffer);
221 std::cout <<
"[Consumer] Final buffer emptied.\n";
226 const size_t max_queue_size = 20;
227 const size_t buffer_pool_size = 30;
228 const size_t buffer_size =
229 60 * 1024 * 1024 /
sizeof(int16_t);
232 BufferPool pool(buffer_pool_size, buffer_size);
243 std::atomic<uint64_t> totalBytesKernel{0};
244 std::atomic<uint64_t> totalBytesCheck{0};
245 std::atomic<uint64_t> totalBytesForm{0};
246 std::atomic<uint64_t> totalBytesRaw{0};
247 std::atomic<uint64_t> totalBytesZS{0};
248 std::atomic<uint64_t> totalBytesMWD{0};
251 std::thread kernelThread(
252 kernelBufferWorker, std::ref(kernelBufferQueue), std::ref(pool), 0);
253 std::thread checkDataThread(workerThread,
254 std::ref(kernelBufferQueue),
255 std::ref(checkDataQueue),
258 std::ref(totalBytesCheck),
260 std::thread formEventsThread(workerThread,
261 std::ref(checkDataQueue),
262 std::ref(formEventsQueue),
265 std::ref(totalBytesForm),
267 std::thread rawProcessThread(workerThread,
268 std::ref(formEventsQueue),
269 std::ref(rawBufferQueue),
271 " | Algorithm Processed",
272 std::ref(totalBytesRaw),
274 std::thread zsProcessThread(workerThread,
275 std::ref(rawBufferQueue),
276 std::ref(zsBufferQueue),
279 std::ref(totalBytesZS),
281 std::thread mwdProcessThread(workerThread,
282 std::ref(zsBufferQueue),
283 std::ref(mwdBufferQueue),
286 std::ref(totalBytesMWD),
288 std::thread finalConsumerThread(
289 consumerThread, std::ref(mwdBufferQueue), std::ref(pool), 6);
293 checkDataThread.join();
294 formEventsThread.join();
295 rawProcessThread.join();
296 zsProcessThread.join();
297 mwdProcessThread.join();
298 finalConsumerThread.join();
301 kernelBufferQueue.shutdown();
302 checkDataQueue.shutdown();
303 formEventsQueue.shutdown();
304 rawBufferQueue.shutdown();
305 zsBufferQueue.shutdown();
306 mwdBufferQueue.shutdown();