| 675 | } |
| 676 | |
| 677 | Status KrpcDataStreamSender::Channel::DoEndDataStreamRpc() { |
| 678 | DCHECK(rpc_in_flight_); |
| 679 | EndDataStreamRequestPB eos_req; |
| 680 | rpc_controller_.Reset(); |
| 681 | if (FLAGS_data_stream_sender_eos_timeout_ms > 0) { |
| 682 | // Provide a timeout so EOS RPCs are prioritized over others, as completing a stream |
| 683 | // can help free up resources. |
| 684 | rpc_controller_.set_timeout( |
| 685 | MonoDelta::FromMilliseconds(FLAGS_data_stream_sender_eos_timeout_ms)); |
| 686 | } |
| 687 | *eos_req.mutable_dest_fragment_instance_id() = fragment_instance_id_; |
| 688 | eos_req.set_sender_id(parent_->sender_id_); |
| 689 | eos_req.set_dest_node_id(dest_node_id_); |
| 690 | eos_resp_.Clear(); |
| 691 | rpc_start_time_ns_ = MonotonicNanos(); |
| 692 | proxy_->EndDataStreamAsync(eos_req, &eos_resp_, &rpc_controller_, |
| 693 | boost::bind(&KrpcDataStreamSender::Channel::EndDataStreamCompleteCb, this)); |
| 694 | return Status::OK(); |
| 695 | } |
| 696 | |
| 697 | Status KrpcDataStreamSender::Channel::SendEosAsync() { |
| 698 | { |
nothing calls this directly
no test coverage detected