| 233 | |
| 234 | |
| 235 | void Consumer::TaskDispatch() { |
| 236 | |
| 237 | // set nonblock for CoRead/CoWrite |
| 238 | { |
| 239 | auto flags = fcntl(impl_->consume_fds[1], F_GETFL, 0); |
| 240 | fcntl(impl_->consume_fds[1], F_SETFL, flags | O_NONBLOCK); |
| 241 | } |
| 242 | |
| 243 | comm::RetCode ret; |
| 244 | |
| 245 | char ch; |
| 246 | while (1) { |
| 247 | comm::ConsumerConsumeBP::GetThreadInstance()->OnTaskDispatch(impl_->cc); |
| 248 | |
| 249 | impl_->nhandle_task_finished = 0; |
| 250 | impl_->nbatch_handle_task_finished = 0; |
| 251 | |
| 252 | while (!comm::utils::CoRead(impl_->consume_fds[1], &ch, sizeof(char))) { |
| 253 | QLErr("CoRead fail"); |
| 254 | poll(nullptr, 0, 100); |
| 255 | } |
| 256 | |
| 257 | if (comm::RetCode::RET_OK != (ret = MakeHandleBuckets())) { |
| 258 | QLErr("MakeHandleBuckets ret %d", comm::as_integer(ret)); |
| 259 | } |
| 260 | |
| 261 | for (int i{0}; i < impl_->opt.nbatch_handler; ++i) { |
| 262 | impl_->batch_handle_finish[i] = false; |
| 263 | co_resume(impl_->batch_handle_ctxs[i].co); |
| 264 | } |
| 265 | |
| 266 | for (int i{0}; i < impl_->opt.nhandler; ++i) { |
| 267 | if (!impl_->handle_buckets[i].empty()) { |
| 268 | co_resume(impl_->handle_ctxs[i].co); |
| 269 | } |
| 270 | } |
| 271 | |
| 272 | while (1) { |
| 273 | if (IsAllTaskFinish()) break; |
| 274 | } |
| 275 | |
| 276 | comm::ConsumerConsumeBP::GetThreadInstance()->OnTaskDispatchFinish(impl_->cc); |
| 277 | |
| 278 | |
| 279 | while (!comm::utils::CoWrite(impl_->consume_fds[1], "a", sizeof(char))) { |
| 280 | QLErr("CoWrite fail"); |
| 281 | poll(nullptr, 0, 100); |
| 282 | } |
| 283 | } |
| 284 | } |
| 285 | |
| 286 | void Consumer::HandleTaskFinish() { |
| 287 | ++impl_->nhandle_task_finished; |
no test coverage detected