| 90 | } |
| 91 | |
| 92 | RowVectorPtr TaskQueue::dequeue() { |
| 93 | for (;;) { |
| 94 | RowVectorPtr vector; |
| 95 | std::vector<ContinuePromise> mayContinue; |
| 96 | { |
| 97 | std::lock_guard<std::mutex> l(mutex_); |
| 98 | if (closed_) { |
| 99 | return nullptr; |
| 100 | } |
| 101 | |
| 102 | if (!queue_.empty()) { |
| 103 | auto result = std::move(queue_.front()); |
| 104 | queue_.pop_front(); |
| 105 | totalBytes_ -= result.bytes; |
| 106 | vector = std::move(result.vector); |
| 107 | if (totalBytes_ < maxBytes_ / 2) { |
| 108 | mayContinue = std::move(producerUnblockPromises_); |
| 109 | } |
| 110 | } else if ( |
| 111 | numProducers_.has_value() && producersFinished_ == numProducers_) { |
| 112 | return nullptr; |
| 113 | } |
| 114 | if (!vector) { |
| 115 | consumerBlocked_ = true; |
| 116 | consumerPromise_ = ContinuePromise(); |
| 117 | consumerFuture_ = consumerPromise_.getFuture(); |
| 118 | } |
| 119 | } |
| 120 | // outside of 'mutex_' |
| 121 | for (auto& promise : mayContinue) { |
| 122 | promise.setValue(); |
| 123 | } |
| 124 | if (vector) { |
| 125 | return vector; |
| 126 | } |
| 127 | consumerFuture_.wait(); |
| 128 | } |
| 129 | } |
| 130 | |
| 131 | void TaskQueue::close() { |
| 132 | std::lock_guard<std::mutex> l(mutex_); |