| 156 | } |
| 157 | |
| 158 | void SCQueueSynchronizer::producer_wait() { |
| 159 | auto wait_target = m_tot_task.load(std::memory_order_relaxed); |
| 160 | if (m_worker_started && |
| 161 | m_finished_task.load(std::memory_order_acquire) < wait_target) { |
| 162 | std::unique_lock<std::mutex> lock(m_mtx_finished); |
| 163 | // update wait_target again in this critical section |
| 164 | wait_target = m_tot_task.load(std::memory_order_relaxed); |
| 165 | if (m_waiter_target_queue.empty()) { |
| 166 | m_waiter_target.store(wait_target, std::memory_order_relaxed); |
| 167 | m_waiter_target_queue.push_back(wait_target); |
| 168 | } else { |
| 169 | mgb_assert(wait_target >= m_waiter_target_queue.back()); |
| 170 | if (wait_target > m_waiter_target_queue.back()) { |
| 171 | m_waiter_target_queue.push_back(wait_target); |
| 172 | } |
| 173 | } |
| 174 | |
| 175 | size_t done; |
| 176 | for (;;) { |
| 177 | // ensure that m_waiter_target is visible in consumer |
| 178 | std::atomic_thread_fence(std::memory_order_seq_cst); |
| 179 | |
| 180 | done = m_finished_task.load(std::memory_order_relaxed); |
| 181 | if (done >= wait_target) |
| 182 | break; |
| 183 | m_cv_finished.wait(lock); |
| 184 | } |
| 185 | |
| 186 | if (!m_waiter_target_queue.empty()) { |
| 187 | size_t next_target = 0; |
| 188 | while (done >= (next_target = m_waiter_target_queue.front())) { |
| 189 | m_waiter_target_queue.pop_front(); |
| 190 | if (m_waiter_target_queue.empty()) { |
| 191 | next_target = std::numeric_limits<size_t>::max(); |
| 192 | break; |
| 193 | } |
| 194 | } |
| 195 | m_waiter_target.store(next_target, std::memory_order_release); |
| 196 | // this is necessary in practice, although not needed logically |
| 197 | m_cv_finished.notify_all(); |
| 198 | } |
| 199 | } |
| 200 | m_wait_finish_called = true; |
| 201 | } |
| 202 | |
| 203 | size_t SCQueueSynchronizer::consumer_fetch(size_t max, size_t min) { |
| 204 | mgb_assert(max >= min && min >= 1); |