| 658 | } |
| 659 | |
| 660 | void KrpcDataStreamRecvr::SenderQueue::Cancel() { |
| 661 | vector<std::unique_ptr<TransmitDataCtx>> rpcs_to_close; |
| 662 | { |
| 663 | unique_lock<SpinLock> l(lock_); |
| 664 | if (is_cancelled_) return; |
| 665 | is_cancelled_ = true; |
| 666 | |
| 667 | // Collect and dequeue deferred RPCs. Respond later without holding lock_. |
| 668 | rpcs_to_close.reserve(deferred_rpcs_.size() + pending_deferred_rpcs_.size()); |
| 669 | while (!pending_deferred_rpcs_.empty()) { |
| 670 | unique_ptr<TransmitDataCtx> ctx = DequeuePendingDeferredRpc(); |
| 671 | recvr_->deferred_rpc_tracker()->Release(ctx->rpc_context->GetTransferSize()); |
| 672 | rpcs_to_close.push_back(move(ctx)); |
| 673 | } |
| 674 | while (!deferred_rpcs_.empty()) { |
| 675 | unique_ptr<TransmitDataCtx> ctx = DequeueDeferredRpc(); |
| 676 | recvr_->deferred_rpc_tracker()->Release(ctx->rpc_context->GetTransferSize()); |
| 677 | rpcs_to_close.push_back(move(ctx)); |
| 678 | } |
| 679 | } |
| 680 | for (auto& ctx: rpcs_to_close) { |
| 681 | RespondClosed(ctx.get()); |
| 682 | } |
| 683 | rpcs_to_close.clear(); |
| 684 | VLOG(2) << "cancelled stream: fragment_instance_id=" |
| 685 | << PrintId(recvr_->fragment_instance_id()) |
| 686 | << " node_id=" << recvr_->dest_node_id(); |
| 687 | // Wake up all threads waiting to produce/consume batches. They will all |
| 688 | // notice that the stream is cancelled and handle it. |
| 689 | data_arrival_cv_.notify_all(); |
| 690 | PeriodicCounterUpdater::StopTimeSeriesCounter( |
| 691 | recvr_->bytes_received_time_series_counter_); |
| 692 | } |
| 693 | |
| 694 | void KrpcDataStreamRecvr::SenderQueue::Close() { |
| 695 | unique_lock<SpinLock> l(lock_); |
no test coverage detected