From d72cdb6873275bb5e7f918929c22f5c1c1b1a1d8 Mon Sep 17 00:00:00 2001 From: lotodore Date: Wed, 7 Jan 2009 21:52:16 +0000 Subject: [PATCH] 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. --- src/net/clientcontext.h | 4 + src/net/clientthread.h | 3 - src/net/common/clientcontext.cpp | 18 ++- src/net/common/clientstate.cpp | 6 +- src/net/common/clientthread.cpp | 48 +----- src/net/common/senderthread.cpp | 221 ++++++--------------------- src/net/common/servergamestate.cpp | 10 +- src/net/common/servergamethread.cpp | 16 +- src/net/common/serverlobbythread.cpp | 75 ++++----- src/net/common/sessiondata.cpp | 13 +- src/net/common/sessionmanager.cpp | 12 +- src/net/senderthread.h | 32 +--- src/net/servergamethread.h | 1 - src/net/serverlobbythread.h | 3 - src/net/sessiondata.h | 7 +- src/net/sessionmanager.h | 6 +- 16 files changed, 150 insertions(+), 325 deletions(-) diff --git a/src/net/clientcontext.h b/src/net/clientcontext.h index d630c24e..80757294 100644 --- a/src/net/clientcontext.h +++ b/src/net/clientcontext.h @@ -24,6 +24,9 @@ #include #include #include +#include + +class SenderCallback; class ClientContext : public NetContext { @@ -105,6 +108,7 @@ private: std::string m_cacheDir; bool m_hasSubscribedLobbyMsg; ReceiveBuffer m_receiveBuffer; + boost::shared_ptr m_senderCallback; }; #endif diff --git a/src/net/clientthread.h b/src/net/clientthread.h index 948775b1..f83cc8ca 100644 --- a/src/net/clientthread.h +++ b/src/net/clientthread.h @@ -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 m_context; - boost::shared_ptr m_senderCallback; ClientState *m_curState; GuiInterface &m_gui; AvatarManager &m_avatarManager; - boost::shared_ptr m_sender; boost::shared_ptr m_receiver; GameData m_gameData; diff --git a/src/net/common/clientcontext.cpp b/src/net/common/clientcontext.cpp index 82af1237..cf983935 100644 --- a/src/net/common/clientcontext.cpp +++ b/src/net/common/clientcontext.cpp @@ -18,12 +18,28 @@ ***************************************************************************/ #include +#include + +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 diff --git a/src/net/common/clientstate.cpp b/src/net/common/clientstate.cpp index e5bc1a4d..94eb7cf0 100644 --- a/src/net/common/clientstate.cpp +++ b/src/net/common/clientstate.cpp @@ -627,7 +627,7 @@ ClientStateStartSession::Process(ClientThread &client) boost::shared_ptr 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 startAck(new NetPacketStartEventAck); - client.GetSender().Send(client.GetContext().GetSessionData(), startAck); + client.GetContext().GetSessionData()->GetSender().Send(client.GetContext().GetSessionData(), startAck); // Unsubscribe lobby messages. client.UnsubscribeLobbyMsg(); diff --git a/src/net/common/clientthread.cpp b/src/net/common/clientthread.cpp index 126ee293..664607ad 100644 --- a/src/net/common/clientthread.cpp +++ b/src/net/common/clientthread.cpp @@ -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(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(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 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 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() { diff --git a/src/net/common/senderthread.cpp b/src/net/common/senderthread.cpp index 1433025d..737f8e92 100644 --- a/src/net/common/senderthread.cpp +++ b/src/net/common/senderthread.cpp @@ -25,7 +25,6 @@ #include #include -#include #include 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 session, boost::shared_ptrGetId()) == 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 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 session, boost::shared_ptr 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 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(); - } } } diff --git a/src/net/common/servergamestate.cpp b/src/net/common/servergamestate.cpp index c6b1ccd5..eec896a7 100644 --- a/src/net/common/servergamestate.cpp +++ b/src/net/common/servergamestate.cpp @@ -323,7 +323,7 @@ ServerGameStateInit::HandleNewSession(ServerGameThread &server, SessionWrapper s joinGameAckData.prights = session.playerData->GetRights(); joinGameAckData.gameData = server.GetGameData(); static_cast(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(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(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(reject.get())->SetData(rejectData); - server.GetSender().Send(session.sessionData, reject); + session.sessionData->GetSender().Send(session.sessionData, reject); } } diff --git a/src/net/common/servergamethread.cpp b/src/net/common/servergamethread.cpp index 8f3db82a..fe08ebef 100644 --- a/src/net/common/servergamethread.cpp +++ b/src/net/common/servergamethread.cpp @@ -91,7 +91,7 @@ ServerGameThread::GetCurRound() const void ServerGameThread::SendToAllPlayers(boost::shared_ptr 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(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(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 player, int rea NetPacketGameAdminChanged::Data adminChangedData; adminChangedData.playerId = newAdmin->GetUniqueId(); // Choose next player as admin. static_cast(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 player, int rea thisPlayerLeftData.playerId = player->GetUniqueId(); thisPlayerLeftData.removeReason = reason; static_cast(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() { diff --git a/src/net/common/serverlobbythread.cpp b/src/net/common/serverlobbythread.cpp index 166ee679..519c1bdf 100644 --- a/src/net/common/serverlobbythread.cpp +++ b/src/net/common/serverlobbythread.cpp @@ -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(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(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(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(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 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 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(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(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(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(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(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(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(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 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 game) game->RemoveAllSessions(); // Notify all players. boost::shared_ptr 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(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 connData) //} // Create a new session. - boost::shared_ptr sessionData(new SessionData(connData->ReleaseSocket(), m_curSessionId++)); + boost::shared_ptr 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(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 s, int errorCode) NetPacketError::Data errorData; errorData.errorCode = errorCode; static_cast(packet.get())->SetData(errorData); - GetSender().Send(s, packet); + s->GetSender().Send(s, packet); } void @@ -1173,7 +1167,7 @@ ServerLobbyThread::SendJoinGameFailed(boost::shared_ptr s, int reas NetPacketJoinGameFailed::Data failedData; failedData.failureCode = reason; static_cast(packet.get())->SetData(failedData); - GetSender().Send(s, packet); + s->GetSender().Send(s, packet); } void @@ -1183,7 +1177,7 @@ ServerLobbyThread::SendGameList(boost::shared_ptr 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(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() { diff --git a/src/net/common/sessiondata.cpp b/src/net/common/sessiondata.cpp index 6eb0ce2c..dd633a3a 100644 --- a/src/net/common/sessiondata.cpp +++ b/src/net/common/sessiondata.cpp @@ -18,15 +18,20 @@ ***************************************************************************/ #include +#include -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() { diff --git a/src/net/common/sessionmanager.cpp b/src/net/common/sessionmanager.cpp index 75daa77d..5d1cb3fe 100644 --- a/src/net/common/sessionmanager.cpp +++ b/src/net/common/sessionmanager.cpp @@ -357,7 +357,7 @@ SessionManager::GetRawSessionCount() } void -SessionManager::SendToAllSessions(SenderInterface &sender, boost::shared_ptr packet, SessionData::State state) +SessionManager::SendToAllSessions(boost::shared_ptr packet, SessionData::State state) { boost::recursive_mutex::scoped_lock lock(m_sessionMapMutex); @@ -371,13 +371,13 @@ SessionManager::SendToAllSessions(SenderInterface &sender, boost::shared_ptrsecond.sessionData->GetState() == state) - sender.Send(i->second.sessionData, boost::shared_ptr(packet->Clone())); + i->second.sessionData->GetSender().Send(i->second.sessionData, boost::shared_ptr(packet->Clone())); ++i; } } void -SessionManager::SendLobbyMsgToAllSessions(SenderInterface &sender, boost::shared_ptr packet, SessionData::State state) +SessionManager::SendLobbyMsgToAllSessions(boost::shared_ptr 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(packet->Clone())); + i->second.sessionData->GetSender().Send(i->second.sessionData, boost::shared_ptr(packet->Clone())); ++i; } } void -SessionManager::SendToAllButOneSessions(SenderInterface &sender, boost::shared_ptr packet, SessionId except, SessionData::State state) +SessionManager::SendToAllButOneSessions(boost::shared_ptr 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(packet->Clone())); + i->second.sessionData->GetSender().Send(i->second.sessionData, boost::shared_ptr(packet->Clone())); ++i; } } diff --git a/src/net/senderthread.h b/src/net/senderthread.h index caab1e50..3fd3464d 100644 --- a/src/net/senderthread.h +++ b/src/net/senderthread.h @@ -29,7 +29,6 @@ #include #include -#include #define SENDER_THREAD_TERMINATE_TIMEOUT THREAD_WAIT_INFINITE @@ -46,43 +45,22 @@ public: virtual void Send(boost::shared_ptr session, boost::shared_ptr packet); virtual void Send(boost::shared_ptr session, const NetPacketList &packetList); - unsigned GetNumPacketsInQueue() const; - bool operator<(const SenderThread &other) const; - protected: - struct SendData - { - SendData() - : bytesSent(0) {} - SendData(boost::shared_ptr p, boost::shared_ptr s) - : packet(p), session(s), bytesSent(0) {} - boost::shared_ptr packet; - boost::shared_ptr session; - unsigned bytesSent; - }; - typedef std::list SendDataList; - typedef std::list SessionIdList; + typedef std::list > SendDataList; // Main function of the thread. virtual void Main(); - void InternalStore(SendDataList &sendQueue, unsigned maxQueueSize, boost::shared_ptr session, boost::shared_ptr packet); - void InternalStore(SendDataList &sendQueue, unsigned maxQueueSize, boost::shared_ptr 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 m_session; + mutable boost::mutex m_sessionMutex; + boost::shared_ptr m_curPacket; + unsigned m_bytesSent; SenderCallback &m_callback; - - boost::timers::portable::microsec_timer m_logTimer; }; #endif diff --git a/src/net/servergamethread.h b/src/net/servergamethread.h index ca536eab..294b2d29 100644 --- a/src/net/servergamethread.h +++ b/src/net/servergamethread.h @@ -121,7 +121,6 @@ protected: unsigned GetStateTimerFlag() const; void SetStateTimerFlag(unsigned flag); - SenderInterface &GetSender(); ReceiverHelper &GetReceiver(); const StartData &GetStartData() const; diff --git a/src/net/serverlobbythread.h b/src/net/serverlobbythread.h index dd9f922f..22b8a882 100644 --- a/src/net/serverlobbythread.h +++ b/src/net/serverlobbythread.h @@ -87,8 +87,6 @@ public: ServerStats GetStats() const; boost::posix_time::ptime GetStartTime() const; - SenderInterface &GetSender(); - protected: typedef std::deque > ConnectQueue; @@ -190,7 +188,6 @@ private: GameMap m_gameMap; boost::shared_ptr m_receiver; - boost::shared_ptr m_sender; boost::shared_ptr m_senderCallback; GuiInterface &m_gui; AvatarManager &m_avatarManager; diff --git a/src/net/sessiondata.h b/src/net/sessiondata.h index 6922c350..88a2c54a 100644 --- a/src/net/sessiondata.h +++ b/src/net/sessiondata.h @@ -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 m_sender; mutable boost::mutex m_dataMutex; }; diff --git a/src/net/sessionmanager.h b/src/net/sessionmanager.h index 13572b9d..a3243a8c 100644 --- a/src/net/sessionmanager.h +++ b/src/net/sessionmanager.h @@ -74,9 +74,9 @@ public: void Clear(); unsigned GetRawSessionCount(); - void SendToAllSessions(SenderInterface &sender, boost::shared_ptr packet, SessionData::State state); - void SendLobbyMsgToAllSessions(SenderInterface &sender, boost::shared_ptr packet, SessionData::State state); - void SendToAllButOneSessions(SenderInterface &sender, boost::shared_ptr packet, SessionId except, SessionData::State state); + void SendToAllSessions(boost::shared_ptr packet, SessionData::State state); + void SendLobbyMsgToAllSessions(boost::shared_ptr packet, SessionData::State state); + void SendToAllButOneSessions(boost::shared_ptr packet, SessionId except, SessionData::State state); protected: