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 2335 : TcpMessageWriter::TcpMessageWriter(TcpSession *session, 25 2335 : size_t buffer_send_size) : 26 2329 : offset_(0), last_write_(0), buffer_send_size_(buffer_send_size), 27 2335 : session_(session) { 28 2329 : } 29 : 30 2361 : TcpMessageWriter::~TcpMessageWriter() { 31 2361 : for (BufferQueue::iterator iter = buffer_queue_.begin(); 32 2363 : iter != buffer_queue_.end(); ++iter) { 33 2 : DeleteBuffer(*iter); 34 : } 35 2361 : buffer_queue_.clear(); 36 2361 : } 37 : 38 7525 : int TcpMessageWriter::AsyncSend(const uint8_t *data, size_t len, error_code *ec) { 39 : 40 7525 : int write = len; 41 : 42 7525 : if (buffer_queue_.empty()) { 43 7099 : BufferAppend(data, len); 44 7053 : if (session_->io_strand_) { 45 7052 : boost::asio::detail::recycling_allocator<void> allocator; 46 14167 : session_->io_strand_->post(bind(&TcpSession::AsyncWriteInternal, 47 14173 : session_, TcpSessionPtr(session_)), allocator); 48 : } 49 : } else { 50 423 : BufferAppend(data, len); 51 : } 52 : 53 7551 : 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 7531 : return write; 65 : } 66 : 67 9784 : void TcpMessageWriter::TriggerAsyncWrite() { 68 : 69 : /* assert if there is an async write in progress */ 70 9784 : assert(last_write_ == 0); 71 9784 : assert(!buffer_queue_.empty()); 72 : 73 9784 : boost::asio::mutable_buffer head = buffer_queue_.front(); 74 9784 : size_t remaining = buffer_size(head) - offset_; 75 9784 : last_write_ = min(buffer_send_size_, remaining); 76 : 77 : // Update socket write call statistics. 78 9784 : session_->stats_.write_calls++; 79 9784 : session_->server_->stats_.write_calls++; 80 : 81 9784 : const uint8_t *data = buffer_cast<const uint8_t *>(head) + offset_; 82 9784 : session_->AsyncWrite(data, last_write_); 83 9784 : } 84 : 85 9783 : bool TcpMessageWriter::UpdateBufferQueue(size_t wrote, bool *send_ready) { 86 : 87 9783 : assert(last_write_ == wrote); 88 9783 : assert(!buffer_queue_.empty()); 89 : 90 9783 : bool more_write = true; 91 9783 : last_write_ = 0; 92 9783 : *send_ready = false; 93 : 94 9783 : boost::asio::mutable_buffer head = buffer_queue_.front(); 95 9783 : if ((offset_ + wrote) == buffer_size(head)) { 96 7552 : offset_ = 0; 97 7552 : DeleteBuffer(head); 98 7552 : buffer_queue_.pop_front(); 99 : } else { 100 2231 : offset_ += wrote; 101 : } 102 : 103 9783 : if (session_->write_blocked_ && ((GetBufferQueueSize() - offset_) < 104 9783 : 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 9783 : if (buffer_queue_.empty()) { 114 7126 : buffer_queue_.clear(); 115 7126 : more_write = false; 116 : } 117 : 118 9783 : return more_write; 119 : } 120 : 121 7518 : void TcpMessageWriter::BufferAppend(const uint8_t *src, int bytes) { 122 7518 : uint8_t *data = new uint8_t[bytes]; 123 7521 : memcpy(data, src, bytes); 124 7521 : mutable_buffer buffer = mutable_buffer(data, bytes); 125 7530 : buffer_queue_.push_back(buffer); 126 7479 : } 127 : 128 7554 : void TcpMessageWriter::DeleteBuffer(mutable_buffer buffer) { 129 7554 : const uint8_t *data = buffer_cast<const uint8_t *>(buffer); 130 7554 : delete[] data; 131 7554 : return; 132 : } 133 :