| 90 | } |
| 91 | |
| 92 | int ThreadPoolExecutor::start() noexcept { |
| 93 | if (_running.load(::std::memory_order_acquire)) { |
| 94 | return -1; |
| 95 | } |
| 96 | _running.store(true, ::std::memory_order_release); |
| 97 | _global_task_queue.reserve_and_clear(_global_capacity * 2); |
| 98 | _local_task_queues.set_constructor( |
| 99 | [this](ConcurrentBoundedQueue<Task>* queue) { |
| 100 | new (queue) ConcurrentBoundedQueue<Task>; |
| 101 | queue->reserve_and_clear(_local_capacity * 2); |
| 102 | }); |
| 103 | _threads.reserve(_worker_number); |
| 104 | for (size_t i = 0; i < _worker_number; ++i) { |
| 105 | _threads.emplace_back(&ThreadPoolExecutor::keep_execute, this); |
| 106 | } |
| 107 | if (_balance_interval.count() >= 0) { |
| 108 | _balance_thread = ::std::thread(&ThreadPoolExecutor::keep_balance, this); |
| 109 | } |
| 110 | return 0; |
| 111 | } |
| 112 | |
| 113 | void ThreadPoolExecutor::wakeup_one_worker() noexcept { |
| 114 | _global_task_queue.push<true, false, true>( |
nothing calls this directly
no test coverage detected