| 998 | } |
| 999 | |
| 1000 | void Consumer::AfterConsume(const comm::proto::ConsumerContext &cc, |
| 1001 | const std::vector<std::shared_ptr<comm::proto::QItem> > &items, |
| 1002 | const std::vector<comm::HandleResult> &handle_results) { |
| 1003 | |
| 1004 | phxqueue::comm::RetCode ret; |
| 1005 | |
| 1006 | if (handle_results.size() < items.size()) { |
| 1007 | QLErr("handle_results.size() %d < items.size() %d", (int)handle_results.size(), (int)items.size()); |
| 1008 | return; |
| 1009 | } |
| 1010 | |
| 1011 | vector<shared_ptr<phxqueue::comm::proto::QItem> > retry_items; |
| 1012 | for (int i{0}; i < items.size(); ++i) { |
| 1013 | auto &&item = items[i]; |
| 1014 | auto &&handle_result = handle_results[i]; |
| 1015 | if (phxqueue::comm::HandleResult::RES_ERROR == handle_result) { |
| 1016 | retry_items.push_back(item); |
| 1017 | } |
| 1018 | } |
| 1019 | |
| 1020 | if (0 == retry_items.size()) return; |
| 1021 | |
| 1022 | vector<unique_ptr<phxqueue::comm::proto::AddRequest> > reqs; |
| 1023 | |
| 1024 | uint64_t retry_sub_ids = (1ULL << (cc.sub_id() - 1)); |
| 1025 | { |
| 1026 | ret = phxqueue::producer::Producer::MakeAddRequests(cc.topic_id(), retry_items, reqs, |
| 1027 | [&cc, retry_sub_ids](phxqueue::comm::proto::QItem &item)->void { |
| 1028 | item.set_count(item.count() + 1); |
| 1029 | item.set_sub_ids(retry_sub_ids); |
| 1030 | |
| 1031 | auto now = comm::utils::Time::GetTimestampMS(); |
| 1032 | item.set_atime(now / 1000); |
| 1033 | item.set_atime_ms(now % 1000); |
| 1034 | }); |
| 1035 | if (phxqueue::comm::RetCode::RET_OK != ret) { |
| 1036 | QLErr("MakeAddRequests ret %d", as_integer(ret)); |
| 1037 | return; |
| 1038 | } |
| 1039 | } |
| 1040 | |
| 1041 | for (auto &&req : reqs) { |
| 1042 | if (!req || 0 == req->items_size()) continue; |
| 1043 | |
| 1044 | int retry_pub_id = req->items(0).pub_id(); |
| 1045 | |
| 1046 | BeforeAdd(*req); |
| 1047 | |
| 1048 | comm::proto::AddResponse resp; |
| 1049 | while (true) { // retry forever |
| 1050 | if (impl_->opt.use_store_master_client_on_add) { |
| 1051 | store::StoreMasterClient<comm::proto::AddRequest, comm::proto::AddResponse> store_master_client; |
| 1052 | ret = store_master_client.ClientCall(*req, resp, bind(&Consumer::Add, this, placeholders::_1, placeholders::_2)); |
| 1053 | } else { |
| 1054 | ret = Add(*req, resp); |
| 1055 | } |
| 1056 | if (comm::RetCode::RET_OK != ret) { |
| 1057 | QLErr("Retry ret %d topic_id %d retry_pub_id %d sub_id %d store_id %d queue_id %d item_size %zu retry_item_size %d", |
nothing calls this directly
no test coverage detected