4 #include <condition_variable>
16 std::vector<int16_t>
data;
22 std::fill(std::begin(metadata), std::end(metadata),
'\0');
31 std::vector<std::shared_ptr<T>> buffer;
35 std::atomic<size_t> size = 0;
37 std::condition_variable cv_push;
38 std::condition_variable cv_pop;
42 explicit RingBuffer(
size_t max_size) : capacity(max_size), buffer(max_size) {}
44 void push(
const std::shared_ptr<T>& item)
46 std::unique_lock<std::mutex> lock(mutex);
47 cv_push.wait(lock, [
this]() {
return stop || size < capacity; });
51 head = (head + 1) % capacity;
56 std::shared_ptr<T> pop()
58 std::unique_lock<std::mutex> lock(mutex);
59 cv_pop.wait(lock, [
this]() {
return stop || size > 0; });
62 auto item = buffer[tail];
63 tail = (tail + 1) % capacity;
72 std::lock_guard<std::mutex> lock(mutex);
84 std::vector<std::shared_ptr<DataStruct>> pool;
88 BufferPool(
size_t pool_size,
size_t buffer_size)
90 for(
size_t i = 0; i < pool_size; ++i)
92 pool.emplace_back(std::make_shared<DataStruct>(buffer_size));
96 std::shared_ptr<DataStruct> acquire()
98 std::lock_guard<std::mutex> lock(mutex);
101 return std::make_shared<DataStruct>(60 * 1024 * 1024 /
sizeof(int16_t));
105 auto buffer = pool.back();
111 void release(std::shared_ptr<DataStruct> buffer)
113 std::lock_guard<std::mutex> lock(mutex);
114 pool.push_back(buffer);
122 auto buffer = pool.acquire();
123 buffer->size = (rand() % (60 * 1024 * 1024 /
sizeof(int16_t))) +
126 buffer->data.begin(), buffer->data.begin() + buffer->size, rand() % 100);
127 strncpy(buffer->metadata,
128 "Generated by kernelBufferThread",
129 sizeof(buffer->metadata));
130 bufferQueue.push(buffer);
131 std::this_thread::sleep_for(std::chrono::milliseconds(5));
141 auto buffer = inputQueue.pop();
144 strncpy(buffer->metadata + strlen(buffer->metadata),
146 sizeof(buffer->metadata) - strlen(buffer->metadata));
147 outputQueue.push(buffer);
157 auto buffer = inputQueue.pop();
160 strncpy(buffer->metadata + strlen(buffer->metadata),
162 sizeof(buffer->metadata) - strlen(buffer->metadata));
163 outputQueue.push(buffer);
173 auto buffer = inputQueue.pop();
176 strncpy(buffer->metadata + strlen(buffer->metadata),
177 " | Algorithm Processed",
178 sizeof(buffer->metadata) - strlen(buffer->metadata));
179 outputQueue.push(buffer);
185 const size_t max_queue_size = 10;
186 const size_t buffer_pool_size = 20;
187 const size_t buffer_size =
188 60 * 1024 * 1024 /
sizeof(int16_t);
191 BufferPool pool(buffer_pool_size, buffer_size);
202 std::thread kernelThread(
203 kernelBufferThread, std::ref(kernelBufferQueue), std::ref(pool));
205 std::ref(kernelBufferQueue),
206 std::ref(checkDataQueue),
209 std::ref(checkDataQueue),
210 std::ref(formEventsQueue),
212 std::thread rawProcess(processAlgorithmThread,
213 std::ref(formEventsQueue),
214 std::ref(rawBufferQueue),
216 std::thread zsProcess(processAlgorithmThread,
217 std::ref(rawBufferQueue),
218 std::ref(zsBufferQueue),
220 std::thread mwdProcess(processAlgorithmThread,
221 std::ref(zsBufferQueue),
222 std::ref(mwdBufferQueue),
226 std::this_thread::sleep_for(std::chrono::seconds(10));
229 kernelBufferQueue.shutdown();
230 checkDataQueue.shutdown();
231 formEventsQueue.shutdown();
232 rawBufferQueue.shutdown();
233 zsBufferQueue.shutdown();
234 mwdBufferQueue.shutdown();