MCPcopy Create free account
hub / github.com/apache/incubator-pegasus / call_task

Method call_task

src/client/partition_resolver.cpp:57–130  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

55}
56
57void 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 }

Callers 1

callMethod · 0.80

Calls 15

dsn_now_msFunction · 0.85
error_retryFunction · 0.85
enum_to_stringFunction · 0.85
enqueueFunction · 0.85
dsn_rpc_callFunction · 0.85
fetch_current_handlerMethod · 0.80
replace_callbackMethod · 0.80
set_retryMethod · 0.80
stateMethod · 0.80
thread_hashMethod · 0.80
getMethod · 0.65
get_requestMethod · 0.45

Tested by

no test coverage detected