| 710 | } |
| 711 | |
| 712 | int KrpcDataStreamRecvr::SenderQueue::CollectPendingDeferredRpcs() { |
| 713 | lock_.DCheckLocked(); |
| 714 | DCHECK(!deferred_rpcs_.empty()); |
| 715 | |
| 716 | // Try dequeuing multiple entries from 'deferred_rpcs_' to parallelize the CPU |
| 717 | // bound deserialization work. No point in dequeuing more than number of |
| 718 | // deserialization threads available. |
| 719 | DCHECK_EQ(pending_deferred_rpcs_.size(), num_deserialize_tasks_pending_); |
| 720 | int free_deserialization_threads = |
| 721 | recvr_->mgr_->num_deserialization_threads() - num_deserialize_tasks_pending_; |
| 722 | DCHECK_GE(free_deserialization_threads, 0); |
| 723 | int num_to_try_dequeue = |
| 724 | min(free_deserialization_threads, static_cast<int>(deferred_rpcs_.size())); |
| 725 | int num_to_dequeue = 0; |
| 726 | for(; num_to_dequeue < num_to_try_dequeue; num_to_dequeue++) { |
| 727 | int64_t size = deferred_rpcs_.front()->deserialized_size; |
| 728 | if (!CanEnqueue(size)) break; |
| 729 | pending_deserialized_size_ += size; |
| 730 | if (pending_deferred_rpcs_.empty()) { |
| 731 | has_pending_deferred_rpcs_start_time_ns_ = MonotonicNanos(); |
| 732 | } |
| 733 | pending_deferred_rpcs_.push(DequeueDeferredRpc()); |
| 734 | num_deserialize_tasks_pending_++; |
| 735 | } |
| 736 | DCHECK_EQ(pending_deferred_rpcs_.size(), num_deserialize_tasks_pending_); |
| 737 | DCHECK(!batch_queue_.empty() || !pending_deferred_rpcs_.empty() |
| 738 | || num_pending_enqueue_ > 0); |
| 739 | return num_to_dequeue; |
| 740 | } |
| 741 | |
| 742 | void KrpcDataStreamRecvr::SenderQueue::RespondClosed(TransmitDataCtx* ctx) { |
| 743 | Status cancel_status = Status::Expected(TErrorCode::DATASTREAM_RECVR_CLOSED, |
nothing calls this directly
no test coverage detected