| 190 | } |
| 191 | |
| 192 | void DequeueAllRunner(TLockFreeQueue<int>& queue, bool singleConsumer) { |
| 193 | size_t threadsNum = 4; |
| 194 | size_t enqueuesPerThread = 10'000; |
| 195 | TThreadPool p; |
| 196 | p.Start(threadsNum, 0); |
| 197 | |
| 198 | TVector<NThreading::TFuture<void>> futures; |
| 199 | |
| 200 | for (size_t i = 0; i < threadsNum; ++i) { |
| 201 | NThreading::TPromise<void> promise = NThreading::NewPromise(); |
| 202 | futures.emplace_back(promise.GetFuture()); |
| 203 | |
| 204 | p.SafeAddFunc([enqueuesPerThread, &queue, promise]() mutable { |
| 205 | for (size_t i = 0; i != enqueuesPerThread; ++i) { |
| 206 | queue.Enqueue(i); |
| 207 | } |
| 208 | |
| 209 | promise.SetValue(); |
| 210 | }); |
| 211 | } |
| 212 | |
| 213 | std::atomic<size_t> elementsLeft = threadsNum * enqueuesPerThread; |
| 214 | |
| 215 | ui64 numOfConsumers = singleConsumer ? 1 : threadsNum; |
| 216 | |
| 217 | TVector<TVector<int>> dataBuckets(numOfConsumers); |
| 218 | |
| 219 | for (size_t i = 0; i < numOfConsumers; ++i) { |
| 220 | NThreading::TPromise<void> promise = NThreading::NewPromise(); |
| 221 | futures.emplace_back(promise.GetFuture()); |
| 222 | |
| 223 | p.SafeAddFunc([&queue, &elementsLeft, promise, consumerData{&dataBuckets[i]}]() mutable { |
| 224 | TVector<int> vec; |
| 225 | while (static_cast<i64>(elementsLeft.load()) > 0) { |
| 226 | for (size_t i = 0; i != 100; ++i) { |
| 227 | vec.clear(); |
| 228 | queue.DequeueAll(&vec); |
| 229 | |
| 230 | elementsLeft -= vec.size(); |
| 231 | consumerData->insert(consumerData->end(), vec.begin(), vec.end()); |
| 232 | } |
| 233 | } |
| 234 | |
| 235 | promise.SetValue(); |
| 236 | }); |
| 237 | } |
| 238 | |
| 239 | NThreading::WaitExceptionOrAll(futures).GetValueSync(); |
| 240 | p.Stop(); |
| 241 | |
| 242 | TVector<int> left; |
| 243 | queue.DequeueAll(&left); |
| 244 | |
| 245 | UNIT_ASSERT(left.empty()); |
| 246 | |
| 247 | TVector<int> data; |
| 248 | for (auto& dataBucket : dataBuckets) { |
| 249 | data.insert(data.end(), dataBucket.begin(), dataBucket.end()); |
no test coverage detected