| 200 | } |
| 201 | |
| 202 | void worker() { |
| 203 | int error_count = 0; |
| 204 | bool delete_recv = false; |
| 205 | bool delete_send = false; |
| 206 | while (true) { |
| 207 | { |
| 208 | std::unique_lock lock(queue_mutex_); |
| 209 | |
| 210 | if (delete_recv) { |
| 211 | recvs_.front().promise.set_value(); |
| 212 | recvs_.pop_front(); |
| 213 | delete_recv = false; |
| 214 | } |
| 215 | if (delete_send) { |
| 216 | sends_.front().promise.set_value(); |
| 217 | sends_.pop_front(); |
| 218 | delete_send = false; |
| 219 | } |
| 220 | |
| 221 | if (stop_) { |
| 222 | return; |
| 223 | } |
| 224 | |
| 225 | if (!have_tasks()) { |
| 226 | condition_.wait(lock, [this] { return stop_ || have_tasks(); }); |
| 227 | if (stop_) { |
| 228 | return; |
| 229 | } |
| 230 | } |
| 231 | } |
| 232 | |
| 233 | if (!recvs_.empty()) { |
| 234 | auto& task = recvs_.front(); |
| 235 | ssize_t r = ::recv(fd_, task.buffer, task.size, 0); |
| 236 | if (r > 0) { |
| 237 | task.buffer = static_cast<char*>(task.buffer) + r; |
| 238 | task.size -= r; |
| 239 | delete_recv = task.size == 0; |
| 240 | error_count = 0; |
| 241 | } else if (errno != EAGAIN) { |
| 242 | error_count++; |
| 243 | log_info( |
| 244 | true, "Receiving from socket", fd_, "failed with errno", errno); |
| 245 | } |
| 246 | } |
| 247 | if (!sends_.empty()) { |
| 248 | auto& task = sends_.front(); |
| 249 | ssize_t r = ::send(fd_, task.buffer, task.size, 0); |
| 250 | if (r > 0) { |
| 251 | task.buffer = static_cast<char*>(task.buffer) + r; |
| 252 | task.size -= r; |
| 253 | delete_send = task.size == 0; |
| 254 | error_count = 0; |
| 255 | } else if (errno != EAGAIN) { |
| 256 | error_count++; |
| 257 | log_info(true, "Sending to socket", fd_, "failed with errno", errno); |
| 258 | } |
| 259 | } |