2 #include <condition_variable>
14 int16_t
data[60 * 1024 * 1024 /
sizeof(int16_t)];
24 std::queue<std::shared_ptr<T>>
queue;
26 std::condition_variable cv;
33 void push(
const std::shared_ptr<T>& item)
36 std::unique_lock<std::mutex> lock(mutex);
37 cv.wait(lock, [
this]() {
return stop ||
queue.size() < max_size; });
45 std::shared_ptr<T> pop()
47 std::unique_lock<std::mutex> lock(mutex);
48 cv.wait(lock, [
this]() {
return stop || !
queue.empty(); });
53 auto item =
queue.front();
62 std::lock_guard<std::mutex> lock(mutex);
73 auto buffer = std::make_shared<DataStruct>();
74 buffer->size = (rand() % (60 * 1024 * 1024 /
sizeof(int16_t))) +
76 std::fill(buffer->data,
77 buffer->data + buffer->size,
79 buffer->metadata =
"Generated by kernelBufferThread";
80 bufferQueue.push(buffer);
81 std::this_thread::sleep_for(std::chrono::milliseconds(10));
90 auto buffer = inputQueue.pop();
94 buffer->metadata +=
" | Checked";
95 outputQueue.push(buffer);
104 auto buffer = inputQueue.pop();
108 buffer->metadata +=
" | Events Formed";
109 outputQueue.push(buffer);
118 auto buffer = inputQueue.pop();
122 buffer->metadata +=
" | Algorithm Processed";
123 outputQueue.push(buffer);
129 const size_t max_queue_size = 10;
140 std::thread kernelThread(kernelBufferThread, std::ref(kernelBufferQueue));
142 checkDataThread, std::ref(kernelBufferQueue), std::ref(checkDataQueue));
144 formEventsThread, std::ref(checkDataQueue), std::ref(formEventsQueue));
145 std::thread rawProcess(
146 processAlgorithmThread, std::ref(formEventsQueue), std::ref(rawBufferQueue));
147 std::thread zsProcess(
148 processAlgorithmThread, std::ref(rawBufferQueue), std::ref(zsBufferQueue));
149 std::thread mwdProcess(
150 processAlgorithmThread, std::ref(zsBufferQueue), std::ref(mwdBufferQueue));
153 std::this_thread::sleep_for(std::chrono::seconds(10));
156 kernelBufferQueue.shutdown();
157 checkDataQueue.shutdown();
158 formEventsQueue.shutdown();
159 rawBufferQueue.shutdown();
160 zsBufferQueue.shutdown();
161 mwdBufferQueue.shutdown();