| 129 | } |
| 130 | |
| 131 | std::vector<std::unique_ptr<SerializedPage>> |
| 132 | ExchangeClient::next(uint32_t maxBytes, bool* atEnd, ContinueFuture* future) { |
| 133 | RequestSpec requestSpec; |
| 134 | std::vector<std::unique_ptr<SerializedPage>> pages; |
| 135 | { |
| 136 | std::lock_guard<std::mutex> l(queue_->mutex()); |
| 137 | *atEnd = false; |
| 138 | pages = queue_->dequeueLocked(maxBytes, atEnd, future); |
| 139 | if (*atEnd) { |
| 140 | return pages; |
| 141 | } |
| 142 | |
| 143 | if (!pages.empty() && queue_->totalBytes() > maxQueuedBytes_) { |
| 144 | return pages; |
| 145 | } |
| 146 | |
| 147 | requestSpec = pickSourcesToRequestLocked(); |
| 148 | } |
| 149 | |
| 150 | // Outside of lock |
| 151 | request(requestSpec); |
| 152 | return pages; |
| 153 | } |
| 154 | |
| 155 | void ExchangeClient::request(const RequestSpec& requestSpec) { |
| 156 | auto self = shared_from_this(); |
nothing calls this directly
no test coverage detected