From 46a136bc8b0d3cf0019b254b6e2f8ba4e42fe7d5 Mon Sep 17 00:00:00 2001 From: lotodore Date: Fri, 18 Feb 2011 18:36:30 +0000 Subject: [PATCH] Tremendously improving the speed of the network send code (ticket #52). The code has also been simplified, and uses some kind of double buffering as in graphics programming. There are now a lot less memory allocations and a lot less copying. However, memory usage is also increased. --- src/net/common/senderhelper.cpp | 133 ++++++++++++++++++++------------ 1 file changed, 84 insertions(+), 49 deletions(-) diff --git a/src/net/common/senderhelper.cpp b/src/net/common/senderhelper.cpp index 18382080..f903a315 100644 --- a/src/net/common/senderhelper.cpp +++ b/src/net/common/senderhelper.cpp @@ -20,12 +20,12 @@ #include #include #include +#include #include #include #include #include -#include #include #include #include @@ -34,80 +34,103 @@ using namespace std; using boost::asio::ip::tcp; -#define SEND_ERROR_TIMEOUT_MSEC 20000 -#define SEND_TIMEOUT_MSEC 10 -#define SEND_QUEUE_SIZE 10000000 -#define SEND_LOG_INTERVAL_SEC 60 +#define SEND_BUF_FIRST_ALLOC_CHUNKSIZE 4096 +#define MAX_SEND_BUF_SIZE SEND_BUF_FIRST_ALLOC_CHUNKSIZE * 1000 -typedef std::list > SendDataList; - class SendDataManager : public boost::enable_shared_from_this { public: SendDataManager() - : sendBuf(NULL), writeInProgress(false) { + : sendBuf(NULL), curWriteBuf(NULL), sendBufAllocated(0), sendBufUsed(0), + curWriteBufAllocated(0), curWriteBufUsed(0) + { } - ~SendDataManager() { - delete[] sendBuf; + ~SendDataManager() + { + free(sendBuf); + free(curWriteBuf); + } + + inline size_t GetSendBufLeft() const + { + int bytesLeft = sendBufAllocated - sendBufUsed; + return bytesLeft < 0 ? (size_t)0 : (size_t)bytesLeft; + } + + inline size_t GetAllocated() const + { + return sendBufAllocated; + } + + inline bool ReallocSendBuf() + { + bool retVal = false; + size_t allocAmount = sendBufAllocated * 2; + if (0 == allocAmount) { + allocAmount = (size_t)SEND_BUF_FIRST_ALLOC_CHUNKSIZE; + } + char *tempBuf = (char *)realloc(sendBuf, allocAmount); + if (tempBuf) { + sendBuf = tempBuf; + sendBufAllocated = allocAmount; + retVal = true; + } + return retVal; + } + + inline void AppendToSendBufWithoutCheck(const char *data, size_t size) + { + memcpy(sendBuf + sendBufUsed, data, size); + sendBufUsed += size; } void HandleWrite(boost::shared_ptr socket, const boost::system::error_code &error); - void AsyncSendNextPacket(boost::shared_ptr socket, bool handlerMode = false); + void AsyncSendNextPacket(boost::shared_ptr socket); mutable boost::mutex dataMutex; - SendDataList list; + +private: char *sendBuf; - bool writeInProgress; + char *curWriteBuf; + size_t sendBufAllocated; + size_t sendBufUsed; + size_t curWriteBufAllocated; + size_t curWriteBufUsed; }; void SendDataManager::HandleWrite(boost::shared_ptr socket, const boost::system::error_code &error) { - if (!error) - AsyncSendNextPacket(socket, true); + if (!error) { + // Successfully sent the data. + curWriteBufUsed = 0; + // Send more data, if available. + AsyncSendNextPacket(socket); + } } void -SendDataManager::AsyncSendNextPacket(boost::shared_ptr socket, bool handlerMode) +SendDataManager::AsyncSendNextPacket(boost::shared_ptr socket) { boost::mutex::scoped_lock lock(dataMutex); - if (!writeInProgress || handlerMode) { - delete[] sendBuf; - sendBuf = NULL; - - unsigned bufPos = 0; - unsigned bufSize = 0; - // Count required bytes. - SendDataList::iterator i = list.begin(); - SendDataList::iterator end = list.end(); - while (i != end) { - bufSize += (*i)->GetSize(); - ++i; - } - if (bufSize) { - sendBuf = new char[bufSize]; - i = list.begin(); - end = list.end(); - while (i != end) { - memcpy(sendBuf + bufPos, (*i)->GetData(), (*i)->GetSize()); - bufPos += (*i)->GetSize(); - ++i; - } - list.clear(); + if (!curWriteBufUsed) { + // Swap buffers and send data. + boost::swap(curWriteBuf, sendBuf); + boost::swap(curWriteBufAllocated, sendBufAllocated); + boost::swap(curWriteBufUsed, sendBufUsed); + if (curWriteBufUsed) { boost::asio::async_write( *socket, - boost::asio::buffer(sendBuf, bufSize), + boost::asio::buffer(curWriteBuf, curWriteBufUsed), boost::bind(&SendDataManager::HandleWrite, shared_from_this(), socket, boost::asio::placeholders::error)); - writeInProgress = true; - } else - writeInProgress = false; + } } } @@ -136,7 +159,7 @@ SenderHelper::Send(boost::shared_ptr session, boost::shared_ptrdataMutex); - if (tmpManager->list.size() < SEND_QUEUE_SIZE) { + if (tmpManager->GetAllocated() <= MAX_SEND_BUF_SIZE) { InternalStorePacket(*tmpManager, packet); } } @@ -163,7 +186,7 @@ SenderHelper::Send(boost::shared_ptr session, const NetPacketList & { // Second: Add packets to specific queue. boost::mutex::scoped_lock lock(tmpManager->dataMutex); - if (tmpManager->list.size() + packetList.size() <= SEND_QUEUE_SIZE) { + if (tmpManager->GetAllocated() <= MAX_SEND_BUF_SIZE) { NetPacketList::const_iterator i = packetList.begin(); NetPacketList::const_iterator end = packetList.end(); while (i != end) { @@ -190,15 +213,27 @@ SenderHelper::SignalSessionTerminated(unsigned sessionId) m_sendQueueMap.erase(pos); } +static int der_encode_to_buffer_cb(const void *data, size_t size, void *arg) { + SendDataManager *m = (SendDataManager *)arg; + + // Realloc buffer if necessary. + while (m->GetSendBufLeft() < size) { + if (!m->ReallocSendBuf()) { + return -1; + } + } + + m->AppendToSendBufWithoutCheck((const char*)data, size); + + return 0; +} + void SenderHelper::InternalStorePacket(SendDataManager &tmpManager, boost::shared_ptr packet) { - unsigned char buf[MAX_PACKET_SIZE]; - asn_enc_rval_t e = der_encode_to_buffer(&asn_DEF_PokerTHMessage, packet->GetMsg(), buf, MAX_PACKET_SIZE); + asn_enc_rval_t e = der_encode(&asn_DEF_PokerTHMessage, packet->GetMsg(), der_encode_to_buffer_cb, &tmpManager); //cerr << "OUT:" << endl << packet->ToString() << endl; if (e.encoded == -1) LOG_ERROR("Failed to encode NetPacket: " << packet->GetMsg()->present); - else - tmpManager.list.push_back(boost::shared_ptr(new EncodedPacket(buf, e.encoded))); }