pulsatrix
Loading...
Searching...
No Matches
data_thread_pool.hpp
Go to the documentation of this file.
1
5#pragma once
6
7#include <condition_variable>
8#include <functional>
9#include <future>
10#include <memory>
11#include <mutex>
12#include <queue>
13#include <stdexcept>
14#include <thread>
15#include <type_traits>
16#include <utility>
17#include <vector>
18
19namespace pulsatrix {
20
30public:
31 explicit DataThreadPool(unsigned num_threads);
33
36
44 template <typename F>
45 auto submit(F&& task) -> std::future<std::invoke_result_t<F>> {
46 using ReturnType = std::invoke_result_t<F>;
47 auto packaged = std::make_shared<std::packaged_task<ReturnType()>>(std::forward<F>(task));
48 std::future<ReturnType> result = packaged->get_future();
49 {
50 std::lock_guard<std::mutex> lock(mutex_);
51 if (stop_) {
52 throw std::runtime_error("DataThreadPool::submit: pool is stopped");
53 }
54 tasks_.emplace([packaged]() { (*packaged)(); });
55 }
56 cv_.notify_one();
57 return result;
58 }
59
60private:
61 std::vector<std::thread> workers_;
62 std::queue<std::function<void()>> tasks_;
63 std::mutex mutex_;
64 std::condition_variable cv_;
65 bool stop_ = false;
66};
67
68} // namespace pulsatrix
Minimal, generic thread pool for CPU-side data pipeline work (fetch, decode, transform,...
Definition data_thread_pool.hpp:29
DataThreadPool(unsigned num_threads)
DataThreadPool & operator=(const DataThreadPool &)=delete
auto submit(F &&task) -> std::future< std::invoke_result_t< F > >
Enqueues a task for execution by a worker thread.
Definition data_thread_pool.hpp:45
DataThreadPool(const DataThreadPool &)=delete
Definition acquisition_functions.hpp:16