| 269 | } |
| 270 | |
| 271 | void KrpcDataStreamMgr::DeserializeThreadFn(int thread_id, const DeserializeTask& task) { |
| 272 | int64_t queue_wait_time_ns = MonotonicNanos() - task.enqueue_time_ns; |
| 273 | if (LIKELY(queue_wait_time_ns > 0)) { |
| 274 | deserialize_queue_wait_time_ns_histogram_->Update(queue_wait_time_ns); |
| 275 | } |
| 276 | shared_ptr<KrpcDataStreamRecvr> recvr; |
| 277 | { |
| 278 | bool already_unregistered; |
| 279 | lock_guard<mutex> l(lock_); |
| 280 | recvr = FindRecvr(task.finst_id, task.dest_node_id, &already_unregistered); |
| 281 | DCHECK(recvr != nullptr || already_unregistered); |
| 282 | } |
| 283 | if (recvr != nullptr) recvr->ProcessDeferredRpc(task.sender_id); |
| 284 | } |
| 285 | |
| 286 | void KrpcDataStreamMgr::CloseSender(const EndDataStreamRequestPB* request, |
| 287 | EndDataStreamResponsePB* response, kudu::rpc::RpcContext* rpc_context) { |
nothing calls this directly
no test coverage detected