Line data Source code
1 : /* 2 : * Copyright (c) 2013 Juniper Networks, Inc. All rights reserved. 3 : */ 4 : 5 : #include "boost/asio/detail/recycling_allocator.hpp" 6 : 7 : #include "io/tcp_message_write.h" 8 : 9 : #include "base/util.h" 10 : #include "base/logging.h" 11 : #include "io/tcp_session.h" 12 : #include "io/io_log.h" 13 : 14 : using boost::asio::buffer; 15 : using boost::asio::buffer_cast; 16 : using boost::asio::mutable_buffer; 17 : using boost::system::error_code; 18 : using std::min; 19 : 20 : const int TcpMessageWriter::kDefaultWriteBufferSize; 21 : const int TcpMessageWriter::kMaxPendingBufferSize; 22 : const int TcpMessageWriter::kMinPendingBufferSize; 23 : 24 3575 : TcpMessageWriter::TcpMessageWriter(TcpSession *session, 25 3575 : size_t buffer_send_size) : 26 3575 : offset_(0), last_write_(0), buffer_send_size_(buffer_send_size), 27 3575 : session_(session) { 28 3575 : } 29 : 30 3587 : TcpMessageWriter::~TcpMessageWriter() { 31 3587 : for (BufferQueue::iterator iter = buffer_queue_.begin(); 32 3592 : iter != buffer_queue_.end(); ++iter) { 33 5 : DeleteBuffer(*iter); 34 : } 35 3587 : buffer_queue_.clear(); 36 3586 : } 37 : 38 258646 : int TcpMessageWriter::AsyncSend(const uint8_t *data, size_t len, error_code *ec) { 39 : 40 258646 : int write = len; 41 : 42 258646 : if (buffer_queue_.empty()) { 43 219725 : BufferAppend(data, len); 44 219663 : if (session_->io_strand_) { 45 219663 : boost::asio::detail::recycling_allocator<void> allocator; 46 439400 : session_->io_strand_->post(bind(&TcpSession::AsyncWriteInternal, 47 439335 : session_, TcpSessionPtr(session_)), allocator); 48 : } 49 : } else { 50 38924 : BufferAppend(data, len); 51 : } 52 : 53 258661 : if ((GetBufferQueueSize() - offset_) > TcpMessageWriter::kMaxPendingBufferSize) { 54 14 : if (!session_->write_blocked_) { 55 : /* throttle the sender */ 56 4 : session_->stats_.write_blocked++; 57 4 : session_->server_->stats_.write_blocked++; 58 4 : session_->stats_.write_block_start_time = UTCTimestampUsec(); 59 4 : session_->write_blocked_ = true; 60 : } 61 14 : write = 0; 62 : } 63 : 64 258527 : return write; 65 : } 66 : 67 260915 : void TcpMessageWriter::TriggerAsyncWrite() { 68 : 69 : /* assert if there is an async write in progress */ 70 260915 : assert(last_write_ == 0); 71 260915 : assert(!buffer_queue_.empty()); 72 : 73 260915 : boost::asio::mutable_buffer head = buffer_queue_.front(); 74 260915 : size_t remaining = buffer_size(head) - offset_; 75 260915 : last_write_ = min(buffer_send_size_, remaining); 76 : 77 : // Update socket write call statistics. 78 260915 : session_->stats_.write_calls++; 79 260915 : session_->server_->stats_.write_calls++; 80 : 81 260915 : const uint8_t *data = buffer_cast<const uint8_t *>(head) + offset_; 82 260915 : session_->AsyncWrite(data, last_write_); 83 260915 : } 84 : 85 260912 : bool TcpMessageWriter::UpdateBufferQueue(size_t wrote, bool *send_ready) { 86 : 87 260912 : assert(last_write_ == wrote); 88 260912 : assert(!buffer_queue_.empty()); 89 : 90 260912 : bool more_write = true; 91 260912 : last_write_ = 0; 92 260912 : *send_ready = false; 93 : 94 260912 : boost::asio::mutable_buffer head = buffer_queue_.front(); 95 260912 : if ((offset_ + wrote) == buffer_size(head)) { 96 258681 : offset_ = 0; 97 258681 : DeleteBuffer(head); 98 258681 : buffer_queue_.pop_front(); 99 : } else { 100 2231 : offset_ += wrote; 101 : } 102 : 103 260912 : if (session_->write_blocked_ && ((GetBufferQueueSize() - offset_) < 104 260912 : TcpMessageWriter::kMinPendingBufferSize)) { 105 4 : uint64_t blocked_usecs = UTCTimestampUsec() - 106 4 : session_->stats_.write_block_start_time; 107 4 : session_->stats_.write_blocked_duration_usecs += blocked_usecs; 108 4 : session_->server_->stats_.write_blocked_duration_usecs += blocked_usecs; 109 4 : session_->write_blocked_ = false; 110 4 : *send_ready = true; 111 : } 112 : 113 260912 : if (buffer_queue_.empty()) { 114 219739 : buffer_queue_.clear(); 115 219739 : more_write = false; 116 : } 117 : 118 260912 : return more_write; 119 : } 120 : 121 258629 : void TcpMessageWriter::BufferAppend(const uint8_t *src, int bytes) { 122 258629 : uint8_t *data = new uint8_t[bytes]; 123 258658 : memcpy(data, src, bytes); 124 258658 : mutable_buffer buffer = mutable_buffer(data, bytes); 125 258616 : buffer_queue_.push_back(buffer); 126 258558 : } 127 : 128 258686 : void TcpMessageWriter::DeleteBuffer(mutable_buffer buffer) { 129 258686 : const uint8_t *data = buffer_cast<const uint8_t *>(buffer); 130 258686 : delete[] data; 131 258686 : return; 132 : } 133 :