| 599 | } |
| 600 | |
| 601 | void KrpcDataStreamRecvr::SenderQueue::TakeOverEarlySender( |
| 602 | unique_ptr<TransmitDataCtx> ctx) { |
| 603 | // TakeOverEarlySender() is called by the same thread which calls Close(). |
| 604 | // The receiver cannot be closed while this function is in progress so |
| 605 | // 'recvr_->mgr_' shouldn't be NULL. |
| 606 | DCHECK(TestInfo::is_test() || FragmentInstanceState::IsFragmentExecThread()); |
| 607 | DCHECK(!recvr_->closed_ && recvr_->mgr_ != nullptr); |
| 608 | COUNTER_ADD(recvr_->total_received_batches_counter_, 1); |
| 609 | int64_t serialized_size; |
| 610 | Status status = GetBatchSize(ctx->request, ctx->rpc_context, &serialized_size); |
| 611 | if (UNLIKELY(!status.ok())) { |
| 612 | { |
| 613 | unique_lock<SpinLock> l(lock_); |
| 614 | if (!is_cancelled_) MarkErrorStatus(status, l); |
| 615 | } |
| 616 | TRACE_TO( |
| 617 | ctx->rpc_context->trace(), "Error unpacking request: $0", status.GetDetail()); |
| 618 | DataStreamService::RespondRpc(status, ctx->response, ctx->rpc_context); |
| 619 | return; |
| 620 | } |
| 621 | COUNTER_ADD(recvr_->bytes_received_counter_, serialized_size); |
| 622 | DCHECK_GE(ctx->deserialized_size, 0); |
| 623 | int sender_id = ctx->request->sender_id(); |
| 624 | int num_to_dequeue = 0; |
| 625 | { |
| 626 | unique_lock<SpinLock> l(lock_); |
| 627 | // Only enqueue a deferred RPC if the sender queue is not yet cancelled. |
| 628 | if (UNLIKELY(is_cancelled_)) { |
| 629 | l.unlock(); |
| 630 | TRACE_TO(ctx->rpc_context->trace(), "Recvr closed"); |
| 631 | RespondClosed(ctx.get()); |
| 632 | return; |
| 633 | } |
| 634 | recvr_->deferred_rpc_tracker()->Consume(ctx->rpc_context->GetTransferSize()); |
| 635 | EnqueueDeferredRpc(move(ctx)); |
| 636 | num_to_dequeue = CollectPendingDeferredRpcs(); |
| 637 | } |
| 638 | |
| 639 | if (num_to_dequeue > 0) { |
| 640 | recvr_->mgr_->EnqueueDeserializeTask(recvr_->fragment_instance_id(), |
| 641 | recvr_->dest_node_id(), sender_id, num_to_dequeue); |
| 642 | } |
| 643 | } |
| 644 | |
| 645 | bool KrpcDataStreamRecvr::SenderQueue::DecrementSenders() { |
| 646 | lock_guard<SpinLock> l(lock_); |
no test coverage detected