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 29951 : TcpMessageWriter::TcpMessageWriter(TcpSession *session, 25 29951 : size_t buffer_send_size) : 26 29953 : offset_(0), last_write_(0), buffer_send_size_(buffer_send_size), 27 29951 : session_(session) { 28 29953 : } 29 : 30 29962 : TcpMessageWriter::~TcpMessageWriter() { 31 29962 : for (BufferQueue::iterator iter = buffer_queue_.begin(); 32 30771 : iter != buffer_queue_.end(); ++iter) { 33 811 : DeleteBuffer(*iter); 34 : } 35 29956 : buffer_queue_.clear(); 36 29952 : } 37 : 38 1847269 : int TcpMessageWriter::AsyncSend(const uint8_t *data, size_t len, error_code *ec) { 39 : 40 1847269 : int write = len; 41 : 42 1847269 : if (buffer_queue_.empty()) { 43 1568161 : BufferAppend(data, len); 44 1567916 : if (session_->io_strand_) { 45 1567912 : boost::asio::detail::recycling_allocator<void> allocator; 46 3136179 : session_->io_strand_->post(bind(&TcpSession::AsyncWriteInternal, 47 3135615 : session_, TcpSessionPtr(session_)), allocator); 48 : } 49 : } else { 50 279150 : BufferAppend(data, len); 51 : } 52 : 53 1847318 : if ((GetBufferQueueSize() - offset_) > TcpMessageWriter::kMaxPendingBufferSize) { 54 14 : if (!session_->write_blocked_) { 55 : /* throttle the sender */ 56 5 : session_->stats_.write_blocked++; 57 5 : session_->server_->stats_.write_blocked++; 58 5 : session_->stats_.write_block_start_time = UTCTimestampUsec(); 59 5 : session_->write_blocked_ = true; 60 : } 61 15 : write = 0; 62 : } 63 : 64 1846448 : return write; 65 : } 66 : 67 1850161 : void TcpMessageWriter::TriggerAsyncWrite() { 68 : 69 : /* assert if there is an async write in progress */ 70 1850161 : assert(last_write_ == 0); 71 1850161 : assert(!buffer_queue_.empty()); 72 : 73 1850161 : boost::asio::mutable_buffer head = buffer_queue_.front(); 74 1850161 : size_t remaining = buffer_size(head) - offset_; 75 1850161 : last_write_ = min(buffer_send_size_, remaining); 76 : 77 : // Update socket write call statistics. 78 1850161 : session_->stats_.write_calls++; 79 1850161 : session_->server_->stats_.write_calls++; 80 : 81 1850161 : const uint8_t *data = buffer_cast<const uint8_t *>(head) + offset_; 82 1850161 : session_->AsyncWrite(data, last_write_); 83 1850161 : } 84 : 85 1849424 : bool TcpMessageWriter::UpdateBufferQueue(size_t wrote, bool *send_ready) { 86 : 87 1849424 : assert(last_write_ == wrote); 88 1849424 : assert(!buffer_queue_.empty()); 89 : 90 1849424 : bool more_write = true; 91 1849424 : last_write_ = 0; 92 1849424 : *send_ready = false; 93 : 94 1849424 : boost::asio::mutable_buffer head = buffer_queue_.front(); 95 1849424 : if ((offset_ + wrote) == buffer_size(head)) { 96 1846676 : offset_ = 0; 97 1846676 : DeleteBuffer(head); 98 1846676 : buffer_queue_.pop_front(); 99 : } else { 100 2748 : offset_ += wrote; 101 : } 102 : 103 1849424 : if (session_->write_blocked_ && ((GetBufferQueueSize() - offset_) < 104 1849424 : TcpMessageWriter::kMinPendingBufferSize)) { 105 5 : uint64_t blocked_usecs = UTCTimestampUsec() - 106 5 : session_->stats_.write_block_start_time; 107 5 : session_->stats_.write_blocked_duration_usecs += blocked_usecs; 108 5 : session_->server_->stats_.write_blocked_duration_usecs += blocked_usecs; 109 5 : session_->write_blocked_ = false; 110 5 : *send_ready = true; 111 : } 112 : 113 1849424 : if (buffer_queue_.empty()) { 114 1567539 : buffer_queue_.clear(); 115 1567539 : more_write = false; 116 : } 117 : 118 1849424 : return more_write; 119 : } 120 : 121 1847305 : void TcpMessageWriter::BufferAppend(const uint8_t *src, int bytes) { 122 1847305 : uint8_t *data = new uint8_t[bytes]; 123 1847384 : memcpy(data, src, bytes); 124 1847384 : mutable_buffer buffer = mutable_buffer(data, bytes); 125 1847230 : buffer_queue_.push_back(buffer); 126 1846923 : } 127 : 128 1847486 : void TcpMessageWriter::DeleteBuffer(mutable_buffer buffer) { 129 1847486 : const uint8_t *data = buffer_cast<const uint8_t *>(buffer); 130 1847486 : delete[] data; 131 1847487 : return; 132 : } 133 :