4 #include <condition_variable>
17 std::vector<int16_t>
data;
23 std::fill(std::begin(metadata), std::end(metadata),
'\0');
32 std::vector<std::shared_ptr<T>> buffer;
36 std::atomic<size_t> size = 0;
38 std::condition_variable cv_push;
39 std::condition_variable cv_pop;
43 explicit RingBuffer(
size_t max_size) : capacity(max_size), buffer(max_size) {}
45 void push(
const std::shared_ptr<T>& item)
47 std::unique_lock<std::mutex> lock(mutex);
48 cv_push.wait(lock, [
this]() {
return stop || size < capacity; });
52 head = (head + 1) % capacity;
57 std::shared_ptr<T> pop()
59 std::unique_lock<std::mutex> lock(mutex);
60 cv_pop.wait(lock, [
this]() {
return stop || size > 0; });
63 auto item = buffer[tail];
64 tail = (tail + 1) % capacity;
73 std::lock_guard<std::mutex> lock(mutex);
85 std::vector<std::shared_ptr<DataStruct>> pool;
89 BufferPool(
size_t pool_size,
size_t buffer_size)
91 for(
size_t i = 0; i < pool_size; ++i)
93 pool.emplace_back(std::make_shared<DataStruct>(buffer_size));
97 std::shared_ptr<DataStruct> acquire()
99 std::lock_guard<std::mutex> lock(mutex);
102 return std::make_shared<DataStruct>(60 * 1024 * 1024 /
sizeof(int16_t));
106 auto buffer = pool.back();
112 void release(std::shared_ptr<DataStruct> buffer)
114 std::lock_guard<std::mutex> lock(mutex);
115 pool.push_back(buffer);
123 std::vector<std::thread> workers;
124 std::queue<std::function<void()>> tasks;
126 std::condition_variable cv;
132 for(
size_t i = 0; i < thread_count; ++i)
134 workers.emplace_back([
this] {
137 std::function<void()> task;
139 std::unique_lock<std::mutex> lock(mutex);
140 cv.wait(lock, [
this]() {
return stop || !tasks.empty(); });
141 if(stop && tasks.empty())
143 task = std::move(tasks.front());
152 void enqueue(std::function<
void()> task)
155 std::lock_guard<std::mutex> lock(mutex);
156 tasks.emplace(std::move(task));
164 std::lock_guard<std::mutex> lock(mutex);
168 for(
auto& worker : workers)
178 auto buffer = pool.acquire();
180 (rand() % (60 * 1024 * 1024 /
sizeof(int16_t))) + 1000;
181 std::fill(buffer->data.begin(), buffer->data.begin() + buffer->size, rand() % 100);
183 buffer->metadata,
"Generated by kernelBufferThread",
sizeof(buffer->metadata));
184 bufferQueue.push(buffer);
193 auto buffer = inputQueue.pop();
196 strncat(buffer->metadata,
198 sizeof(buffer->metadata) - strlen(buffer->metadata) - 1);
199 outputQueue.push(buffer);
205 const size_t max_queue_size = 10;
206 const size_t buffer_pool_size = 20;
207 const size_t buffer_size =
208 60 * 1024 * 1024 /
sizeof(int16_t);
209 const size_t thread_pool_size = 6;
212 BufferPool pool(buffer_pool_size, buffer_size);
226 for(
int i = 0; i < 100; ++i)
228 threadPool.enqueue([&]() { kernelBufferTask(kernelBufferQueue, pool); });
229 threadPool.enqueue([&]() {
230 processingTask(kernelBufferQueue, checkDataQueue, pool,
" | Checked");
232 threadPool.enqueue([&]() {
233 processingTask(checkDataQueue, formEventsQueue, pool,
" | Events Formed");
235 threadPool.enqueue([&]() {
237 formEventsQueue, rawBufferQueue, pool,
" | Algorithm Processed");
239 threadPool.enqueue([&]() {
240 processingTask(rawBufferQueue, zsBufferQueue, pool,
" | ZS Processed");
242 threadPool.enqueue([&]() {
243 processingTask(zsBufferQueue, mwdBufferQueue, pool,
" | MWD Processed");
248 std::this_thread::sleep_for(std::chrono::seconds(10));
251 threadPool.shutdown();
254 kernelBufferQueue.shutdown();
255 checkDataQueue.shutdown();
256 formEventsQueue.shutdown();
257 rawBufferQueue.shutdown();
258 zsBufferQueue.shutdown();
259 mwdBufferQueue.shutdown();