From 2541f28ff249afea527aae9c33b7edb225ebc862 Mon Sep 17 00:00:00 2001 From: lotodore Date: Sun, 10 May 2009 17:53:37 +0000 Subject: [PATCH] Sender thread is no longer a thread (part of boost::asio integration). Warning: Network games are still broken! --- src/net/common/clientthread.cpp | 7 +- src/net/common/senderthread.cpp | 125 +++++++++++---------------- src/net/common/serverlobbythread.cpp | 20 +++-- src/net/senderinterface.h | 6 +- src/net/senderthread.h | 13 +-- src/net/serverlobbythread.h | 2 +- 6 files changed, 71 insertions(+), 102 deletions(-) diff --git a/src/net/common/clientthread.cpp b/src/net/common/clientthread.cpp index f7ef5ec3..9b9d5f48 100644 --- a/src/net/common/clientthread.cpp +++ b/src/net/common/clientthread.cpp @@ -378,12 +378,12 @@ void ClientThread::Main() { // Start sub-threads. - m_senderThread->Start(); m_avatarDownloader.reset(new DownloaderThread); m_avatarDownloader->Run(); SetState(CLIENT_INITIAL_STATE::Instance()); // Main loop. + boost::asio::io_service::work ioWork(*m_ioService); try { while (!ShouldTerminate()) @@ -414,6 +414,9 @@ ClientThread::Main() } if (IsSessionEstablished()) SendPacketLoop(); + m_ioService->poll(); + GetContext().GetSessionData()->GetSender().Process(); + Thread::Msleep(10); } } catch (const PokerTHException &e) { @@ -422,8 +425,6 @@ ClientThread::Main() // Terminate sub-threads. m_avatarDownloader->SignalTermination(); m_avatarDownloader->Join(DOWNLOADER_THREAD_TERMINATE_TIMEOUT); - m_senderThread->SignalStop(); - m_senderThread->WaitStop(); } void diff --git a/src/net/common/senderthread.cpp b/src/net/common/senderthread.cpp index dfb6c5ba..05cd941e 100644 --- a/src/net/common/senderthread.cpp +++ b/src/net/common/senderthread.cpp @@ -100,24 +100,6 @@ SenderThread::~SenderThread() { } -void -SenderThread::Start() -{ - Run(); -} - -void -SenderThread::SignalStop() -{ - SignalTermination(); -} - -void -SenderThread::WaitStop() -{ - Join(SENDER_THREAD_TERMINATE_TIMEOUT); -} - void SenderThread::Send(boost::shared_ptr session, boost::shared_ptr packet) { @@ -190,73 +172,66 @@ SenderThread::SignalSessionTerminated(unsigned sessionId) } void -SenderThread::Main() +SenderThread::Process() { - boost::asio::io_service::work ioWork(*m_ioService); - while (!ShouldTerminate()) + // Close sessions if they were destructed. { - // Close sessions if they were destructed. + boost::mutex::scoped_lock lock(m_removedSessionsMutex); + if (!m_removedSessions.empty()) { - boost::mutex::scoped_lock lock(m_removedSessionsMutex); - if (!m_removedSessions.empty()) + SessionIdList newRemovedSessions; + SessionIdList::iterator i = m_removedSessions.begin(); + SessionIdList::iterator end = m_removedSessions.end(); + + boost::mutex::scoped_lock lock(m_sendQueueMapMutex); + + while (i != end) { - SessionIdList newRemovedSessions; - SessionIdList::iterator i = m_removedSessions.begin(); - SessionIdList::iterator end = m_removedSessions.end(); - - boost::mutex::scoped_lock lock(m_sendQueueMapMutex); - - while (i != end) + SendQueueMap::iterator pos = m_sendQueueMap.find(*i); + if (pos != m_sendQueueMap.end()) { - SendQueueMap::iterator pos = m_sendQueueMap.find(*i); - if (pos != m_sendQueueMap.end()) + // Remove session if no write is in progress, else wait. + bool shouldDelete; { - // Remove session if no write is in progress, else wait. - bool shouldDelete; - { - boost::mutex::scoped_lock lock(pos->second->dataMutex); - shouldDelete = (!pos->second->writeInProgress && pos->second->list.empty()); - } - if (shouldDelete) - m_sendQueueMap.erase(pos); - else - newRemovedSessions.push_back(*i); + boost::mutex::scoped_lock lock(pos->second->dataMutex); + shouldDelete = (!pos->second->writeInProgress && pos->second->list.empty()); } - ++i; + if (shouldDelete) + m_sendQueueMap.erase(pos); + else + newRemovedSessions.push_back(*i); } - m_removedSessions = newRemovedSessions; + ++i; + } + m_removedSessions = newRemovedSessions; + } + } + // Iterate through all changed sessions, and send data if needed. + bool sessionValid; + do + { + sessionValid = false; + unsigned sessionId = 0; + + { + boost::mutex::scoped_lock lock(m_changedSessionsMutex); + if (!m_changedSessions.empty()) + { + sessionId = m_changedSessions.front(); + m_changedSessions.pop_front(); + sessionValid = true; } } - // Iterate through all changed sessions, and send data if needed. - bool sessionValid; - do + boost::shared_ptr tmpManager; + if (sessionValid) { - sessionValid = false; - unsigned sessionId = 0; - - { - boost::mutex::scoped_lock lock(m_changedSessionsMutex); - if (!m_changedSessions.empty()) - { - sessionId = m_changedSessions.front(); - m_changedSessions.pop_front(); - sessionValid = true; - } - } - boost::shared_ptr tmpManager; - if (sessionValid) - { - boost::mutex::scoped_lock lock(m_sendQueueMapMutex); - SendQueueMap::iterator pos = m_sendQueueMap.find(sessionId); - if (pos != m_sendQueueMap.end()) - tmpManager = pos->second; - } - if (tmpManager) - tmpManager->AsyncSendNextPacket(); - } while (sessionValid); - - m_ioService->poll(); - Msleep(SEND_TIMEOUT_MSEC); - } + boost::mutex::scoped_lock lock(m_sendQueueMapMutex); + SendQueueMap::iterator pos = m_sendQueueMap.find(sessionId); + if (pos != m_sendQueueMap.end()) + tmpManager = pos->second; + } + if (tmpManager) + tmpManager->AsyncSendNextPacket(); + } while (sessionValid); } diff --git a/src/net/common/serverlobbythread.cpp b/src/net/common/serverlobbythread.cpp index f5b4f4cd..36c02db8 100644 --- a/src/net/common/serverlobbythread.cpp +++ b/src/net/common/serverlobbythread.cpp @@ -42,6 +42,7 @@ #define SERVER_CHECK_SESSION_TIMEOUTS_INTERVAL_MSEC 500 #define SERVER_REMOVE_GAME_INTERVAL_MSEC 100 #define SERVER_REMOVE_PLAYER_INTERVAL_MSEC 100 +#define SERVER_UPDATE_AVATAR_LOCK_INTERVAL_MSEC 1000 #define SERVER_INIT_AVATAR_CLIENT_LOCK_SEC 30 // Forbid a client to send an additional avatar. @@ -493,24 +494,23 @@ ServerLobbyThread::GetNextGameId() void ServerLobbyThread::Main() { + boost::asio::io_service::work ioWork(*m_ioService); try { // Register all timers. RegisterTimers(); - // Start send thread. - m_sender->Start(); - while (!ShouldTerminate()) { // Process re-added sessions. NewSessionLoop(); // Resubscribe Lobby Messages if needed. ResubscribeLobbyMsgLoop(); - // Update avatar limitation lock. - UpdateAvatarClientTimerLoop(); // Process timers. m_timerManager.Process(); + // Process asio service. + m_ioService->poll(); + m_sender->Process(); Thread::Msleep(10); } } catch (const PokerTHException &e) @@ -524,9 +524,6 @@ ServerLobbyThread::Main() // Remove all sessions. m_gameSessionManager.Clear(); m_sessionManager.Clear(); - // Stop sender thread. - m_sender->SignalStop(); - m_sender->WaitStop(); } void @@ -557,6 +554,11 @@ ServerLobbyThread::RegisterTimers() SERVER_SAVE_STATISTICS_INTERVAL_SEC * 1000, boost::bind(&ServerLobbyThread::TimerSaveStatisticsFile, this), true); + // Update the avatar upload locks. + m_timerManager.RegisterTimer( + SERVER_UPDATE_AVATAR_LOCK_INTERVAL_MSEC, + boost::bind(&ServerLobbyThread::TimerUpdateClientAvatarLock, this), + true); } void @@ -1112,7 +1114,7 @@ ServerLobbyThread::ResubscribeLobbyMsgLoop() } void -ServerLobbyThread::UpdateAvatarClientTimerLoop() +ServerLobbyThread::TimerUpdateClientAvatarLock() { boost::mutex::scoped_lock lock(m_timerAvatarClientAddressMapMutex); diff --git a/src/net/senderinterface.h b/src/net/senderinterface.h index 9773fa5c..b2e391f1 100644 --- a/src/net/senderinterface.h +++ b/src/net/senderinterface.h @@ -29,13 +29,11 @@ class SenderInterface public: virtual ~SenderInterface(); - virtual void Start() = 0; - virtual void SignalStop() = 0; - virtual void WaitStop() = 0; - virtual void Send(boost::shared_ptr session, boost::shared_ptr packet) = 0; virtual void Send(boost::shared_ptr session, const NetPacketList &packetList) = 0; + virtual void Process() = 0; + virtual void SignalSessionTerminated(unsigned sessionId) = 0; }; diff --git a/src/net/senderthread.h b/src/net/senderthread.h index db79af27..cd43eb84 100644 --- a/src/net/senderthread.h +++ b/src/net/senderthread.h @@ -21,28 +21,24 @@ #ifndef _SENDERTHREAD_H_ #define _SENDERTHREAD_H_ -#include #include #include #include class SessionData; class SendDataManager; -#define SENDER_THREAD_TERMINATE_TIMEOUT THREAD_WAIT_INFINITE -class SenderThread : public Thread, public SenderInterface +class SenderThread : public SenderInterface { public: SenderThread(SenderCallback &cb, boost::shared_ptr ioService); virtual ~SenderThread(); - virtual void Start(); - virtual void SignalStop(); - virtual void WaitStop(); - virtual void Send(boost::shared_ptr session, boost::shared_ptr packet); virtual void Send(boost::shared_ptr session, const NetPacketList &packetList); + virtual void Process(); + virtual void SignalSessionTerminated(unsigned sessionId); protected: @@ -50,9 +46,6 @@ protected: typedef std::map > SendQueueMap; - // Main function of the thread. - virtual void Main(); - private: SenderCallback &m_callback; diff --git a/src/net/serverlobbythread.h b/src/net/serverlobbythread.h index 2e81d938..5a1cdad7 100644 --- a/src/net/serverlobbythread.h +++ b/src/net/serverlobbythread.h @@ -131,7 +131,7 @@ protected: void TimerRemoveGame(); void TimerRemovePlayer(); void ResubscribeLobbyMsgLoop(); - void UpdateAvatarClientTimerLoop(); + void TimerUpdateClientAvatarLock(); void TimerCheckSessionTimeouts(); void TimerCleanupAvatarCache();