| 207 | } |
| 208 | |
| 209 | void ExecutorTasks::init(size_t num_threads_, size_t use_threads_, const SlotAllocationPtr & cpu_slots_, bool profile_processors, bool trace_processors, ReadProgressCallback * callback) |
| 210 | { |
| 211 | num_threads = num_threads_; |
| 212 | use_threads = use_threads_; |
| 213 | threads_queue.init(num_threads); |
| 214 | task_queue.init(num_threads); |
| 215 | fast_task_queue.init(num_threads); |
| 216 | |
| 217 | { |
| 218 | std::lock_guard lock(mutex); // In case finish() is executed concurrently with init() due to exception |
| 219 | cpu_slots = cpu_slots_; |
| 220 | } |
| 221 | |
| 222 | // Initialize slot counters with zeros up to max_threads |
| 223 | slot_count.resize(num_threads, 0); |
| 224 | |
| 225 | { |
| 226 | std::lock_guard guard(executor_contexts_mutex); |
| 227 | |
| 228 | executor_contexts.reserve(num_threads); |
| 229 | for (size_t i = 0; i < num_threads; ++i) |
| 230 | executor_contexts.emplace_back(std::make_unique<ExecutionThreadContext>(i, profile_processors, trace_processors, callback)); |
| 231 | } |
| 232 | } |
| 233 | |
| 234 | size_t ExecutorTasks::fill(Queue & queue, [[maybe_unused]] Queue & async_queue) |
| 235 | { |