From fee8450a744fb7f8f6a10b6c28a6f17eb574e6a7 Mon Sep 17 00:00:00 2001 From: lotodore Date: Tue, 30 Oct 2007 21:02:52 +0000 Subject: [PATCH] Still trying to fix the SIGPIPE problem. --- src/net/common/senderthread.cpp | 109 +++++++++++++++++---------- src/net/common/serverlobbythread.cpp | 32 ++++---- src/net/common/sessionmanager.cpp | 19 +++++ src/net/senderthread.h | 1 + src/net/sessionmanager.h | 1 + 5 files changed, 105 insertions(+), 57 deletions(-) diff --git a/src/net/common/senderthread.cpp b/src/net/common/senderthread.cpp index da015248..181335c5 100644 --- a/src/net/common/senderthread.cpp +++ b/src/net/common/senderthread.cpp @@ -28,13 +28,14 @@ using namespace std; -#define SEND_ERROR_TIMEOUT_MSEC 20000 +#define SEND_ERROR_NORMAL_TIMEOUT_MSEC 20000 +#define SEND_ERROR_LOW_PRIO_TIMEOUT_MSEC 10000 #define SEND_TIMEOUT_MSEC 10 #define SEND_QUEUE_SIZE 1000 #define SEND_LOW_PRIO_QUEUE_SIZE 50000 SenderThread::SenderThread(SenderCallback &cb) -: m_tmpOutBufSize(0), m_callback(cb) +: m_tmpOutBufSize(0), m_tmpIsLowPrio(false), m_callback(cb) { } @@ -121,11 +122,6 @@ SenderThread::Main() while (!ShouldTerminate()) { - if (sendTimer.is_running() && sendTimer.elapsed().total_milliseconds() > SEND_ERROR_TIMEOUT_MSEC) - { - RemoveCurSendData(); - sendTimer.reset(); - } // Send remaining bytes of output buffer OR // copy ONE packet to output buffer. // For reasons of simplicity, only one packet is sent at a time. @@ -139,6 +135,7 @@ SenderThread::Main() { tmpData = m_outBuf.front(); m_outBuf.pop_front(); + m_tmpIsLowPrio = false; } } @@ -150,6 +147,7 @@ SenderThread::Main() { tmpData = m_lowPrioOutBuf.front(); m_lowPrioOutBuf.pop_front(); + m_tmpIsLowPrio = true; } } @@ -171,63 +169,92 @@ SenderThread::Main() if (m_tmpOutBufSize) { SOCKET tmpSocket = m_curSession->GetSocket(); - fd_set writeSet; - struct timeval timeout; - FD_ZERO(&writeSet); - FD_SET(tmpSocket, &writeSet); + // send next chunk of data + int bytesSent = send(tmpSocket, m_tmpOutBuf, m_tmpOutBufSize, SOCKET_SEND_FLAGS); - timeout.tv_sec = 0; - timeout.tv_usec = SEND_TIMEOUT_MSEC * 1000; - int selectResult = select(tmpSocket + 1, NULL, &writeSet, NULL, &timeout); - if (!IS_VALID_SELECT(selectResult)) + if (!IS_VALID_SEND(bytesSent)) { // Never assume that this is a fatal error. int errCode = SOCKET_ERRNO(); - if (errCode != SOCKET_ERR_WOULDBLOCK) + if (errCode == SOCKET_ERR_WOULDBLOCK) + { + fd_set writeSet; + struct timeval timeout; + + FD_ZERO(&writeSet); + FD_SET(tmpSocket, &writeSet); + + timeout.tv_sec = 0; + timeout.tv_usec = SEND_TIMEOUT_MSEC * 1000; + int selectResult = select(tmpSocket + 1, NULL, &writeSet, NULL, &timeout); + if (!IS_VALID_SELECT(selectResult)) + { + // Never assume that this is a fatal error. + int errCode = SOCKET_ERRNO(); + if (errCode != SOCKET_ERR_WOULDBLOCK) + { + // 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(m_curSession->GetId(), ERR_SOCK_SELECT_FAILED, errCode); + RemoveCurSendData(); + } + Msleep(SEND_TIMEOUT_MSEC); + } + } + else { // 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(m_curSession->GetId(), ERR_SOCK_SELECT_FAILED, errCode); + m_callback.SignalNetError(m_curSession->GetId(), ERR_SOCK_SEND_FAILED, errCode); RemoveCurSendData(); } Msleep(SEND_TIMEOUT_MSEC); } - if (selectResult > 0) // send is possible + else if ((unsigned)bytesSent < m_tmpOutBufSize) { - // send next chunk of data - int bytesSent = send(tmpSocket, m_tmpOutBuf, m_tmpOutBufSize, SOCKET_SEND_FLAGS); - - if (!IS_VALID_SEND(bytesSent)) - { - // Never assume that this is a fatal error. - int errCode = SOCKET_ERRNO(); - if (errCode != SOCKET_ERR_WOULDBLOCK) - { - // 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(m_curSession->GetId(), ERR_SOCK_SEND_FAILED, errCode); - RemoveCurSendData(); - } - Msleep(SEND_TIMEOUT_MSEC); - } - else if ((unsigned)bytesSent < m_tmpOutBufSize) + if (bytesSent) { m_tmpOutBufSize -= (unsigned)bytesSent; memmove(m_tmpOutBuf, m_tmpOutBuf + bytesSent, m_tmpOutBufSize); } else - { - m_tmpOutBufSize = 0; - m_curSession.reset(); - sendTimer.reset(); - } + Msleep(SEND_TIMEOUT_MSEC); + } + else + { + m_tmpOutBufSize = 0; + m_curSession.reset(); + sendTimer.reset(); } } else Msleep(SEND_TIMEOUT_MSEC); + + // Check whether the send timed out. + if (sendTimer.is_running()) + { + bool doRemove = false; + + unsigned msec = static_cast(sendTimer.elapsed().total_milliseconds()); + if (m_tmpIsLowPrio) + { + if (msec > SEND_ERROR_LOW_PRIO_TIMEOUT_MSEC) + doRemove = true; + } + else + { + if (msec > SEND_ERROR_NORMAL_TIMEOUT_MSEC) + doRemove = true; + } + if (doRemove) + { + RemoveCurSendData(); + sendTimer.reset(); + } + } } } diff --git a/src/net/common/serverlobbythread.cpp b/src/net/common/serverlobbythread.cpp index ff96cec8..8afd7cd1 100644 --- a/src/net/common/serverlobbythread.cpp +++ b/src/net/common/serverlobbythread.cpp @@ -138,8 +138,8 @@ ServerLobbyThread::NotifyPlayerJoinedGame(unsigned gameId, unsigned playerId) packetData.gameId = gameId; packetData.playerId = playerId; static_cast(packet.get())->SetData(packetData); - m_sessionManager.SendToAllSessions(GetSender(), packet, SessionData::Established); - m_gameSessionManager.SendToAllSessions(GetSender(), packet, SessionData::Game); + m_sessionManager.SendToAllSessionsLowPrio(GetSender(), packet, SessionData::Established); + m_gameSessionManager.SendToAllSessionsLowPrio(GetSender(), packet, SessionData::Game); } void @@ -151,8 +151,8 @@ ServerLobbyThread::NotifyPlayerLeftGame(unsigned gameId, unsigned playerId) packetData.gameId = gameId; packetData.playerId = playerId; static_cast(packet.get())->SetData(packetData); - m_sessionManager.SendToAllSessions(GetSender(), packet, SessionData::Established); - m_gameSessionManager.SendToAllSessions(GetSender(), packet, SessionData::Game); + m_sessionManager.SendToAllSessionsLowPrio(GetSender(), packet, SessionData::Established); + m_gameSessionManager.SendToAllSessionsLowPrio(GetSender(), packet, SessionData::Game); } void @@ -164,24 +164,24 @@ ServerLobbyThread::NotifyGameAdminChanged(unsigned gameId, unsigned newAdminPlay packetData.gameId = gameId; packetData.newAdminplayerId = newAdminPlayerId; static_cast(packet.get())->SetData(packetData); - m_sessionManager.SendToAllSessions(GetSender(), packet, SessionData::Established); - m_gameSessionManager.SendToAllSessions(GetSender(), packet, SessionData::Game); + m_sessionManager.SendToAllSessionsLowPrio(GetSender(), packet, SessionData::Established); + m_gameSessionManager.SendToAllSessionsLowPrio(GetSender(), packet, SessionData::Game); } void ServerLobbyThread::NotifyStartingGame(unsigned gameId) { boost::shared_ptr packet = CreateNetPacketGameListUpdate(gameId, GAME_MODE_STARTED); - m_sessionManager.SendToAllSessions(GetSender(), packet, SessionData::Established); - m_gameSessionManager.SendToAllSessions(GetSender(), packet, SessionData::Game); + m_sessionManager.SendToAllSessionsLowPrio(GetSender(), packet, SessionData::Established); + m_gameSessionManager.SendToAllSessionsLowPrio(GetSender(), packet, SessionData::Game); } void ServerLobbyThread::NotifyReopeningGame(unsigned gameId) { boost::shared_ptr packet = CreateNetPacketGameListUpdate(gameId, GAME_MODE_CREATED); - m_sessionManager.SendToAllSessions(GetSender(), packet, SessionData::Established); - m_gameSessionManager.SendToAllSessions(GetSender(), packet, SessionData::Game); + m_sessionManager.SendToAllSessionsLowPrio(GetSender(), packet, SessionData::Established); + m_gameSessionManager.SendToAllSessionsLowPrio(GetSender(), packet, SessionData::Game); } void @@ -737,8 +737,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.SendToAllSessions(GetSender(), CreateNetPacketGameListNew(*game), SessionData::Established); - m_gameSessionManager.SendToAllSessions(GetSender(), CreateNetPacketGameListNew(*game), SessionData::Game); + m_sessionManager.SendToAllSessionsLowPrio(GetSender(), CreateNetPacketGameListNew(*game), SessionData::Established); + m_gameSessionManager.SendToAllSessionsLowPrio(GetSender(), CreateNetPacketGameListNew(*game), SessionData::Game); ++m_totalGamesStarted; BroadcastStatisticsUpdate(); @@ -753,8 +753,8 @@ ServerLobbyThread::InternalRemoveGame(boost::shared_ptr game) game->RemoveAllSessions(); // Notify all players. boost::shared_ptr packet = CreateNetPacketGameListUpdate(game->GetId(), GAME_MODE_CLOSED); - m_sessionManager.SendToAllSessions(GetSender(), packet, SessionData::Established); - m_gameSessionManager.SendToAllSessions(GetSender(), packet, SessionData::Game); + m_sessionManager.SendToAllSessionsLowPrio(GetSender(), packet, SessionData::Established); + m_gameSessionManager.SendToAllSessionsLowPrio(GetSender(), packet, SessionData::Game); } void @@ -891,8 +891,8 @@ ServerLobbyThread::BroadcastStatisticsUpdate() try { static_cast(packet.get())->SetData(statData); - m_sessionManager.SendToAllSessions(GetSender(), packet, SessionData::Established); - m_gameSessionManager.SendToAllSessions(GetSender(), packet, SessionData::Game); + m_sessionManager.SendToAllSessionsLowPrio(GetSender(), packet, SessionData::Established); + m_gameSessionManager.SendToAllSessionsLowPrio(GetSender(), packet, SessionData::Game); } catch (const NetException &) { // Ignore errors for now. diff --git a/src/net/common/sessionmanager.cpp b/src/net/common/sessionmanager.cpp index 8b9f6eb0..0aec9e3f 100644 --- a/src/net/common/sessionmanager.cpp +++ b/src/net/common/sessionmanager.cpp @@ -360,6 +360,25 @@ SessionManager::SendToAllSessions(SenderThread &sender, boost::shared_ptr packet, SessionData::State state) +{ + boost::mutex::scoped_lock lock(m_sessionMapMutex); + + SessionMap::iterator i = m_sessionMap.begin(); + SessionMap::iterator end = m_sessionMap.end(); + + while (i != end) + { + assert(i->second.sessionData.get()); + + // Send each client (with a certain state) a copy of the packet. + if (i->second.sessionData->GetState() == state) + sender.SendLowPrio(i->second.sessionData, boost::shared_ptr(packet->Clone())); + ++i; + } +} + void SessionManager::SendToAllButOneSessions(SenderThread &sender, boost::shared_ptr packet, SessionId except, SessionData::State state) { diff --git a/src/net/senderthread.h b/src/net/senderthread.h index efe9b8a4..f050c598 100644 --- a/src/net/senderthread.h +++ b/src/net/senderthread.h @@ -67,6 +67,7 @@ private: char m_tmpOutBuf[MAX_PACKET_SIZE]; unsigned m_tmpOutBufSize; + bool m_tmpIsLowPrio; SenderCallback &m_callback; }; diff --git a/src/net/sessionmanager.h b/src/net/sessionmanager.h index d3d7a126..dab5b94a 100644 --- a/src/net/sessionmanager.h +++ b/src/net/sessionmanager.h @@ -74,6 +74,7 @@ public: unsigned GetRawSessionCount(); void SendToAllSessions(SenderThread &sender, boost::shared_ptr packet, SessionData::State state); + void SendToAllSessionsLowPrio(SenderThread &sender, boost::shared_ptr packet, SessionData::State state); void SendToAllButOneSessions(SenderThread &sender, boost::shared_ptr packet, SessionId except, SessionData::State state); protected: