| 55 | } |
| 56 | |
| 57 | void partition_resolver::call_task(const rpc_response_task_ptr &t) |
| 58 | { |
| 59 | auto &hdr = *(t->get_request()->header); |
| 60 | uint64_t deadline_ms = dsn_now_ms() + hdr.client.timeout_ms; |
| 61 | |
| 62 | rpc_response_handler old_callback; |
| 63 | t->fetch_current_handler(old_callback); |
| 64 | auto new_callback = [ this, deadline_ms, oc = std::move(old_callback) ]( |
| 65 | dsn::error_code err, dsn::message_ex * req, dsn::message_ex * resp) |
| 66 | { |
| 67 | if (req->header->gpid.value() != 0 && err != ERR_OK && error_retry(err)) { |
| 68 | on_access_failure(req->header->gpid.get_partition_index(), err); |
| 69 | // still got time, retry |
| 70 | uint64_t nms = dsn_now_ms(); |
| 71 | uint64_t gap = 8 << req->send_retry_count; |
| 72 | if (gap > 1000) |
| 73 | gap = 1000; |
| 74 | if (nms + gap < deadline_ms) { |
| 75 | req->send_retry_count++; |
| 76 | req->header->client.timeout_ms = static_cast<int>(deadline_ms - nms - gap); |
| 77 | |
| 78 | rpc_response_task_ptr ctask = |
| 79 | dynamic_cast<rpc_response_task *>(task::get_current_task()); |
| 80 | partition_resolver_ptr r(this); |
| 81 | |
| 82 | CHECK_NOTNULL(ctask, "current task must be rpc_response_task"); |
| 83 | ctask->replace_callback(std::move(oc)); |
| 84 | CHECK(ctask->set_retry(false), |
| 85 | "rpc_response_task set retry failed, state = {}", |
| 86 | enum_to_string(ctask->state())); |
| 87 | |
| 88 | // sleep gap milliseconds before retry |
| 89 | tasking::enqueue(LPC_RPC_DELAY_CALL, |
| 90 | nullptr, |
| 91 | [r, ctask]() { r->call_task(ctask); }, |
| 92 | 0, |
| 93 | std::chrono::milliseconds(gap)); |
| 94 | return; |
| 95 | } else { |
| 96 | LOG_ERROR("service access failed ({}), no more time for further tries, set error " |
| 97 | "= ERR_TIMEOUT, trace_id = {:#018x}", |
| 98 | err, |
| 99 | req->header->trace_id); |
| 100 | err = ERR_TIMEOUT; |
| 101 | } |
| 102 | } |
| 103 | |
| 104 | if (oc) |
| 105 | oc(err, req, resp); |
| 106 | }; |
| 107 | t->replace_callback(std::move(new_callback)); |
| 108 | |
| 109 | resolve(hdr.client.partition_hash, |
| 110 | [t](resolve_result &&result) mutable { |
| 111 | if (result.err != ERR_OK) { |
| 112 | t->enqueue(result.err, nullptr); |
| 113 | return; |
| 114 | } |
no test coverage detected