8 #include <condition_variable>
17 #include <unordered_map>
23 std::vector<int16_t>
data;
29 std::fill(std::begin(metadata), std::end(metadata),
'\0');
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;
49 explicit RingBuffer(
size_t max_size) : capacity(max_size), buffer(max_size) {}
51 void push(
const std::shared_ptr<T>& item)
53 std::unique_lock<std::mutex> lock(mutex);
54 cv_push.wait(lock, [
this]() {
return stop || size < capacity; });
58 head = (head + 1) % capacity;
63 std::shared_ptr<T> pop()
65 std::unique_lock<std::mutex> lock(mutex);
66 cv_pop.wait(lock, [
this]() {
return stop || size > 0; });
69 auto item = buffer[tail];
70 tail = (tail + 1) % capacity;
79 std::lock_guard<std::mutex> lock(mutex);
88 std::lock_guard<std::mutex> lock(mutex);
97 std::vector<std::shared_ptr<DataStruct>> pool;
101 BufferPool(
size_t pool_size,
size_t buffer_size)
103 for(
size_t i = 0; i < pool_size; ++i)
105 pool.emplace_back(std::make_shared<DataStruct>(buffer_size));
109 std::shared_ptr<DataStruct> acquire()
111 std::lock_guard<std::mutex> lock(mutex);
114 return std::make_shared<DataStruct>(60 * 1024 * 1024 /
sizeof(int16_t));
118 auto buffer = pool.back();
124 void release(std::shared_ptr<DataStruct> buffer)
126 std::lock_guard<std::mutex> lock(mutex);
127 pool.push_back(buffer);
132 std::atomic<uint64_t> globalTotalBytes{0};
133 const uint64_t targetTotalBytes = 10L * 1024 * 1024 * 1024;
134 std::atomic<bool> kernelDone{
false};
135 std::unordered_map<size_t, std::vector<int16_t>>
137 std::mutex dataMapMutex;
138 std::vector<std::shared_ptr<DataStruct>> finalBuffer;
139 std::mutex finalBufferMutex;
142 void generateData(std::shared_ptr<DataStruct>& buffer)
144 std::lock_guard<std::mutex> lock(dataMapMutex);
145 for(
size_t i = 0; i < buffer->size; ++i)
147 buffer->data[i] =
static_cast<int16_t
>(i % 32768);
149 preGeneratedData[
reinterpret_cast<size_t>(buffer.get())] = buffer->data;
153 bool validateFinalBuffer()
155 std::lock_guard<std::mutex> lock(dataMapMutex);
156 for(
const auto& buffer : finalBuffer)
158 auto it = preGeneratedData.find(
reinterpret_cast<size_t>(buffer.get()));
159 if(it == preGeneratedData.end() ||
160 !std::equal(buffer->data.begin(), buffer->data.end(), it->second.begin()))
169 void pin_thread_to_core(
size_t core_id)
173 CPU_SET(core_id, &cpuset);
174 pthread_setaffinity_np(pthread_self(),
sizeof(cpu_set_t), &cpuset);
182 std::atomic<uint64_t>& totalBytesProcessed,
185 pin_thread_to_core(core_id);
186 auto start_time = std::chrono::high_resolution_clock::now();
187 while(!kernelDone || !inputQueue.is_empty())
189 auto buffer = inputQueue.pop();
194 strncat(buffer->metadata,
196 sizeof(buffer->metadata) - strlen(buffer->metadata) - 1);
199 uint64_t bytesProcessed = buffer->size *
sizeof(int16_t);
200 totalBytesProcessed += bytesProcessed;
203 outputQueue.push(buffer);
205 auto end_time = std::chrono::high_resolution_clock::now();
206 std::chrono::duration<double> elapsed = end_time - start_time;
207 double gbits = totalBytesProcessed.load() * 8.0 / 1e9;
208 std::cout <<
"[" << stage <<
"] Processing speed: " << std::fixed
209 << std::setprecision(2) << (gbits / elapsed.count()) <<
" Gbit/s\n";
217 pin_thread_to_core(core_id);
218 while(globalTotalBytes < targetTotalBytes)
220 auto buffer = pool.acquire();
221 buffer->size = (rand() % (60 * 1024 * 1024 /
sizeof(int16_t))) +
223 generateData(buffer);
224 strncpy(buffer->metadata,
225 "Generated by kernelBufferWorker",
226 sizeof(buffer->metadata));
228 uint64_t bytesGenerated = buffer->size *
sizeof(int16_t);
229 globalTotalBytes += bytesGenerated;
231 outputQueue.push(buffer);
239 pin_thread_to_core(core_id);
240 while(!kernelDone || !inputQueue.is_empty())
242 auto buffer = inputQueue.pop();
251 pool.release(buffer);
253 std::cout <<
"[Consumer] Final buffer collected.\n";
258 const size_t max_queue_size = 20;
259 const size_t buffer_pool_size = 30;
260 const size_t buffer_size =
261 60 * 1024 * 1024 /
sizeof(int16_t);
264 BufferPool pool(buffer_pool_size, buffer_size);
275 std::atomic<uint64_t> totalBytesKernel{0};
276 std::atomic<uint64_t> totalBytesCheck{0};
277 std::atomic<uint64_t> totalBytesForm{0};
278 std::atomic<uint64_t> totalBytesRaw{0};
279 std::atomic<uint64_t> totalBytesZS{0};
280 std::atomic<uint64_t> totalBytesMWD{0};
283 std::thread kernelThread(
284 kernelBufferWorker, std::ref(kernelBufferQueue), std::ref(pool), 0);
285 std::thread checkDataThread(workerThread,
286 std::ref(kernelBufferQueue),
287 std::ref(checkDataQueue),
290 std::ref(totalBytesCheck),
292 std::thread formEventsThread(workerThread,
293 std::ref(checkDataQueue),
294 std::ref(formEventsQueue),
297 std::ref(totalBytesForm),
299 std::thread rawProcessThread(workerThread,
300 std::ref(formEventsQueue),
301 std::ref(rawBufferQueue),
303 " | Algorithm Processed",
304 std::ref(totalBytesRaw),
306 std::thread zsProcessThread(workerThread,
307 std::ref(rawBufferQueue),
308 std::ref(zsBufferQueue),
311 std::ref(totalBytesZS),
313 std::thread mwdProcessThread(workerThread,
314 std::ref(zsBufferQueue),
315 std::ref(mwdBufferQueue),
318 std::ref(totalBytesMWD),
320 std::thread finalConsumerThread(
321 consumerThread, std::ref(mwdBufferQueue), std::ref(pool), 6);
325 checkDataThread.join();
326 formEventsThread.join();
327 rawProcessThread.join();
328 zsProcessThread.join();
329 mwdProcessThread.join();
330 finalConsumerThread.join();
333 bool success = validateFinalBuffer();
336 std::cout <<
"[Validation] All data validated successfully.\n";
340 std::cerr <<
"[Validation] Data validation failed.\n";
344 kernelBufferQueue.shutdown();
345 checkDataQueue.shutdown();
346 formEventsQueue.shutdown();
347 rawBufferQueue.shutdown();
348 zsBufferQueue.shutdown();
349 mwdBufferQueue.shutdown();