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 32076 : TcpMessageWriter::TcpMessageWriter(TcpSession *session, 25 32076 : size_t buffer_send_size) : 26 32079 : offset_(0), last_write_(0), buffer_send_size_(buffer_send_size), 27 32076 : session_(session) { 28 32079 : } 29 : 30 32090 : TcpMessageWriter::~TcpMessageWriter() { 31 32090 : for (BufferQueue::iterator iter = buffer_queue_.begin(); 32 32764 : iter != buffer_queue_.end(); ++iter) { 33 677 : DeleteBuffer(*iter); 34 : } 35 32084 : buffer_queue_.clear(); 36 32078 : } 37 : 38 1881790 : int TcpMessageWriter::AsyncSend(const uint8_t *data, size_t len, error_code *ec) { 39 : 40 1881790 : int write = len; 41 : 42 1881790 : if (buffer_queue_.empty()) { 43 1560812 : BufferAppend(data, len); 44 1560556 : if (session_->io_strand_) { 45 1560570 : boost::asio::detail::recycling_allocator<void> allocator; 46 3121496 : session_->io_strand_->post(bind(&TcpSession::AsyncWriteInternal, 47 3121042 : session_, TcpSessionPtr(session_)), allocator); 48 : } 49 : } else { 50 320996 : BufferAppend(data, len); 51 : } 52 : 53 1881823 : if ((GetBufferQueueSize() - offset_) > TcpMessageWriter::kMaxPendingBufferSize) { 54 32 : 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 1881175 : return write; 65 : } 66 : 67 1884546 : void TcpMessageWriter::TriggerAsyncWrite() { 68 : 69 : /* assert if there is an async write in progress */ 70 1884546 : assert(last_write_ == 0); 71 1884546 : assert(!buffer_queue_.empty()); 72 : 73 1884546 : boost::asio::mutable_buffer head = buffer_queue_.front(); 74 1884546 : size_t remaining = buffer_size(head) - offset_; 75 1884546 : last_write_ = min(buffer_send_size_, remaining); 76 : 77 : // Update socket write call statistics. 78 1884546 : session_->stats_.write_calls++; 79 1884546 : session_->server_->stats_.write_calls++; 80 : 81 1884546 : const uint8_t *data = buffer_cast<const uint8_t *>(head) + offset_; 82 1884546 : session_->AsyncWrite(data, last_write_); 83 1884546 : } 84 : 85 1883937 : bool TcpMessageWriter::UpdateBufferQueue(size_t wrote, bool *send_ready) { 86 : 87 1883937 : assert(last_write_ == wrote); 88 1883937 : assert(!buffer_queue_.empty()); 89 : 90 1883937 : bool more_write = true; 91 1883937 : last_write_ = 0; 92 1883937 : *send_ready = false; 93 : 94 1883937 : boost::asio::mutable_buffer head = buffer_queue_.front(); 95 1883937 : if ((offset_ + wrote) == buffer_size(head)) { 96 1881328 : offset_ = 0; 97 1881328 : DeleteBuffer(head); 98 1881328 : buffer_queue_.pop_front(); 99 : } else { 100 2609 : offset_ += wrote; 101 : } 102 : 103 1883937 : if (session_->write_blocked_ && ((GetBufferQueueSize() - offset_) < 104 1883937 : 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 1883937 : if (buffer_queue_.empty()) { 114 1560304 : buffer_queue_.clear(); 115 1560304 : more_write = false; 116 : } 117 : 118 1883937 : return more_write; 119 : } 120 : 121 1881798 : void TcpMessageWriter::BufferAppend(const uint8_t *src, int bytes) { 122 1881798 : uint8_t *data = new uint8_t[bytes]; 123 1881898 : memcpy(data, src, bytes); 124 1881898 : mutable_buffer buffer = mutable_buffer(data, bytes); 125 1881769 : buffer_queue_.push_back(buffer); 126 1881452 : } 127 : 128 1882005 : void TcpMessageWriter::DeleteBuffer(mutable_buffer buffer) { 129 1882005 : const uint8_t *data = buffer_cast<const uint8_t *>(buffer); 130 1882005 : delete[] data; 131 1882005 : return; 132 : } 133 :