Use boost::asio timer callbacks instead of self-made ones. There are still some termination issues (terminating a network game will result in a crash). However, it should be quite effective now with a large number of users.

This commit is contained in:
lotodore
2009-06-01 19:49:16 +00:00
parent cee2d787d0
commit 47648392fb
7 changed files with 124 additions and 146 deletions
+50 -18
View File
@@ -23,47 +23,79 @@
using namespace std;
TimerManager::TimerManager()
: m_curTimerId(0)
TimerManager::TimerManager(boost::shared_ptr<boost::asio::io_service> ioService)
: m_ioService(ioService), m_curTimerId(0)
{
}
unsigned
TimerManager::RegisterTimer(unsigned timeoutMsec, boost::function<void()> timerHandler, bool timerRepeat)
TimerManager::RegisterTimer(unsigned timeoutMsec, boost::function<void()> timerHandler, bool autoRestart)
{
// Register a new timer callback.
boost::recursive_mutex::scoped_lock lock(m_timerMutex);
// Use a unique id for each timer.
unsigned id = GetNextTimerId();
// Sort all timers by tick on which the callback occurs.
unsigned absoluteTimer = static_cast<unsigned>(m_softwareTimer.elapsed().total_milliseconds()) + timeoutMsec;
m_timerMap.insert(
TimerMap::value_type(absoluteTimer, TimerData(id, timeoutMsec, timerHandler, timerRepeat)));
TimerData data;
data.timer.reset(
new boost::asio::deadline_timer(*m_ioService, boost::posix_time::milliseconds(timeoutMsec)));
data.userHandler = timerHandler;
data.durationMsec = timeoutMsec;
data.autoRestart = autoRestart;
data.timer->async_wait(boost::bind(&TimerManager::Handler, boost::asio::placeholders::error, data));
m_timerMap.insert(TimerMap::value_type(id, data));
return id;
}
bool
TimerManager::RestartTimer(unsigned timerId, unsigned timeoutMsec, boost::function<void()> timerHandler)
{
boost::recursive_mutex::scoped_lock lock(m_timerMutex);
bool restarted = false;
TimerMap::iterator pos = m_timerMap.find(timerId);
if (pos != m_timerMap.end())
{
pos->second.userHandler = timerHandler;
boost::shared_ptr<boost::asio::deadline_timer> timer(pos->second.timer);
timer->cancel();
timer->expires_from_now(boost::posix_time::milliseconds(timeoutMsec));
timer->async_wait(boost::bind(&TimerManager::Handler, boost::asio::placeholders::error, pos->second));
restarted = true;
}
return restarted;
}
bool
TimerManager::UnregisterTimer(unsigned timerId)
{
// Remove a timer callback from the map.
boost::recursive_mutex::scoped_lock lock(m_timerMutex);
bool unregistered = false;
TimerMap::iterator i = m_timerMap.begin();
TimerMap::iterator end = m_timerMap.end();
while (i != end)
TimerMap::iterator pos = m_timerMap.find(timerId);
if (pos != m_timerMap.end())
{
if (i->second.id == timerId)
{
m_timerMap.erase(i);
unregistered = true;
break;
}
++i;
pos->second.timer->cancel();
m_timerMap.erase(pos);
unregistered = true;
}
return unregistered;
}
void
TimerManager::Handler(boost::system::error_code ec, TimerManager::TimerData data)
{
if (!ec)
{
data.userHandler();
if (data.autoRestart)
{
data.timer->expires_from_now(boost::posix_time::milliseconds(data.durationMsec));
data.timer->async_wait(boost::bind(&TimerManager::Handler, boost::asio::placeholders::error, data));
}
}
}
/*void
TimerManager::Process()
{
boost::recursive_mutex::scoped_lock lock(m_timerMutex);
@@ -101,7 +133,7 @@ TimerManager::Process()
}
}
} while (timerOccurred);
}
}*/
unsigned
TimerManager::GetNextTimerId()
+11 -14
View File
@@ -24,40 +24,37 @@
#include <map>
#include <boost/thread.hpp>
#include <boost/function.hpp>
#include <third_party/boost/timers.hpp>
#include <boost/asio.hpp>
class TimerManager
{
public:
TimerManager();
TimerManager(boost::shared_ptr<boost::asio::io_service> ioService);
unsigned RegisterTimer(unsigned timeoutMsec, boost::function<void()> timerHandler, bool timerRepeat = false);
unsigned RegisterTimer(unsigned timeoutMsec, boost::function<void()> timerHandler, bool autoRestart = false);
bool RestartTimer(unsigned timerId, unsigned timeoutMsec, boost::function<void()> timerHandler);
bool UnregisterTimer(unsigned timerId);
void Process();
protected:
struct TimerData
{
TimerData(unsigned timerId, unsigned timeoutMsec, boost::function<void()> timerHandler, bool timerRepeat)
: id(timerId), msec(timeoutMsec), handler(timerHandler), repeat(timerRepeat) {}
unsigned id;
unsigned msec;
boost::function<void()> handler;
bool repeat;
boost::shared_ptr<boost::asio::deadline_timer> timer;
boost::function<void()> userHandler;
unsigned durationMsec;
bool autoRestart;
};
typedef std::multimap<unsigned, TimerData> TimerMap;
static void Handler(boost::system::error_code ec, TimerData data);
typedef std::map<unsigned, TimerData> TimerMap;
unsigned GetNextTimerId();
private:
boost::shared_ptr<boost::asio::io_service> m_ioService;
mutable boost::recursive_mutex m_timerMutex;
TimerMap m_timerMap;
unsigned m_curTimerId;
boost::timers::portable::microsec_timer m_softwareTimer;
};
#endif
-1
View File
@@ -57,7 +57,6 @@ ServerGame::ServerGame(ServerLobbyThread &lobbyThread, u_int32_t id, const strin
ServerGame::~ServerGame()
{
GetLobbyThread().GetTimerManager().UnregisterTimer(m_removePlayerTimerId);
GetLobbyThread().GetTimerManager().UnregisterTimer(m_voteKickTimerId);
GetLobbyThread().GetTimerManager().UnregisterTimer(m_stateTimerId);
+32 -36
View File
@@ -488,10 +488,10 @@ ServerGameStateInit::TimerAdminWarning(ServerGame &server)
session.sessionData->GetSender().Send(session.sessionData, warning);
}
// Start timeout timer.
server.SetStateTimerId(
server.GetLobbyThread().GetTimerManager().RegisterTimer(
SERVER_GAME_ADMIN_WARNING_REMAINING_SEC * 1000,
boost::bind(&ServerGameStateInit::TimerAdminTimeout, this, boost::ref(server))));
server.GetLobbyThread().GetTimerManager().RestartTimer(
server.GetStateTimerId(),
SERVER_GAME_ADMIN_WARNING_REMAINING_SEC * 1000,
boost::bind(&ServerGameStateInit::TimerAdminTimeout, this, boost::ref(server)));
}
void
@@ -605,7 +605,6 @@ ServerGameStateStartGame::Enter(ServerGame &server)
void
ServerGameStateStartGame::Exit(ServerGame &server)
{
// The Id might be invalid, but this is not a problem.
server.GetLobbyThread().GetTimerManager().UnregisterTimer(server.GetStateTimerId());
server.SetStateTimerId(0);
}
@@ -722,7 +721,6 @@ ServerGameStateHand::Enter(ServerGame &server)
void
ServerGameStateHand::Exit(ServerGame &server)
{
// The Id might be invalid, but this is not a problem.
server.GetLobbyThread().GetTimerManager().UnregisterTimer(server.GetStateTimerId());
server.SetStateTimerId(0);
}
@@ -744,7 +742,6 @@ ServerGameStateHand::InternalProcessPacket(ServerGame &/*server*/, SessionWrappe
void
ServerGameStateHand::TimerLoop(ServerGame &server)
{
server.SetStateTimerId(0);
Game &curGame = server.GetGame();
// Main game loop.
@@ -792,19 +789,19 @@ ServerGameStateHand::TimerLoop(ServerGame &server)
server.SendToAllPlayers(allIn, SessionData::Game);
curGame.getCurrentHand()->setCardsShown(true);
server.SetStateTimerId(
server.GetLobbyThread().GetTimerManager().RegisterTimer(
SERVER_SHOW_CARDS_DELAY_SEC * 1000,
boost::bind(&ServerGameStateHand::TimerLoop, this, boost::ref(server))));
server.GetLobbyThread().GetTimerManager().RestartTimer(
server.GetStateTimerId(),
SERVER_SHOW_CARDS_DELAY_SEC * 1000,
boost::bind(&ServerGameStateHand::TimerLoop, this, boost::ref(server)));
}
else
{
SendNewRoundCards(server, curGame, newRound);
server.SetStateTimerId(
server.GetLobbyThread().GetTimerManager().RegisterTimer(
GetDealCardsDelaySec(server) * 1000,
boost::bind(&ServerGameStateHand::TimerLoop, this, boost::ref(server))));
server.GetLobbyThread().GetTimerManager().RestartTimer(
server.GetStateTimerId(),
GetDealCardsDelaySec(server) * 1000,
boost::bind(&ServerGameStateHand::TimerLoop, this, boost::ref(server)));
}
}
else
@@ -832,19 +829,19 @@ ServerGameStateHand::TimerLoop(ServerGame &server)
// If the player is computer controlled, let the engine act.
if (curPlayer->getMyType() == PLAYER_TYPE_COMPUTER)
{
server.SetStateTimerId(
server.GetLobbyThread().GetTimerManager().RegisterTimer(
SERVER_COMPUTER_ACTION_DELAY_SEC * 1000,
boost::bind(&ServerGameStateHand::TimerComputerAction, this, boost::ref(server))));
server.GetLobbyThread().GetTimerManager().RestartTimer(
server.GetStateTimerId(),
SERVER_COMPUTER_ACTION_DELAY_SEC * 1000,
boost::bind(&ServerGameStateHand::TimerComputerAction, this, boost::ref(server)));
}
// If the player we are waiting for left, continue without him.
else if (!server.GetSessionManager().IsPlayerConnected(curPlayer->getMyName()))
{
PerformPlayerAction(server, curPlayer, PLAYER_ACTION_FOLD, 0);
server.SetStateTimerId(
server.GetLobbyThread().GetTimerManager().RegisterTimer(
SERVER_LOOP_DELAY_MSEC,
boost::bind(&ServerGameStateHand::TimerLoop, this, boost::ref(server))));
server.GetLobbyThread().GetTimerManager().RestartTimer(
server.GetStateTimerId(),
SERVER_LOOP_DELAY_MSEC,
boost::bind(&ServerGameStateHand::TimerLoop, this, boost::ref(server)));
}
else
{
@@ -924,17 +921,17 @@ ServerGameStateHand::TimerLoop(ServerGame &server)
else if (playersWithCash.size() == 1)
{
// View a dialog for a new game - delayed.
server.SetStateTimerId(
server.GetLobbyThread().GetTimerManager().RegisterTimer(
SERVER_DELAY_NEXT_GAME_SEC * 1000,
boost::bind(&ServerGameStateHand::TimerNextGame, this, boost::ref(server))));
server.GetLobbyThread().GetTimerManager().RestartTimer(
server.GetStateTimerId(),
SERVER_DELAY_NEXT_GAME_SEC * 1000,
boost::bind(&ServerGameStateHand::TimerNextGame, this, boost::ref(server)));
}
else
{
server.SetStateTimerId(
server.GetLobbyThread().GetTimerManager().RegisterTimer(
SERVER_DELAY_NEXT_HAND_SEC * 1000,
boost::bind(&ServerGameStateHand::TimerNextHand, this, boost::ref(server))));
server.GetLobbyThread().GetTimerManager().RestartTimer(
server.GetStateTimerId(),
SERVER_DELAY_NEXT_HAND_SEC * 1000,
boost::bind(&ServerGameStateHand::TimerNextHand, this, boost::ref(server)));
}
}
}
@@ -946,10 +943,10 @@ ServerGameStateHand::TimerShowCards(ServerGame &server)
Game &curGame = server.GetGame();
SendNewRoundCards(server, curGame, curGame.getCurrentHand()->getCurrentRound());
server.SetStateTimerId(
server.GetLobbyThread().GetTimerManager().RegisterTimer(
GetDealCardsDelaySec(server) * 1000,
boost::bind(&ServerGameStateHand::TimerLoop, this, boost::ref(server))));
server.GetLobbyThread().GetTimerManager().RestartTimer(
server.GetStateTimerId(),
GetDealCardsDelaySec(server) * 1000,
boost::bind(&ServerGameStateHand::TimerLoop, this, boost::ref(server)));
}
void
@@ -1058,7 +1055,6 @@ ServerGameStateWaitPlayerAction::Enter(ServerGame &server)
void
ServerGameStateWaitPlayerAction::Exit(ServerGame &server)
{
// The Id might be invalid, but this is not a problem.
server.GetLobbyThread().GetTimerManager().UnregisterTimer(server.GetStateTimerId());
server.SetStateTimerId(0);
}
+29 -67
View File
@@ -1,5 +1,5 @@
/***************************************************************************
* Copyright (C) 2007 by Lothar May *
* Copyright (C) 2007-2009 by Lothar May *
* *
* This program is free software; you can redistribute it and/or modify *
* it under the terms of the GNU General Public License as published by *
@@ -43,6 +43,7 @@
#define SERVER_REMOVE_GAME_INTERVAL_MSEC 500
#define SERVER_REMOVE_PLAYER_INTERVAL_MSEC 100
#define SERVER_UPDATE_AVATAR_LOCK_INTERVAL_MSEC 1000
#define SERVER_PROCESS_SEND_INTERVAL_MSEC 10
#define SERVER_INIT_AVATAR_CLIENT_LOCK_SEC 30 // Forbid a client to send an additional avatar.
@@ -88,10 +89,11 @@ private:
ServerLobbyThread::ServerLobbyThread(GuiInterface &gui, ConfigFile *playerConfig, AvatarManager &avatarManager,
boost::shared_ptr<boost::asio::io_service> ioService)
: m_ioService(ioService), m_curBanId(0), m_gui(gui), m_avatarManager(avatarManager), m_playerConfig(playerConfig),
m_curGameId(0), m_curUniquePlayerId(0), m_curSessionId(INVALID_SESSION + 1),
: m_ioService(ioService), m_timerManager(ioService), m_curBanId(0), m_gui(gui), m_avatarManager(avatarManager),
m_playerConfig(playerConfig), m_curGameId(0), m_curUniquePlayerId(0), m_curSessionId(INVALID_SESSION + 1),
m_statDataChanged(false), m_startTime(boost::posix_time::second_clock::local_time())
{
m_work.reset(new boost::asio::io_service::work(*m_ioService));
m_senderCallback.reset(new ServerSenderCallback(*this));
m_sender.reset(new SenderHelper(*m_senderCallback, m_ioService));
m_receiver.reset(new ReceiverHelper);
@@ -118,6 +120,14 @@ ServerLobbyThread::Init(const string &pwd, const string &logDir)
}
}
void
ServerLobbyThread::SignalTermination()
{
Thread::SignalTermination();
m_work.reset();
m_ioService->stop();
}
void
ServerLobbyThread::AddConnection(boost::shared_ptr<tcp::socket> sock)
{
@@ -177,14 +187,16 @@ ServerLobbyThread::AddConnection(boost::shared_ptr<tcp::socket> sock)
void
ServerLobbyThread::ReAddSession(SessionWrapper session, int reason)
{
boost::shared_ptr<NetPacket> packet(new NetPacketRemovedFromGame);
NetPacketRemovedFromGame::Data removedData;
removedData.removeReason = reason;
static_cast<NetPacketRemovedFromGame *>(packet.get())->SetData(removedData);
session.sessionData->GetSender().Send(session.sessionData, packet);
if (session.sessionData.get() && session.playerData.get())
{
boost::shared_ptr<NetPacket> packet(new NetPacketRemovedFromGame);
NetPacketRemovedFromGame::Data removedData;
removedData.removeReason = reason;
static_cast<NetPacketRemovedFromGame *>(packet.get())->SetData(removedData);
session.sessionData->GetSender().Send(session.sessionData, packet);
boost::mutex::scoped_lock lock(m_sessionQueueMutex);
m_sessionQueue.push_back(session);
HandleReAddedSession(session);
}
}
void
@@ -225,8 +237,7 @@ ServerLobbyThread::CloseSession(SessionWrapper session)
void
ServerLobbyThread::ResubscribeLobbyMsg(SessionWrapper session)
{
boost::mutex::scoped_lock lock(m_resubscribeListMutex);
m_resubscribeList.push_back(session.sessionData->GetId());
InternalResubscribeMsg(session);
}
void
@@ -496,25 +507,12 @@ ServerLobbyThread::GetNextGameId()
void
ServerLobbyThread::Main()
{
boost::asio::io_service::work ioWork(*m_ioService);
try
{
// Register all timers.
RegisterTimers();
while (!ShouldTerminate())
{
// Process re-added sessions.
NewSessionLoop();
// Resubscribe Lobby Messages if needed.
ResubscribeLobbyMsgLoop();
// Process timers.
m_timerManager.Process();
// Process asio service.
m_ioService->poll();
m_sender->Process();
Thread::Msleep(10);
}
m_ioService->run();
} catch (const PokerTHException &e)
{
GetCallback().SignalNetServerError(e.GetErrorId(), e.GetOsErrorCode());
@@ -559,6 +557,11 @@ ServerLobbyThread::RegisterTimers()
SERVER_UPDATE_AVATAR_LOCK_INTERVAL_MSEC,
boost::bind(&ServerLobbyThread::TimerUpdateClientAvatarLock, this),
true);
// Check if new data needs to be sent.
m_timerManager.RegisterTimer(
SERVER_PROCESS_SEND_INTERVAL_MSEC,
boost::bind(&SenderInterface::Process, m_sender),
true);
}
void
@@ -572,7 +575,6 @@ ServerLobbyThread::HandleRead(SessionId sessionId, const boost::system::error_co
{
if (!error)
{
ReceiveBuffer &buf = session.sessionData->GetReceiveBuffer();
buf.recvBufUsed += bytesRead;
GetReceiver().ScanPackets(buf);
@@ -1031,23 +1033,6 @@ ServerLobbyThread::RequestPlayerAvatar(SessionWrapper session)
session.sessionData->GetSender().Send(session.sessionData, retrieveAvatar);
}
void
ServerLobbyThread::NewSessionLoop()
{
// Handle one incoming session at a time.
SessionWrapper tmpSession;
{
boost::mutex::scoped_lock lock(m_sessionQueueMutex);
if (!m_sessionQueue.empty())
{
tmpSession = m_sessionQueue.front();
m_sessionQueue.pop_front();
}
}
if (tmpSession.sessionData.get() && tmpSession.playerData.get())
HandleReAddedSession(tmpSession);
}
void
ServerLobbyThread::TimerRemoveGame()
{
@@ -1084,29 +1069,6 @@ ServerLobbyThread::TimerRemovePlayer()
}
}
void
ServerLobbyThread::ResubscribeLobbyMsgLoop()
{
boost::mutex::scoped_lock lock(m_resubscribeListMutex);
if (!m_resubscribeList.empty())
{
SessionIdList::iterator i = m_resubscribeList.begin();
SessionIdList::iterator end = m_resubscribeList.end();
while (i != end)
{
SessionWrapper tmpSession = m_gameSessionManager.GetSessionById(*i);
if (!tmpSession.sessionData.get())
tmpSession = m_sessionManager.GetSessionById(*i);
if (tmpSession.sessionData.get())
InternalResubscribeMsg(tmpSession);
++i;
}
m_resubscribeList.clear();
}
}
void
ServerLobbyThread::TimerUpdateClientAvatarLock()
{
-1
View File
@@ -153,7 +153,6 @@ private:
ConfigFile *m_playerConfig;
unsigned m_gameNum;
unsigned m_curPetitionId;
unsigned m_removePlayerTimerId;
unsigned m_voteKickTimerId;
unsigned m_stateTimerId;
+2 -9
View File
@@ -53,6 +53,7 @@ public:
virtual ~ServerLobbyThread();
void Init(const std::string &pwd, const std::string &logDir);
virtual void SignalTermination();
void AddConnection(boost::shared_ptr<boost::asio::ip::tcp::socket> sock);
void ReAddSession(SessionWrapper session, int reason);
@@ -126,10 +127,8 @@ protected:
void HandleNetPacketJoinGame(SessionWrapper session, const NetPacketJoinGame &tmpPacket);
void EstablishSession(SessionWrapper session);
void RequestPlayerAvatar(SessionWrapper session);
void NewSessionLoop();
void TimerRemoveGame();
void TimerRemovePlayer();
void ResubscribeLobbyMsgLoop();
void TimerUpdateClientAvatarLock();
void TimerCheckSessionTimeouts();
void TimerCleanupAvatarCache();
@@ -173,12 +172,9 @@ protected:
private:
boost::shared_ptr<boost::asio::io_service> m_ioService;
boost::shared_ptr<boost::asio::io_service::work> m_work;
TimerManager m_timerManager;
SessionQueue m_sessionQueue;
mutable boost::mutex m_sessionQueueMutex;
SessionManager m_sessionManager;
SessionManager m_gameSessionManager;
@@ -194,9 +190,6 @@ private:
PlayerDataMap m_computerPlayers;
mutable boost::mutex m_computerPlayersMutex;
SessionIdList m_resubscribeList;
mutable boost::mutex m_resubscribeListMutex;
RegexMap m_banPlayerNameMap;
IPAddressMap m_banIPAddressMap;
unsigned m_curBanId;