One sender thread per client. Requires loads of resources. Kind of hacked, because the internal dependencies are very ugly now. Requires loads of testing. A problem might be the last packet sent by the server to a specific client before closing the connection, for example an error message. It might never even reach the client, because the sender thread is terminated before. I have not yet found a good way to fix this.

This commit is contained in:
lotodore
2009-01-07 21:52:16 +00:00
parent 606158c4f0
commit d72cdb6873
16 changed files with 150 additions and 325 deletions
+4
View File
@@ -24,6 +24,9 @@
#include <net/netcontext.h>
#include <net/receivebuffer.h>
#include <net/sessiondata.h>
#include <boost/shared_ptr.hpp>
class SenderCallback;
class ClientContext : public NetContext
{
@@ -105,6 +108,7 @@ private:
std::string m_cacheDir;
bool m_hasSubscribedLobbyMsg;
ReceiveBuffer m_receiveBuffer;
boost::shared_ptr<SenderCallback> m_senderCallback;
};
#endif
-3
View File
@@ -114,7 +114,6 @@ protected:
ClientState &GetState();
void SetState(ClientState &newState);
SenderThread &GetSender();
ReceiverHelper &GetReceiver();
void SetGameId(unsigned id);
@@ -167,12 +166,10 @@ private:
mutable boost::mutex m_outPacketListMutex;
boost::shared_ptr<ClientContext> m_context;
boost::shared_ptr<ClientSenderCallback> m_senderCallback;
ClientState *m_curState;
GuiInterface &m_gui;
AvatarManager &m_avatarManager;
boost::shared_ptr<SenderThread> m_sender;
boost::shared_ptr<ReceiverHelper> m_receiver;
GameData m_gameData;
+17 -1
View File
@@ -18,12 +18,28 @@
***************************************************************************/
#include <net/clientcontext.h>
#include <net/sendercallback.h>
class ClientSenderCallback : public SenderCallback
{
public:
ClientSenderCallback() {}
virtual ~ClientSenderCallback() {}
virtual void SignalNetError(SessionId /*session*/, int /*errorID*/, int /*osErrorID*/)
{
}
private:
};
ClientContext::ClientContext()
: m_protocol(0), m_addrFamily(AF_INET), m_useServerList(false), m_serverPort(0),
m_hasSubscribedLobbyMsg(true)
{
bzero(&m_clientSockaddr, sizeof(m_clientSockaddr));
m_senderCallback.reset(new ClientSenderCallback());
}
ClientContext::~ClientContext()
@@ -40,7 +56,7 @@ ClientContext::GetSocket() const
void
ClientContext::SetSocket(SOCKET sockfd)
{
m_sessionData.reset(new SessionData(sockfd, SESSION_ID_GENERIC));
m_sessionData.reset(new SessionData(sockfd, SESSION_ID_GENERIC, *m_senderCallback));
}
boost::shared_ptr<SessionData>
+3 -3
View File
@@ -627,7 +627,7 @@ ClientStateStartSession::Process(ClientThread &client)
boost::shared_ptr<NetPacket> packet(new NetPacketInit);
((NetPacketInit *)packet.get())->SetData(initData);
client.GetSender().Send(context.GetSessionData(), packet);
context.GetSessionData()->GetSender().Send(context.GetSessionData(), packet);
client.SetState(ClientStateWaitSession::Instance());
@@ -903,7 +903,7 @@ ClientStateWaitSession::InternalProcess(ClientThread &client, boost::shared_ptr<
tmpList);
if (!avatarError)
client.GetSender().Send(client.GetContext().GetSessionData(), tmpList);
client.GetContext().GetSessionData()->GetSender().Send(client.GetContext().GetSessionData(), tmpList);
else
throw ClientException(__FILE__, __LINE__, avatarError, 0);
}
@@ -1055,7 +1055,7 @@ ClientStateSynchronizeStart::Process(ClientThread &client)
{
// Acknowledge start.
boost::shared_ptr<NetPacket> startAck(new NetPacketStartEventAck);
client.GetSender().Send(client.GetContext().GetSessionData(), startAck);
client.GetContext().GetSessionData()->GetSender().Send(client.GetContext().GetSessionData(), startAck);
// Unsubscribe lobby messages.
client.UnsubscribeLobbyMsg();
+5 -43
View File
@@ -39,31 +39,11 @@
using namespace std;
class ClientSenderCallback : public SenderCallback
{
public:
ClientSenderCallback(ClientThread &client) : m_client(client) {}
virtual ~ClientSenderCallback() {}
virtual void SignalNetError(SessionId /*session*/, int errorID, int osErrorID)
{
// Just signal the error.
// We assume that the client thread will be terminated.
m_client.GetCallback().SignalNetClientError(errorID, osErrorID);
}
private:
ClientThread &m_client;
};
ClientThread::ClientThread(GuiInterface &gui, AvatarManager &avatarManager)
: m_curState(NULL), m_gui(gui), m_avatarManager(avatarManager),
m_curGameId(0), m_curGameNum(1), m_guiPlayerId(0), m_sessionEstablished(false)
{
m_context.reset(new ClientContext);
m_senderCallback.reset(new ClientSenderCallback(*this));
m_sender.reset(new SenderThread(GetSenderCallback()));
m_receiver.reset(new ReceiverHelper);
myQtToolsInterface.reset(CreateQtToolsWrapper());
}
@@ -351,8 +331,6 @@ ClientThread::Main()
{
SetState(CLIENT_INITIAL_STATE::Instance());
GetSender().Run();
try
{
while (!ShouldTerminate())
@@ -388,8 +366,6 @@ ClientThread::Main()
{
GetCallback().SignalNetClientError(e.GetErrorId(), e.GetOsErrorCode());
}
GetSender().SignalTermination();
GetSender().Join(SENDER_THREAD_TERMINATE_TIMEOUT);
}
void
@@ -411,7 +387,7 @@ ClientThread::SendPacketLoop()
while (i != end)
{
GetSender().Send(GetContext().GetSessionData(), *i);
GetContext().GetSessionData()->GetSender().Send(GetContext().GetSessionData(), *i);
++i;
}
m_outPacketList.clear();
@@ -442,7 +418,7 @@ ClientThread::RequestPlayerInfo(unsigned id, bool requestAvatar)
NetPacketRetrievePlayerInfo::Data reqData;
reqData.playerId = id;
static_cast<NetPacketRetrievePlayerInfo *>(req.get())->SetData(reqData);
GetSender().Send(GetContext().GetSessionData(), req);
GetContext().GetSessionData()->GetSender().Send(GetContext().GetSessionData(), req);
m_playerInfoRequestList.push_back(id);
@@ -546,7 +522,7 @@ ClientThread::RetrieveAvatarIfNeeded(unsigned id, const PlayerInfo &info)
retrieveAvatarData.requestId = id;
retrieveAvatarData.avatar = info.avatar;
static_cast<NetPacketRetrieveAvatar *>(retrieveAvatar.get())->SetData(retrieveAvatarData);
GetSender().Send(GetContext().GetSessionData(), retrieveAvatar);
GetContext().GetSessionData()->GetSender().Send(GetContext().GetSessionData(), retrieveAvatar);
}
}
}
@@ -621,7 +597,7 @@ ClientThread::UnsubscribeLobbyMsg()
{
// Send unsubscribe request.
boost::shared_ptr<NetPacket> unsubscr(new NetPacketUnsubscribeGameList);
GetSender().Send(GetContext().GetSessionData(), unsubscr);
GetContext().GetSessionData()->GetSender().Send(GetContext().GetSessionData(), unsubscr);
GetContext().SetSubscribeLobbyMsg(false);
}
}
@@ -635,7 +611,7 @@ ClientThread::ResubscribeLobbyMsg()
ClearGameInfoMap();
// Send resubscribe request.
boost::shared_ptr<NetPacket> resubscr(new NetPacketResubscribeGameList);
GetSender().Send(GetContext().GetSessionData(), resubscr);
GetContext().GetSessionData()->GetSender().Send(GetContext().GetSessionData(), resubscr);
GetContext().SetSubscribeLobbyMsg(true);
}
}
@@ -667,13 +643,6 @@ ClientThread::SetState(ClientState &newState)
m_curState = &newState;
}
SenderThread &
ClientThread::GetSender()
{
assert(m_sender.get());
return *m_sender;
}
ReceiverHelper &
ClientThread::GetReceiver()
{
@@ -737,13 +706,6 @@ ClientThread::GetGame()
return m_game;
}
ClientSenderCallback &
ClientThread::GetSenderCallback()
{
assert(m_senderCallback.get());
return *m_senderCallback;
}
QtToolsInterface &
ClientThread::GetQtToolsInterface()
{
+48 -173
View File
@@ -25,7 +25,6 @@
#include <cstring>
#include <cassert>
#include <third_party/boost/timers.hpp>
#include <boost/bind.hpp>
using namespace std;
@@ -36,7 +35,7 @@ using namespace std;
#define SEND_LOG_INTERVAL_SEC 60
SenderThread::SenderThread(SenderCallback &cb)
: m_callback(cb)
: m_callback(cb), m_bytesSent(0)
{
}
@@ -67,17 +66,16 @@ SenderThread::Send(boost::shared_ptr<SessionData> session, boost::shared_ptr<Net
{
if (packet.get() && session.get())
{
boost::mutex::scoped_lock lock0(m_sessionsStalledMutex);
if (find(m_sessionsStalled.begin(), m_sessionsStalled.end(), session->GetId()) == m_sessionsStalled.end())
{
boost::mutex::scoped_lock lock1(m_sendQueueMutex);
InternalStore(m_sendQueue, SEND_QUEUE_SIZE, session, packet);
boost::mutex::scoped_lock lock(m_sendQueueMutex);
if (m_sendQueue.size() < SEND_QUEUE_SIZE)
{
m_sendQueue.push_back(packet);
}
}
else
{
// This session is stalled.
boost::mutex::scoped_lock lock2(m_stalledQueueMutex);
InternalStore(m_stalledQueue, SEND_QUEUE_SIZE, session, packet);
boost::mutex::scoped_lock lock(m_sessionMutex);
m_session = session;
}
}
}
@@ -87,135 +85,54 @@ SenderThread::Send(boost::shared_ptr<SessionData> session, const NetPacketList &
{
if (!packetList.empty() && session.get())
{
boost::mutex::scoped_lock lock0(m_sessionsStalledMutex);
if (find(m_sessionsStalled.begin(), m_sessionsStalled.end(), session->GetId()) == m_sessionsStalled.end())
boost::mutex::scoped_lock lock(m_sendQueueMutex);
if (m_sendQueue.size() + packetList.size() <= SEND_QUEUE_SIZE)
{
boost::mutex::scoped_lock lock1(m_sendQueueMutex);
InternalStore(m_sendQueue, SEND_QUEUE_SIZE, session, packetList);
NetPacketList::const_iterator i = packetList.begin();
NetPacketList::const_iterator end = packetList.end();
while (i != end)
{
m_sendQueue.push_back(*i);
++i;
}
}
else
{
// This session is stalled.
boost::mutex::scoped_lock lock2(m_stalledQueueMutex);
InternalStore(m_stalledQueue, SEND_QUEUE_SIZE, session, packetList);
boost::mutex::scoped_lock lock(m_sessionMutex);
m_session = session;
}
}
}
unsigned
SenderThread::GetNumPacketsInQueue() const
{
unsigned numPackets;
{
boost::mutex::scoped_lock lock1(m_sendQueueMutex);
numPackets = m_sendQueue.size();
}
{
boost::mutex::scoped_lock lock2(m_stalledQueueMutex);
numPackets += m_stalledQueue.size();
}
return numPackets;
}
bool
SenderThread::operator<(const SenderThread &other) const
{
return GetNumPacketsInQueue() < other.GetNumPacketsInQueue();
}
void
SenderThread::InternalStore(SendDataList &sendQueue, unsigned maxQueueSize, boost::shared_ptr<SessionData> session, boost::shared_ptr<NetPacket> packet)
{
if (sendQueue.size() < maxQueueSize) // Queue is limited in size.
sendQueue.push_back(SendData(packet, session));
// TODO: Throw exception if failed.
}
void
SenderThread::InternalStore(SendDataList &sendQueue, unsigned maxQueueSize, boost::shared_ptr<SessionData> session, const NetPacketList &packetList)
{
if (sendQueue.size() + packetList.size() <= maxQueueSize)
{
NetPacketList::const_iterator i = packetList.begin();
NetPacketList::const_iterator end = packetList.end();
while (i != end)
{
sendQueue.push_back(SendData(*i, session));
++i;
}
}
// TODO: Throw exception if failed.
}
void
SenderThread::Main()
{
while (!ShouldTerminate())
{
/*
* To prevent stalling of the sender, keeping the order
* of the packets in the sender queue is not guaranteed.
* Instead, only the order of the packets for a single
* target session is maintained.
*
* This could also be done by using one sender thread
* for each session, but that would require too many
* resources.
*
* The send queue is a list of packets. Whenever a
* select timeout occurs, the session of the current
* packet is placed in the stalled list, and all other
* packets from the queue for that session are also
* attached to the stalled list. This process is
* continued with the next packet in the send queue.
*
* When the send queue is empty, the list of stalled
* packets is copied back to the send queue, and
* everything starts from the beginning.
*
* No packets are lost in this algorithm, and at the same
* time, the send process is never fully stalled,
* except when only packets for stalled sessions are
* present.
*
* Note: If someone keeps putting in packets to send, the
* stalled packets will never be sent, but this is
* considered more a theoretical problem.
*/
// For reasons of simplicity, only one packet is sent at a time.
SendData tmpData;
// Check main queue.
if (!m_curPacket)
{
boost::mutex::scoped_lock lock0(m_sessionsStalledMutex);
boost::mutex::scoped_lock lock1(m_sendQueueMutex);
if (m_sendQueue.empty())
{
// Check stalled queue.
// Attention: TRIPLE lock (on purpose).
boost::mutex::scoped_lock lock2(m_stalledQueueMutex);
if (!m_stalledQueue.empty())
{
m_sendQueue.swap(m_stalledQueue);
m_sessionsStalled.clear(); // No more sessions stalled, all in send list.
}
}
boost::mutex::scoped_lock lock(m_sendQueueMutex);
if (!m_sendQueue.empty())
{
tmpData = m_sendQueue.front();
m_curPacket = m_sendQueue.front();
m_sendQueue.pop_front();
}
}
if (tmpData.packet && tmpData.session)
if (m_curPacket)
{
const unsigned tmpLen = tmpData.packet->GetLen();
if (tmpLen <= MAX_PACKET_SIZE)
const unsigned tmpLen = m_curPacket->GetLen();
if (tmpLen > MAX_PACKET_SIZE)
m_curPacket.reset(); // TODO log
else
{
SOCKET tmpSocket = tmpData.session->GetSocket();
SOCKET tmpSocket;
{
boost::mutex::scoped_lock lock(m_sessionMutex);
tmpSocket = m_session->GetSocket();
}
// send next chunk of data
int bytesSent = send(tmpSocket, ((const char *)tmpData.packet->GetRawData()) + tmpData.bytesSent, tmpLen - tmpData.bytesSent, SOCKET_SEND_FLAGS);
int bytesSent = send(tmpSocket, ((const char *)m_curPacket->GetRawData()) + m_bytesSent, tmpLen - m_bytesSent, SOCKET_SEND_FLAGS);
if (!IS_VALID_SEND(bytesSent))
{
@@ -241,41 +158,12 @@ SenderThread::Main()
// Skip this packet - this is bad, and is therefore reported.
// Ignore invalid or not connected sockets.
if (errCode != SOCKET_ERR_NOTCONN && errCode != SOCKET_ERR_NOTSOCK)
m_callback.SignalNetError(tmpData.session->GetId(), ERR_SOCK_SELECT_FAILED, errCode);
}
Msleep(SEND_TIMEOUT_MSEC);
}
else if (selectResult == 0)
{
// A timeout occured - don't block the thread.
{
// Attention: TRIPLE lock (on purpose).
// Stall all packets for that sender.
boost::mutex::scoped_lock lock0(m_sessionsStalledMutex);
boost::mutex::scoped_lock lock1(m_sendQueueMutex);
boost::mutex::scoped_lock lock2(m_stalledQueueMutex);
m_stalledQueue.push_back(tmpData);
m_sessionsStalled.push_back(tmpData.session->GetId());
SendDataList::iterator i = m_sendQueue.begin();
SendDataList::iterator end = m_sendQueue.end();
while (i != end)
{
SendDataList::iterator next = i;
++next;
if ((*i).session && (*i).session->GetId() == tmpData.session->GetId())
{
m_stalledQueue.push_back(*i);
m_sendQueue.erase(i);
}
i = next;
boost::mutex::scoped_lock lock(m_sessionMutex);
m_callback.SignalNetError(m_session->GetId(), ERR_SOCK_SELECT_FAILED, errCode);
}
}
}
else
{
// Select was successful - store the packet back in the main queue.
boost::mutex::scoped_lock lock1(m_sendQueueMutex);
m_sendQueue.push_front(tmpData);
Msleep(SEND_TIMEOUT_MSEC);
}
}
else // other errors than would block
@@ -283,45 +171,32 @@ SenderThread::Main()
// Skip this packet - this is bad, and is therefore reported.
// Ignore invalid or not connected sockets.
if (errCode != SOCKET_ERR_NOTCONN && errCode != SOCKET_ERR_NOTSOCK)
m_callback.SignalNetError(tmpData.session->GetId(), ERR_SOCK_SEND_FAILED, errCode);
{
boost::mutex::scoped_lock lock(m_sessionMutex);
m_callback.SignalNetError(m_session->GetId(), ERR_SOCK_SEND_FAILED, errCode);
}
Msleep(SEND_TIMEOUT_MSEC);
}
}
else if ((unsigned)bytesSent + tmpData.bytesSent < tmpLen)
else if ((unsigned)bytesSent + m_bytesSent < tmpLen)
{
if (bytesSent)
{
tmpData.bytesSent += bytesSent;
// Send was partly successful - store the packet back in the main queue.
boost::mutex::scoped_lock lock1(m_sendQueueMutex);
m_sendQueue.push_front(tmpData);
// Send was partly successful.
m_bytesSent += bytesSent;
}
else
Msleep(SEND_TIMEOUT_MSEC);
}
else //if ((unsigned)bytesSent + m_bytesSent == tmpLen)
{
m_curPacket.reset();
m_bytesSent = 0;
}
}
else
Msleep(SEND_TIMEOUT_MSEC);
}
else
Msleep(SEND_TIMEOUT_MSEC);
if (m_logTimer.elapsed().total_seconds() >= SEND_LOG_INTERVAL_SEC)
{
unsigned sendQueueSize = 0;
unsigned stalledQueueSize = 0;
{
boost::mutex::scoped_lock lock1(m_sendQueueMutex);
sendQueueSize = m_sendQueue.size();
}
{
boost::mutex::scoped_lock lock2(m_stalledQueueMutex);
stalledQueueSize = m_stalledQueue.size();
}
LOG_VERBOSE("Sender TICK - send queue " << sendQueueSize << ", stalled " << stalledQueueSize << ".");
m_logTimer.reset();
m_logTimer.start();
}
}
}
+5 -5
View File
@@ -323,7 +323,7 @@ ServerGameStateInit::HandleNewSession(ServerGameThread &server, SessionWrapper s
joinGameAckData.prights = session.playerData->GetRights();
joinGameAckData.gameData = server.GetGameData();
static_cast<NetPacketJoinGameAck *>(joinGameAck.get())->SetData(joinGameAckData);
server.GetSender().Send(session.sessionData, joinGameAck);
session.sessionData->GetSender().Send(session.sessionData, joinGameAck);
// Send notifications for connected players to client.
PlayerDataList tmpPlayerList = server.GetFullPlayerDataList();
@@ -331,7 +331,7 @@ ServerGameStateInit::HandleNewSession(ServerGameThread &server, SessionWrapper s
PlayerDataList::iterator player_end = tmpPlayerList.end();
while (player_i != player_end)
{
server.GetSender().Send(session.sessionData, CreateNetPacketPlayerJoined(*(*player_i)));
session.sessionData->GetSender().Send(session.sessionData, CreateNetPacketPlayerJoined(*(*player_i)));
++player_i;
}
@@ -364,7 +364,7 @@ ServerGameStateInit::Process(ServerGameThread &server)
warningData.timeoutReason = NETWORK_TIMEOUT_GAME_ADMIN_IDLE;
warningData.remainingSeconds = SERVER_GAME_ADMIN_WARNING_REMAINING_SEC;
static_cast<NetPacketTimeoutWarning *>(warning.get())->SetData(warningData);
server.GetSender().Send(session.sessionData, warning);
session.sessionData->GetSender().Send(session.sessionData, warning);
}
}
else if (server.GetStateTimer().elapsed().total_seconds() >= SERVER_GAME_ADMIN_TIMEOUT_SEC)
@@ -638,7 +638,7 @@ ServerGameStateStartHand::Process(ServerGameThread &server)
handStartData.smallBlind = curGame.getCurrentHand()->getSmallBlind();
static_cast<NetPacketHandStart *>(notifyCards.get())->SetData(handStartData);
server.GetSender().Send(tmpPlayer->getNetSessionData(), notifyCards);
tmpPlayer->getNetSessionData()->GetSender().Send(tmpPlayer->getNetSessionData(), notifyCards);
}
++i;
}
@@ -1021,7 +1021,7 @@ ServerGameStateWaitPlayerAction::InternalProcess(ServerGameThread &server, Sessi
rejectData.playerBet = actionData.playerBet;
rejectData.rejectionReason = code;
static_cast<NetPacketPlayersActionRejected *>(reject.get())->SetData(rejectData);
server.GetSender().Send(session.sessionData, reject);
session.sessionData->GetSender().Send(session.sessionData, reject);
}
}
+5 -11
View File
@@ -91,7 +91,7 @@ ServerGameThread::GetCurRound() const
void
ServerGameThread::SendToAllPlayers(boost::shared_ptr<NetPacket> packet, SessionData::State state)
{
GetSessionManager().SendToAllSessions(GetSender(), packet, state);
GetSessionManager().SendToAllSessions(packet, state);
}
void
@@ -379,7 +379,7 @@ ServerGameThread::InternalDenyAskVoteKick(SessionWrapper byWhom, unsigned player
denyPetitionData.playerId = playerIdWho;
denyPetitionData.denyReason = reason;
static_cast<NetPacketAskKickPlayerDenied *>(denyPetition.get())->SetData(denyPetitionData);
GetSender().Send(byWhom.sessionData, denyPetition);
byWhom.sessionData->GetSender().Send(byWhom.sessionData, denyPetition);
}
void
@@ -429,7 +429,7 @@ ServerGameThread::InternalDenyVoteKick(SessionWrapper byWhom, unsigned petitionI
denyVoteData.petitionId = petitionId;
denyVoteData.denyReason = reason;
static_cast<NetPacketVoteKickPlayerDenied *>(denyVote.get())->SetData(denyVoteData);
GetSender().Send(byWhom.sessionData, denyVote);
byWhom.sessionData->GetSender().Send(byWhom.sessionData, denyVote);
}
PlayerDataList
@@ -611,7 +611,7 @@ ServerGameThread::RemovePlayerData(boost::shared_ptr<PlayerData> player, int rea
NetPacketGameAdminChanged::Data adminChangedData;
adminChangedData.playerId = newAdmin->GetUniqueId(); // Choose next player as admin.
static_cast<NetPacketGameAdminChanged *>(adminChanged.get())->SetData(adminChangedData);
GetSessionManager().SendToAllSessions(GetSender(), adminChanged, SessionData::Game);
GetSessionManager().SendToAllSessions(adminChanged, SessionData::Game);
GetLobbyThread().NotifyGameAdminChanged(GetId(), newAdmin->GetUniqueId());
}
@@ -625,7 +625,7 @@ ServerGameThread::RemovePlayerData(boost::shared_ptr<PlayerData> player, int rea
thisPlayerLeftData.playerId = player->GetUniqueId();
thisPlayerLeftData.removeReason = reason;
static_cast<NetPacketPlayerLeft *>(thisPlayerLeft.get())->SetData(thisPlayerLeftData);
GetSessionManager().SendToAllSessions(GetSender(), thisPlayerLeft, SessionData::Game);
GetSessionManager().SendToAllSessions(thisPlayerLeft, SessionData::Game);
GetLobbyThread().NotifyPlayerLeftGame(GetId(), player->GetUniqueId());
}
@@ -772,12 +772,6 @@ ServerGameThread::SetStateTimerFlag(unsigned flag)
m_stateTimerFlag = flag;
}
SenderInterface &
ServerGameThread::GetSender()
{
return GetLobbyThread().GetSender();
}
ReceiverHelper &
ServerGameThread::GetReceiver()
{
+31 -44
View File
@@ -81,7 +81,6 @@ ServerLobbyThread::ServerLobbyThread(GuiInterface &gui, ConfigFile *playerConfig
m_statDataChanged(false), m_startTime(boost::posix_time::second_clock::local_time())
{
m_senderCallback.reset(new ServerSenderCallback(*this));
m_sender.reset(new SenderThread(GetSenderCallback()));
m_receiver.reset(new ReceiverHelper);
}
@@ -121,7 +120,7 @@ ServerLobbyThread::ReAddSession(SessionWrapper session, int reason)
NetPacketRemovedFromGame::Data removedData;
removedData.removeReason = reason;
static_cast<NetPacketRemovedFromGame *>(packet.get())->SetData(removedData);
GetSender().Send(session.sessionData, packet);
session.sessionData->GetSender().Send(session.sessionData, packet);
boost::mutex::scoped_lock lock(m_sessionQueueMutex);
m_sessionQueue.push_back(session);
@@ -175,8 +174,8 @@ ServerLobbyThread::NotifyPlayerJoinedGame(unsigned gameId, unsigned playerId)
packetData.gameId = gameId;
packetData.playerId = playerId;
static_cast<NetPacketGameListPlayerJoined *>(packet.get())->SetData(packetData);
m_sessionManager.SendLobbyMsgToAllSessions(GetSender(), packet, SessionData::Established);
m_gameSessionManager.SendLobbyMsgToAllSessions(GetSender(), packet, SessionData::Game);
m_sessionManager.SendLobbyMsgToAllSessions(packet, SessionData::Established);
m_gameSessionManager.SendLobbyMsgToAllSessions(packet, SessionData::Game);
}
void
@@ -188,8 +187,8 @@ ServerLobbyThread::NotifyPlayerLeftGame(unsigned gameId, unsigned playerId)
packetData.gameId = gameId;
packetData.playerId = playerId;
static_cast<NetPacketGameListPlayerLeft *>(packet.get())->SetData(packetData);
m_sessionManager.SendLobbyMsgToAllSessions(GetSender(), packet, SessionData::Established);
m_gameSessionManager.SendLobbyMsgToAllSessions(GetSender(), packet, SessionData::Game);
m_sessionManager.SendLobbyMsgToAllSessions(packet, SessionData::Established);
m_gameSessionManager.SendLobbyMsgToAllSessions(packet, SessionData::Game);
}
void
@@ -201,24 +200,24 @@ ServerLobbyThread::NotifyGameAdminChanged(unsigned gameId, unsigned newAdminPlay
packetData.gameId = gameId;
packetData.newAdminplayerId = newAdminPlayerId;
static_cast<NetPacketGameListAdminChanged *>(packet.get())->SetData(packetData);
m_sessionManager.SendLobbyMsgToAllSessions(GetSender(), packet, SessionData::Established);
m_gameSessionManager.SendLobbyMsgToAllSessions(GetSender(), packet, SessionData::Game);
m_sessionManager.SendLobbyMsgToAllSessions(packet, SessionData::Established);
m_gameSessionManager.SendLobbyMsgToAllSessions(packet, SessionData::Game);
}
void
ServerLobbyThread::NotifyStartingGame(unsigned gameId)
{
boost::shared_ptr<NetPacket> packet = CreateNetPacketGameListUpdate(gameId, GAME_MODE_STARTED);
m_sessionManager.SendLobbyMsgToAllSessions(GetSender(), packet, SessionData::Established);
m_gameSessionManager.SendLobbyMsgToAllSessions(GetSender(), packet, SessionData::Game);
m_sessionManager.SendLobbyMsgToAllSessions(packet, SessionData::Established);
m_gameSessionManager.SendLobbyMsgToAllSessions(packet, SessionData::Game);
}
void
ServerLobbyThread::NotifyReopeningGame(unsigned gameId)
{
boost::shared_ptr<NetPacket> packet = CreateNetPacketGameListUpdate(gameId, GAME_MODE_CREATED);
m_sessionManager.SendLobbyMsgToAllSessions(GetSender(), packet, SessionData::Established);
m_gameSessionManager.SendLobbyMsgToAllSessions(GetSender(), packet, SessionData::Game);
m_sessionManager.SendLobbyMsgToAllSessions(packet, SessionData::Established);
m_gameSessionManager.SendLobbyMsgToAllSessions(packet, SessionData::Game);
}
void
@@ -267,7 +266,7 @@ ServerLobbyThread::SendGlobalChat(const string &message)
outChatData.playerId = 0;
outChatData.text = message;
static_cast<NetPacketChatText *>(outChat.get())->SetData(outChatData);
m_gameSessionManager.SendToAllSessions(GetSender(), outChat, SessionData::Game);
m_gameSessionManager.SendToAllSessions(outChat, SessionData::Game);
}
void
@@ -277,7 +276,7 @@ ServerLobbyThread::SendGlobalMsgBox(const string &message)
NetPacketMsgBoxText::Data outMsgData;
outMsgData.text = message;
static_cast<NetPacketMsgBoxText *>(outMsg.get())->SetData(outMsgData);
m_gameSessionManager.SendToAllSessions(GetSender(), outMsg, SessionData::Game);
m_gameSessionManager.SendToAllSessions(outMsg, SessionData::Game);
}
void
@@ -344,8 +343,6 @@ ServerLobbyThread::GetNextGameId()
void
ServerLobbyThread::Main()
{
GetSender().Start();
try
{
while (!ShouldTerminate())
@@ -379,9 +376,6 @@ ServerLobbyThread::Main()
TerminateGames();
GetSender().SignalStop();
GetSender().WaitStop();
CleanupConnectQueue();
}
@@ -642,7 +636,7 @@ ServerLobbyThread::HandleNetPacketRetrievePlayerInfo(SessionWrapper session, con
if (infoData.playerInfo.hasAvatar)
infoData.playerInfo.avatar = tmpPlayer->GetAvatarMD5();
static_cast<NetPacketPlayerInfo *>(info.get())->SetData(infoData);
GetSender().Send(session.sessionData, info);
session.sessionData->GetSender().Send(session.sessionData, info);
}
else
{
@@ -651,7 +645,7 @@ ServerLobbyThread::HandleNetPacketRetrievePlayerInfo(SessionWrapper session, con
NetPacketUnknownPlayerId::Data unknownData;
unknownData.playerId = request.playerId;
static_cast<NetPacketUnknownPlayerId *>(unknown.get())->SetData(unknownData);
GetSender().Send(session.sessionData, unknown);
session.sessionData->GetSender().Send(session.sessionData, unknown);
}
}
@@ -669,7 +663,7 @@ ServerLobbyThread::HandleNetPacketRetrieveAvatar(SessionWrapper session, const N
if (GetAvatarManager().AvatarFileToNetPackets(tmpFile, request.requestId, tmpPackets) == 0)
{
avatarFound = true;
GetSender().Send(session.sessionData, tmpPackets);
session.sessionData->GetSender().Send(session.sessionData, tmpPackets);
}
else
LOG_ERROR("Failed to read avatar file for network transmission.");
@@ -682,7 +676,7 @@ ServerLobbyThread::HandleNetPacketRetrieveAvatar(SessionWrapper session, const N
NetPacketUnknownAvatar::Data unknownData;
unknownData.requestId = request.requestId;
static_cast<NetPacketUnknownAvatar *>(unknown.get())->SetData(unknownData);
GetSender().Send(session.sessionData, unknown);
session.sessionData->GetSender().Send(session.sessionData, unknown);
}
}
@@ -755,7 +749,7 @@ ServerLobbyThread::EstablishSession(SessionWrapper session)
initAckData.sessionId = session.sessionData->GetId(); // TODO: currently unused.
initAckData.playerId = session.playerData->GetUniqueId();
static_cast<NetPacketInitAck *>(initAck.get())->SetData(initAckData);
GetSender().Send(session.sessionData, initAck);
session.sessionData->GetSender().Send(session.sessionData, initAck);
// Send the game list to the client.
SendGameList(session.sessionData);
@@ -787,7 +781,7 @@ ServerLobbyThread::RequestPlayerAvatar(SessionWrapper session)
retrieveAvatarData.requestId = session.playerData->GetUniqueId();
retrieveAvatarData.avatar = session.playerData->GetAvatarMD5();
static_cast<NetPacketRetrieveAvatar *>(retrieveAvatar.get())->SetData(retrieveAvatarData);
GetSender().Send(session.sessionData, retrieveAvatar);
session.sessionData->GetSender().Send(session.sessionData, retrieveAvatar);
}
void
@@ -942,8 +936,8 @@ ServerLobbyThread::InternalAddGame(boost::shared_ptr<ServerGameThread> game)
// Add game to list.
m_gameMap.insert(GameMap::value_type(game->GetId(), game));
// Notify all players.
m_sessionManager.SendLobbyMsgToAllSessions(GetSender(), CreateNetPacketGameListNew(*game), SessionData::Established);
m_gameSessionManager.SendLobbyMsgToAllSessions(GetSender(), CreateNetPacketGameListNew(*game), SessionData::Game);
m_sessionManager.SendLobbyMsgToAllSessions(CreateNetPacketGameListNew(*game), SessionData::Established);
m_gameSessionManager.SendLobbyMsgToAllSessions(CreateNetPacketGameListNew(*game), SessionData::Game);
{
boost::mutex::scoped_lock lock(m_statMutex);
@@ -973,8 +967,8 @@ ServerLobbyThread::InternalRemoveGame(boost::shared_ptr<ServerGameThread> game)
game->RemoveAllSessions();
// Notify all players.
boost::shared_ptr<NetPacket> packet = CreateNetPacketGameListUpdate(game->GetId(), GAME_MODE_CLOSED);
m_sessionManager.SendLobbyMsgToAllSessions(GetSender(), packet, SessionData::Established);
m_gameSessionManager.SendLobbyMsgToAllSessions(GetSender(), packet, SessionData::Game);
m_sessionManager.SendLobbyMsgToAllSessions(packet, SessionData::Established);
m_gameSessionManager.SendLobbyMsgToAllSessions(packet, SessionData::Game);
}
void
@@ -1016,7 +1010,7 @@ ServerLobbyThread::InternalResubscribeMsg(SessionWrapper session)
try {
static_cast<NetPacketStatisticsChanged *>(packet.get())->SetData(statData);
GetSender().Send(session.sessionData, packet);
session.sessionData->GetSender().Send(session.sessionData, packet);
} catch (const NetException &)
{
// Ignore errors for now.
@@ -1053,7 +1047,7 @@ ServerLobbyThread::HandleNewConnection(boost::shared_ptr<ConnectData> connData)
//}
// Create a new session.
boost::shared_ptr<SessionData> sessionData(new SessionData(connData->ReleaseSocket(), m_curSessionId++));
boost::shared_ptr<SessionData> sessionData(new SessionData(connData->ReleaseSocket(), m_curSessionId++, *m_senderCallback));
m_sessionManager.AddSession(sessionData);
LOG_VERBOSE("Accepted connection - session #" << sessionData->GetId() << ".");
@@ -1117,7 +1111,7 @@ ServerLobbyThread::InternalCheckSessionTimeouts(SessionWrapper session)
warningData.timeoutReason = NETWORK_TIMEOUT_GENERIC;
warningData.remainingSeconds = SERVER_TIMEOUT_WARNING_REMAINING_SEC;
static_cast<NetPacketTimeoutWarning *>(packet.get())->SetData(warningData);
GetSender().Send(session.sessionData, packet);
session.sessionData->GetSender().Send(session.sessionData, packet);
}
else if (session.sessionData->GetActivityTimerElapsedSec() >= SERVER_SESSION_ACTIVITY_TIMEOUT_SEC)
{
@@ -1163,7 +1157,7 @@ ServerLobbyThread::SendError(boost::shared_ptr<SessionData> s, int errorCode)
NetPacketError::Data errorData;
errorData.errorCode = errorCode;
static_cast<NetPacketError *>(packet.get())->SetData(errorData);
GetSender().Send(s, packet);
s->GetSender().Send(s, packet);
}
void
@@ -1173,7 +1167,7 @@ ServerLobbyThread::SendJoinGameFailed(boost::shared_ptr<SessionData> s, int reas
NetPacketJoinGameFailed::Data failedData;
failedData.failureCode = reason;
static_cast<NetPacketJoinGameFailed *>(packet.get())->SetData(failedData);
GetSender().Send(s, packet);
s->GetSender().Send(s, packet);
}
void
@@ -1183,7 +1177,7 @@ ServerLobbyThread::SendGameList(boost::shared_ptr<SessionData> s)
GameMap::const_iterator game_end = m_gameMap.end();
while (game_i != game_end)
{
GetSender().Send(s, CreateNetPacketGameListNew(*game_i->second));
s->GetSender().Send(s, CreateNetPacketGameListNew(*game_i->second));
++game_i;
}
}
@@ -1218,8 +1212,8 @@ ServerLobbyThread::BroadcastStatisticsUpdate(const ServerStats &stats)
try {
static_cast<NetPacketStatisticsChanged *>(packet.get())->SetData(statData);
m_sessionManager.SendLobbyMsgToAllSessions(GetSender(), packet, SessionData::Established);
m_gameSessionManager.SendLobbyMsgToAllSessions(GetSender(), packet, SessionData::Game);
m_sessionManager.SendLobbyMsgToAllSessions(packet, SessionData::Established);
m_gameSessionManager.SendLobbyMsgToAllSessions(packet, SessionData::Game);
} catch (const NetException &)
{
// Ignore errors for now.
@@ -1290,13 +1284,6 @@ ServerLobbyThread::GetCallback()
return m_gui;
}
SenderInterface &
ServerLobbyThread::GetSender()
{
assert(m_sender.get());
return *m_sender;
}
ReceiverHelper &
ServerLobbyThread::GetReceiver()
{
+12 -1
View File
@@ -18,15 +18,20 @@
***************************************************************************/
#include <net/sessiondata.h>
#include <net/senderthread.h>
SessionData::SessionData(SOCKET sockfd, SessionId id)
SessionData::SessionData(SOCKET sockfd, SessionId id, SenderCallback &cb)
: m_sockfd(sockfd), m_id(id), m_state(SessionData::Init), m_readyFlag(false),
m_wantsLobbyMsg(true), m_activityTimeoutNoticeSent(false)
{
m_sender.reset(new SenderThread(cb));
m_sender->Start();
}
SessionData::~SessionData()
{
m_sender->SignalStop();
m_sender->WaitStop();
if (m_sockfd != INVALID_SOCKET)
CLOSESOCKET(m_sockfd);
}
@@ -122,6 +127,12 @@ SessionData::GetReceiveBuffer()
return m_receiveBuffer;
}
SenderInterface &
SessionData::GetSender()
{
return *m_sender;
}
void
SessionData::ResetActivityTimer()
{
+6 -6
View File
@@ -357,7 +357,7 @@ SessionManager::GetRawSessionCount()
}
void
SessionManager::SendToAllSessions(SenderInterface &sender, boost::shared_ptr<NetPacket> packet, SessionData::State state)
SessionManager::SendToAllSessions(boost::shared_ptr<NetPacket> packet, SessionData::State state)
{
boost::recursive_mutex::scoped_lock lock(m_sessionMapMutex);
@@ -371,13 +371,13 @@ SessionManager::SendToAllSessions(SenderInterface &sender, boost::shared_ptr<Net
// Send each client (with a certain state) a copy of the packet.
if (i->second.sessionData->GetState() == state)
sender.Send(i->second.sessionData, boost::shared_ptr<NetPacket>(packet->Clone()));
i->second.sessionData->GetSender().Send(i->second.sessionData, boost::shared_ptr<NetPacket>(packet->Clone()));
++i;
}
}
void
SessionManager::SendLobbyMsgToAllSessions(SenderInterface &sender, boost::shared_ptr<NetPacket> packet, SessionData::State state)
SessionManager::SendLobbyMsgToAllSessions(boost::shared_ptr<NetPacket> packet, SessionData::State state)
{
boost::recursive_mutex::scoped_lock lock(m_sessionMapMutex);
@@ -391,13 +391,13 @@ SessionManager::SendLobbyMsgToAllSessions(SenderInterface &sender, boost::shared
// Send each client (with a certain state) a copy of the packet.
if (i->second.sessionData->GetState() == state && i->second.sessionData->WantsLobbyMsg())
sender.Send(i->second.sessionData, boost::shared_ptr<NetPacket>(packet->Clone()));
i->second.sessionData->GetSender().Send(i->second.sessionData, boost::shared_ptr<NetPacket>(packet->Clone()));
++i;
}
}
void
SessionManager::SendToAllButOneSessions(SenderInterface &sender, boost::shared_ptr<NetPacket> packet, SessionId except, SessionData::State state)
SessionManager::SendToAllButOneSessions(boost::shared_ptr<NetPacket> packet, SessionId except, SessionData::State state)
{
boost::recursive_mutex::scoped_lock lock(m_sessionMapMutex);
@@ -409,7 +409,7 @@ SessionManager::SendToAllButOneSessions(SenderInterface &sender, boost::shared_p
// Send each fully connected client but one a copy of the packet.
if (i->second.sessionData->GetState() == state)
if (i->first != except)
sender.Send(i->second.sessionData, boost::shared_ptr<NetPacket>(packet->Clone()));
i->second.sessionData->GetSender().Send(i->second.sessionData, boost::shared_ptr<NetPacket>(packet->Clone()));
++i;
}
}
+5 -27
View File
@@ -29,7 +29,6 @@
#include <list>
#include <boost/shared_ptr.hpp>
#include <third_party/boost/timers.hpp>
#define SENDER_THREAD_TERMINATE_TIMEOUT THREAD_WAIT_INFINITE
@@ -46,43 +45,22 @@ public:
virtual void Send(boost::shared_ptr<SessionData> session, boost::shared_ptr<NetPacket> packet);
virtual void Send(boost::shared_ptr<SessionData> session, const NetPacketList &packetList);
unsigned GetNumPacketsInQueue() const;
bool operator<(const SenderThread &other) const;
protected:
struct SendData
{
SendData()
: bytesSent(0) {}
SendData(boost::shared_ptr<NetPacket> p, boost::shared_ptr<SessionData> s)
: packet(p), session(s), bytesSent(0) {}
boost::shared_ptr<NetPacket> packet;
boost::shared_ptr<SessionData> session;
unsigned bytesSent;
};
typedef std::list<SendData> SendDataList;
typedef std::list<SessionId> SessionIdList;
typedef std::list<boost::shared_ptr<NetPacket> > SendDataList;
// Main function of the thread.
virtual void Main();
void InternalStore(SendDataList &sendQueue, unsigned maxQueueSize, boost::shared_ptr<SessionData> session, boost::shared_ptr<NetPacket> packet);
void InternalStore(SendDataList &sendQueue, unsigned maxQueueSize, boost::shared_ptr<SessionData> session, const NetPacketList &packetList);
private:
SendDataList m_sendQueue;
mutable boost::mutex m_sendQueueMutex;
SendDataList m_stalledQueue;
mutable boost::mutex m_stalledQueueMutex;
SessionIdList m_sessionsStalled; // Cache
mutable boost::mutex m_sessionsStalledMutex;
boost::shared_ptr<SessionData> m_session;
mutable boost::mutex m_sessionMutex;
boost::shared_ptr<NetPacket> m_curPacket;
unsigned m_bytesSent;
SenderCallback &m_callback;
boost::timers::portable::microsec_timer m_logTimer;
};
#endif
-1
View File
@@ -121,7 +121,6 @@ protected:
unsigned GetStateTimerFlag() const;
void SetStateTimerFlag(unsigned flag);
SenderInterface &GetSender();
ReceiverHelper &GetReceiver();
const StartData &GetStartData() const;
-3
View File
@@ -87,8 +87,6 @@ public:
ServerStats GetStats() const;
boost::posix_time::ptime GetStartTime() const;
SenderInterface &GetSender();
protected:
typedef std::deque<boost::shared_ptr<ConnectData> > ConnectQueue;
@@ -190,7 +188,6 @@ private:
GameMap m_gameMap;
boost::shared_ptr<ReceiverHelper> m_receiver;
boost::shared_ptr<SenderInterface> m_sender;
boost::shared_ptr<ServerSenderCallback> m_senderCallback;
GuiInterface &m_gui;
AvatarManager &m_avatarManager;
+6 -1
View File
@@ -32,13 +32,16 @@
#define SESSION_ID_GENERIC 0xFFFFFFFF
typedef unsigned SessionId;
class SenderThread;
class SenderInterface;
class SenderCallback;
class SessionData
{
public:
enum State { Init, ReceivingAvatar, Established, Game };
SessionData(SOCKET sockfd, SessionId id);
SessionData(SOCKET sockfd, SessionId id, SenderCallback &cb);
~SessionData();
SessionId GetId() const;
@@ -58,6 +61,7 @@ public:
void SetClientAddr(const std::string &addr);
ReceiveBuffer &GetReceiveBuffer();
SenderInterface &GetSender();
void ResetActivityTimer();
unsigned GetActivityTimerElapsedSec() const;
@@ -76,6 +80,7 @@ private:
boost::timers::portable::microsec_timer m_activityTimer;
bool m_activityTimeoutNoticeSent;
boost::timers::portable::microsec_timer m_autoDisconnectTimer;
boost::shared_ptr<SenderThread> m_sender;
mutable boost::mutex m_dataMutex;
};
+3 -3
View File
@@ -74,9 +74,9 @@ public:
void Clear();
unsigned GetRawSessionCount();
void SendToAllSessions(SenderInterface &sender, boost::shared_ptr<NetPacket> packet, SessionData::State state);
void SendLobbyMsgToAllSessions(SenderInterface &sender, boost::shared_ptr<NetPacket> packet, SessionData::State state);
void SendToAllButOneSessions(SenderInterface &sender, boost::shared_ptr<NetPacket> packet, SessionId except, SessionData::State state);
void SendToAllSessions(boost::shared_ptr<NetPacket> packet, SessionData::State state);
void SendLobbyMsgToAllSessions(boost::shared_ptr<NetPacket> packet, SessionData::State state);
void SendToAllButOneSessions(boost::shared_ptr<NetPacket> packet, SessionId except, SessionData::State state);
protected: