| 131 | |
| 132 | |
| 133 | void OldThreadPool::ThreadMain(int thread_id, int device_id, bool set_affinity, |
| 134 | const std::string &name) { |
| 135 | this_thread_idx_ = thread_id; |
| 136 | SetThreadName(name.c_str()); |
| 137 | DeviceGuard g(device_id); |
| 138 | try { |
| 139 | #if NVML_ENABLED |
| 140 | if (set_affinity) { |
| 141 | const char *env_affinity = std::getenv("DALI_AFFINITY_MASK"); |
| 142 | int core = -1; |
| 143 | if (env_affinity) { |
| 144 | const auto &vec = string_split(env_affinity, ','); |
| 145 | if ((size_t)thread_id < vec.size()) { |
| 146 | core = std::stoi(vec[thread_id]); |
| 147 | } else { |
| 148 | DALI_WARN("DALI_AFFINITY_MASK environment variable is set, " |
| 149 | "but does not have enough entries: thread_id (", thread_id, |
| 150 | ") vs #entries (", vec.size(), "). Ignoring..."); |
| 151 | } |
| 152 | } |
| 153 | nvml::SetCPUAffinity(core); |
| 154 | } |
| 155 | #endif |
| 156 | } catch (...) { |
| 157 | tl_errors_[thread_id].push(std::current_exception()); |
| 158 | } |
| 159 | |
| 160 | while (running_) { |
| 161 | // Wait for something to do |
| 162 | queue_semaphore_.acquire(); |
| 163 | |
| 164 | // This lock guards only the queue, not the condition - that's handled by the semaphore |
| 165 | std::unique_lock lock(queue_lock_); |
| 166 | |
| 167 | if (!running_) |
| 168 | break; |
| 169 | |
| 170 | // Get work from the queue. |
| 171 | WorkWithThreadIdx work = std::move(work_queue_.top().second); |
| 172 | work_queue_.pop(); |
| 173 | // Unlock the lock |
| 174 | lock.unlock(); |
| 175 | |
| 176 | // If an error occurs, we save it in tl_errors_. When |
| 177 | // WaitForWork is called, we will check for any errors |
| 178 | // in the threads and return an error if one occured. |
| 179 | try { |
| 180 | work(thread_id); |
| 181 | } catch (...) { |
| 182 | tl_errors_[thread_id].push(std::current_exception()); |
| 183 | } |
| 184 | |
| 185 | // The task is now complete - we can atomically decrement the number of outstanding work. |
| 186 | // If it reaches zero, we must safely notify the potential threads waiting for the work |
| 187 | // to complete. |
| 188 | // NOTE: We don't have to acquire the mutex until the number of waiting threads reaches 0. |
| 189 | if (--outstanding_work_ == 0) { |
| 190 | // We don't need to guard the modification of the atomic value with a mutex - |
nothing calls this directly
no test coverage detected