7 #include <condition_variable>
21 std::vector<int16_t>
data;
27 std::fill(std::begin(metadata),
38 std::vector<std::shared_ptr<T>> buffer;
42 std::atomic<size_t> size = 0;
44 std::condition_variable cv_push;
45 std::condition_variable cv_pop;
50 explicit RingBuffer(
size_t max_size) : capacity(max_size), buffer(max_size) {}
53 ~
RingBuffer() { std::cout <<
"RingBuffer destructor called.\n"; }
56 void push(
const std::shared_ptr<T>& item)
58 std::unique_lock<std::mutex> lock(mutex);
59 cv_push.wait(lock, [
this]() {
60 return stop || size < capacity;
65 head = (head + 1) % capacity;
71 std::shared_ptr<T> pop()
73 std::unique_lock<std::mutex> lock(mutex);
75 [
this]() {
return stop || size > 0; });
78 auto item = buffer[tail];
79 tail = (tail + 1) % capacity;
89 std::lock_guard<std::mutex> lock(mutex);
98 std::lock_guard<std::mutex> lock(mutex);
107 std::vector<std::shared_ptr<DataStruct>> pool;
112 BufferPool(
size_t pool_size,
size_t buffer_size)
114 for(
size_t i = 0; i < pool_size; ++i)
117 std::make_shared<DataStruct>(buffer_size));
122 ~
BufferPool() { std::cout <<
"BufferPool destructor called.\n"; }
125 std::shared_ptr<DataStruct> acquire()
127 std::lock_guard<std::mutex> lock(mutex);
130 return std::make_shared<DataStruct>(60 * 1024 * 1024 /
135 auto buffer = pool.back();
142 void release(std::shared_ptr<DataStruct> buffer)
144 std::lock_guard<std::mutex> lock(mutex);
145 pool.push_back(buffer);
150 std::atomic<uint64_t> globalTotalBytes{0};
151 const uint64_t targetTotalBytes = 1000L * 1024 * 1024 * 1024;
152 std::atomic<bool> kernelDone{
false};
155 void pin_thread_to_core(
size_t core_id)
159 CPU_SET(core_id, &cpuset);
160 pthread_setaffinity_np(
161 pthread_self(),
sizeof(cpu_set_t), &cpuset);
169 std::atomic<uint64_t>& totalBytesProcessed,
172 pin_thread_to_core(core_id);
173 auto start_time = std::chrono::high_resolution_clock::now();
174 while(!kernelDone || !inputQueue.is_empty())
176 auto buffer = inputQueue.pop();
179 strncat(buffer->metadata,
181 sizeof(buffer->metadata) - strlen(buffer->metadata) -
183 uint64_t bytesProcessed =
184 buffer->size *
sizeof(int16_t);
185 totalBytesProcessed += bytesProcessed;
186 outputQueue.push(buffer);
189 auto end_time = std::chrono::high_resolution_clock::now();
190 std::chrono::duration<double> elapsed =
191 end_time - start_time;
192 double gbits = totalBytesProcessed.load() * 8.0 / 1e9;
193 std::cout <<
"[" << stage <<
"] Processing speed: " << std::fixed
194 << std::setprecision(2) << (gbits / elapsed.count())
203 pin_thread_to_core(core_id);
204 while(globalTotalBytes < targetTotalBytes)
206 auto buffer = pool.acquire();
207 buffer->size = (rand() % (60 * 1024 * 1024 /
sizeof(int16_t))) +
209 memset(buffer->data.data(),
211 buffer->size *
sizeof(int16_t));
212 strncpy(buffer->metadata,
213 "Generated by kernelBufferWorker",
214 sizeof(buffer->metadata));
215 uint64_t bytesGenerated =
216 buffer->size *
sizeof(int16_t);
217 globalTotalBytes += bytesGenerated;
218 outputQueue.push(buffer);
226 pin_thread_to_core(core_id);
227 while(!kernelDone || !inputQueue.is_empty())
229 auto buffer = inputQueue.pop();
232 pool.release(buffer);
235 std::cout <<
"[Consumer] Final buffer emptied.\n";
240 const size_t max_queue_size = 20;
241 const size_t buffer_pool_size = 30;
242 const size_t buffer_size =
243 60 * 1024 * 1024 /
sizeof(int16_t);
245 BufferPool pool(buffer_pool_size, buffer_size);
254 std::atomic<uint64_t> totalBytesKernel{
256 std::atomic<uint64_t> totalBytesCheck{
258 std::atomic<uint64_t> totalBytesForm{
260 std::atomic<uint64_t> totalBytesRaw{
262 std::atomic<uint64_t> totalBytesZS{
264 std::atomic<uint64_t> totalBytesMWD{
267 std::thread kernelThread(kernelBufferWorker,
268 std::ref(kernelBufferQueue),
271 std::thread checkDataThread(workerThread,
272 std::ref(kernelBufferQueue),
273 std::ref(checkDataQueue),
276 std::ref(totalBytesCheck),
278 std::thread formEventsThread(workerThread,
279 std::ref(checkDataQueue),
280 std::ref(formEventsQueue),
283 std::ref(totalBytesForm),
285 std::thread rawProcessThread(workerThread,
286 std::ref(formEventsQueue),
287 std::ref(rawBufferQueue),
289 " | Algorithm Processed",
290 std::ref(totalBytesRaw),
292 std::thread zsProcessThread(workerThread,
293 std::ref(rawBufferQueue),
294 std::ref(zsBufferQueue),
297 std::ref(totalBytesZS),
299 std::thread mwdProcessThread(workerThread,
300 std::ref(zsBufferQueue),
301 std::ref(mwdBufferQueue),
304 std::ref(totalBytesMWD),
306 std::thread finalConsumerThread(consumerThread,
307 std::ref(mwdBufferQueue),
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();
329 <<
"[Main] All threads completed successfully.\n";