From 82a1be763b86fdfb121ef9e8ccbd45ba8d906300 Mon Sep 17 00:00:00 2001 From: lotodore Date: Wed, 14 Nov 2007 22:11:04 +0000 Subject: [PATCH] Use several sender threads with load balancing for sending the avatars. Remember last failed send operation (not yet a list, just the last) and cancel any more low prio sends to that sender. --- src/net/common/senderthread.cpp | 55 +++++++++++++++++----------- src/net/common/serverlobbythread.cpp | 19 ++++++++-- src/net/senderthread.h | 4 ++ src/net/serverlobbythread.h | 3 ++ 4 files changed, 57 insertions(+), 24 deletions(-) diff --git a/src/net/common/senderthread.cpp b/src/net/common/senderthread.cpp index d60e39cf..6afb3b1a 100644 --- a/src/net/common/senderthread.cpp +++ b/src/net/common/senderthread.cpp @@ -21,6 +21,7 @@ #include #include #include +#include #include #include @@ -29,8 +30,7 @@ using namespace std; -#define SEND_ERROR_NORMAL_TIMEOUT_MSEC 20000 -#define SEND_ERROR_LOW_PRIO_TIMEOUT_MSEC 10000 +#define SEND_ERROR_TIMEOUT_MSEC 20000 #define SEND_TIMEOUT_MSEC 10 #ifdef POKERTH_DEDICATED_SERVER #define SEND_QUEUE_SIZE 10000 @@ -41,7 +41,7 @@ using namespace std; #endif SenderThread::SenderThread(SenderCallback &cb) -: m_tmpOutBufSize(0), m_tmpIsLowPrio(false), m_callback(cb) +: m_tmpOutBufSize(0), m_lastInvalidSessionId(INVALID_SESSION), m_callback(cb) { } @@ -89,6 +89,27 @@ SenderThread::SendLowPrio(boost::shared_ptr session, const NetPacke } } +unsigned +SenderThread::GetNumPacketsInQueue() const +{ + unsigned numPackets; + { + boost::mutex::scoped_lock lock(m_lowPrioOutBufMutex); + numPackets = m_lowPrioOutBuf.size(); + } + { + boost::mutex::scoped_lock lock(m_outBufMutex); + numPackets += m_outBuf.size(); + } + return numPackets; +} + +bool +SenderThread::operator<(const SenderThread &other) const +{ + return GetNumPacketsInQueue() < other.GetNumPacketsInQueue(); +} + void SenderThread::InternalStore(SendDataDeque &sendQueue, unsigned maxQueueSize, boost::shared_ptr session, boost::shared_ptr packet) { @@ -133,6 +154,7 @@ SenderThread::Main() // For reasons of simplicity, only one packet is sent at a time. if (!m_tmpOutBufSize) { + bool isLowPrio = false; SendData tmpData; // Check main queue first. { @@ -141,7 +163,6 @@ SenderThread::Main() { tmpData = m_outBuf.front(); m_outBuf.pop_front(); - m_tmpIsLowPrio = false; } } @@ -153,13 +174,13 @@ SenderThread::Main() { tmpData = m_lowPrioOutBuf.front(); m_lowPrioOutBuf.pop_front(); - m_tmpIsLowPrio = true; + isLowPrio = true; } } - if (tmpData.first.get()) + if (tmpData.first.get() && tmpData.second.get()) { - if (tmpData.second.get()) + if (!isLowPrio || tmpData.second->GetId() != m_lastInvalidSessionId) { u_int16_t tmpLen = tmpData.first->GetLen(); if (tmpLen <= MAX_PACKET_SIZE) @@ -243,21 +264,13 @@ SenderThread::Main() // Check whether the send timed out. if (sendTimer.is_running()) { - bool doRemove = false; - - unsigned msec = static_cast(sendTimer.elapsed().total_milliseconds()); - if (m_tmpIsLowPrio) - { - if (msec > SEND_ERROR_LOW_PRIO_TIMEOUT_MSEC) - doRemove = true; - } - else - { - if (msec > SEND_ERROR_NORMAL_TIMEOUT_MSEC) - doRemove = true; - } - if (doRemove) + if (sendTimer.elapsed().total_milliseconds() > SEND_ERROR_TIMEOUT_MSEC) { + if (m_curSession.get()) + { + m_lastInvalidSessionId = m_curSession->GetId(); + LOG_MSG("Send operation for session " << m_lastInvalidSessionId << " timed out."); + } RemoveCurSendData(); sendTimer.reset(); } diff --git a/src/net/common/serverlobbythread.cpp b/src/net/common/serverlobbythread.cpp index 5ca89e9d..240b6af0 100644 --- a/src/net/common/serverlobbythread.cpp +++ b/src/net/common/serverlobbythread.cpp @@ -29,14 +29,18 @@ #include #include +#include #include #include +#include #define SERVER_MAX_NUM_SESSIONS 512 // Maximum number of idle users in lobby. #define SERVER_CACHE_CLEANUP_INTERVAL_SEC 86400 // 1 day #define SERVER_SAVE_STATISTICS_INTERVAL_SEC 60 #define SERVER_INIT_SESSION_TIMEOUT_SEC 20 +#define SERVER_NUM_AVATAR_SENDER_THREADS 5 + #define SERVER_COMPUTER_PLAYER_NAME "Computer" #define SERVER_STATISTICS_FILE_NAME "server_statistics.log" @@ -73,6 +77,8 @@ ServerLobbyThread::ServerLobbyThread(GuiInterface &gui, ConfigFile *playerConfig { m_senderCallback.reset(new ServerSenderCallback(*this)); m_sender.reset(new SenderThread(GetSenderCallback())); + for (int i = 0; i < SERVER_NUM_AVATAR_SENDER_THREADS; i++) + m_avatarSenderThreadPool.push_back(boost::shared_ptr(new SenderThread(GetSenderCallback()))); m_receiver.reset(new ReceiverHelper); } @@ -273,6 +279,7 @@ void ServerLobbyThread::Main() { GetSender().Run(); + for_each(m_avatarSenderThreadPool.begin(), m_avatarSenderThreadPool.end(), boost::mem_fn(&SenderThread::Run)); try { @@ -302,8 +309,10 @@ ServerLobbyThread::Main() TerminateGames(); GetSender().SignalTermination(); - if (!GetSender().Join(SENDER_THREAD_TERMINATE_TIMEOUT)) - LOG_ERROR("Fatal error: Unable to terminated Sender Thread in Lobby."); + for_each(m_avatarSenderThreadPool.begin(), m_avatarSenderThreadPool.end(), boost::mem_fn(&SenderThread::SignalTermination)); + + GetSender().Join(SENDER_THREAD_TERMINATE_TIMEOUT); + for_each(m_avatarSenderThreadPool.begin(), m_avatarSenderThreadPool.end(), boost::bind(&SenderThread::Join, _1, SENDER_THREAD_TERMINATE_TIMEOUT)); CleanupConnectQueue(); } @@ -564,7 +573,11 @@ ServerLobbyThread::HandleNetPacketRetrieveAvatar(SessionWrapper session, const N if (GetAvatarManager().AvatarFileToNetPackets(tmpFile, request.requestId, tmpPackets) == 0) { avatarFound = true; - GetSender().SendLowPrio(session.sessionData, tmpPackets); + SenderThreadList::iterator pos = min_element(m_avatarSenderThreadPool.begin(), m_avatarSenderThreadPool.end(), *boost::lambda::_1 < *boost::lambda::_2); + if (pos != m_avatarSenderThreadPool.end()) + (*pos)->SendLowPrio(session.sessionData, tmpPackets); + else + LOG_ERROR("Load balancing for avatar sender threads failed."); } else LOG_ERROR("Failed to read avatar file for network transmission."); diff --git a/src/net/senderthread.h b/src/net/senderthread.h index f050c598..0beda89d 100644 --- a/src/net/senderthread.h +++ b/src/net/senderthread.h @@ -43,6 +43,9 @@ public: void SendLowPrio(boost::shared_ptr session, boost::shared_ptr packet); void SendLowPrio(boost::shared_ptr session, const NetPacketList &packetList); + unsigned GetNumPacketsInQueue() const; + bool operator<(const SenderThread &other) const; + protected: typedef std::pair, boost::shared_ptr > SendData; typedef std::deque SendDataDeque; @@ -68,6 +71,7 @@ private: char m_tmpOutBuf[MAX_PACKET_SIZE]; unsigned m_tmpOutBufSize; bool m_tmpIsLowPrio; + unsigned m_lastInvalidSessionId; SenderCallback &m_callback; }; diff --git a/src/net/serverlobbythread.h b/src/net/serverlobbythread.h index b1670120..9f5c6919 100644 --- a/src/net/serverlobbythread.h +++ b/src/net/serverlobbythread.h @@ -84,6 +84,7 @@ protected: typedef std::map InitTimerSessionMap; typedef std::map > GameMap; typedef std::list RemoveGameList; + typedef std::list > SenderThreadList; // Main function of the thread. virtual void Main(); @@ -166,6 +167,7 @@ private: boost::shared_ptr m_receiver; boost::shared_ptr m_sender; + SenderThreadList m_avatarSenderThreadPool; boost::shared_ptr m_senderCallback; GuiInterface &m_gui; AvatarManager &m_avatarManager; @@ -186,6 +188,7 @@ private: boost::timers::portable::microsec_timer m_cacheCleanupTimer; boost::timers::portable::microsec_timer m_saveStatisticsTimer; + boost::timers::portable::microsec_timer m_uptimeTimer; }; #endif