Rewriting network receive functions. This is still work in progress to make it more efficient.

This commit is contained in:
lotodore
2011-02-19 23:47:24 +00:00
parent e4679d58e7
commit 5730cd321d
20 changed files with 232 additions and 313 deletions
-1
View File
@@ -21,7 +21,6 @@
#include <net/clientthread.h>
#include <net/clientcontext.h>
#include <net/senderhelper.h>
#include <net/receiverhelper.h>
#include <net/netpacket.h>
#include <net/resolverthread.h>
#include <net/clientexception.h>
+5 -53
View File
@@ -23,7 +23,6 @@
#include <net/clientstate.h>
#include <net/clientcontext.h>
#include <net/senderhelper.h>
#include <net/receiverhelper.h>
#include <net/downloaderthread.h>
#include <net/clientexception.h>
#include <net/socket_msg.h>
@@ -48,22 +47,6 @@
using namespace std;
using boost::asio::ip::tcp;
class ClientSenderCallback : public SenderCallback, public SessionDataCallback
{
public:
ClientSenderCallback() {}
virtual ~ClientSenderCallback() {}
virtual void SignalNetError(SessionId /*session*/, int /*errorID*/, int /*osErrorID*/) {
}
virtual void SignalSessionTerminated(unsigned /*session*/) {
}
private:
};
ClientThread::ClientThread(GuiInterface &gui, AvatarManager &avatarManager)
: m_ioService(new boost::asio::io_service), m_curState(NULL), m_gui(gui),
m_avatarManager(avatarManager), m_isServerSelected(false),
@@ -72,10 +55,8 @@ ClientThread::ClientThread(GuiInterface &gui, AvatarManager &avatarManager)
{
m_clientLog.reset(new Log("", 0));
m_context.reset(new ClientContext);
m_receiver.reset(new ReceiverHelper);
myQtToolsInterface.reset(CreateQtToolsWrapper());
m_senderCallback.reset(new ClientSenderCallback());
m_senderHelper.reset(new SenderHelper(*m_senderCallback, m_ioService));
m_senderHelper.reset(new SenderHelper(m_ioService));
}
ClientThread::~ClientThread()
@@ -381,35 +362,13 @@ ClientThread::SendReportAvatar(unsigned reportedPlayerId, const std::string &ava
void
ClientThread::StartAsyncRead()
{
ReceiveBuffer &buf = GetContext().GetSessionData()->GetReceiveBuffer();
GetContext().GetSessionData()->GetAsioSocket()->async_read_some(
boost::asio::buffer(buf.recvBuf + buf.recvBufUsed, RECV_BUF_SIZE - buf.recvBufUsed),
boost::bind(
&ClientThread::HandleRead,
shared_from_this(),
boost::asio::placeholders::error,
boost::asio::placeholders::bytes_transferred));
GetContext().GetSessionData()->GetReceiveBuffer().StartAsyncRead(GetContext().GetSessionData());
}
void
ClientThread::HandleRead(const boost::system::error_code& ec, size_t bytesRead)
ClientThread::HandlePacket(boost::shared_ptr<SessionData> /*session*/, boost::shared_ptr<NetPacket> packet)
{
if (!ec) {
ReceiveBuffer &buf = GetContext().GetSessionData()->GetReceiveBuffer();
buf.recvBufUsed += bytesRead;
GetReceiver().ScanPackets(buf);
while (!buf.receivedPackets.empty()) {
boost::shared_ptr<NetPacket> packet = buf.receivedPackets.front();
buf.receivedPackets.pop_front();
if (packet)
GetState().HandlePacket(shared_from_this(), packet);
}
StartAsyncRead();
} else {
if (ec != boost::asio::error::operation_aborted)
throw NetException(__FILE__, __LINE__, ERR_SOCK_CONN_RESET, 0);
}
GetState().HandlePacket(shared_from_this(), packet);
}
void
@@ -926,7 +885,7 @@ ClientThread::CreateContextSession()
GetContext().SetSessionData(boost::shared_ptr<SessionData>(new SessionData(
newSock,
SESSION_ID_GENERIC,
*m_senderCallback)));
*this)));
GetContext().SetResolver(boost::shared_ptr<boost::asio::ip::tcp::resolver>(
new boost::asio::ip::tcp::resolver(*m_ioService)));
validSocket = true;
@@ -965,13 +924,6 @@ ClientThread::GetSender()
return *m_senderHelper;
}
ReceiverHelper &
ClientThread::GetReceiver()
{
assert(m_receiver);
return *m_receiver;
}
unsigned
ClientThread::GetGameId() const
{
+126
View File
@@ -0,0 +1,126 @@
/***************************************************************************
* Copyright (C) 2011 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 *
* the Free Software Foundation; either version 2 of the License, or *
* (at your option) any later version. *
* *
* This program is distributed in the hope that it will be useful, *
* but WITHOUT ANY WARRANTY; without even the implied warranty of *
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the *
* GNU General Public License for more details. *
* *
* You should have received a copy of the GNU General Public License *
* along with this program; if not, write to the *
* Free Software Foundation, Inc., *
* 59 Temple Place - Suite 330, Boston, MA 02111-1307, USA. *
***************************************************************************/
#include <boost/asio.hpp>
#include <boost/bind.hpp>
#include <net/receivebuffer.h>
#include <net/sessiondata.h>
#include <core/loghelper.h>
#include <core/pokerthexception.h>
#include <boost/swap.hpp>
using namespace std;
using boost::asio::ip::tcp;
ReceiveBuffer::ReceiveBuffer()
: recvBufUsed(0)
{
recvBuf[0] = 0;
}
void
ReceiveBuffer::StartAsyncRead(boost::shared_ptr<SessionData> session)
{
session->GetAsioSocket()->async_read_some(
boost::asio::buffer(recvBuf + recvBufUsed, RECV_BUF_SIZE - recvBufUsed),
boost::bind(
&ReceiveBuffer::HandleRead,
shared_from_this(),
session,
boost::asio::placeholders::error,
boost::asio::placeholders::bytes_transferred));
}
void
ReceiveBuffer::HandleRead(boost::shared_ptr<SessionData> session, const boost::system::error_code &error, size_t bytesRead)
{
if (error != boost::asio::error::operation_aborted) {
try {
if (!error) {
recvBufUsed += bytesRead;
ScanPackets();
ProcessPackets(session);
StartAsyncRead(session);
} else if (error == boost::asio::error::interrupted || error == boost::asio::error::try_again) {
LOG_ERROR("Session " << session->GetId() << " - recv interrupted: " << error);
StartAsyncRead(session);
} else {
LOG_ERROR("Session " << session->GetId() << " - Connection closed: " << error);
// On error: Close this session.
session->Close();
}
} catch (const exception &e) {
LOG_ERROR("Session " << session->GetId() << " - unhandled exception in HandleRead: " << e.what());
}
}
}
void
ReceiveBuffer::ScanPackets()
{
bool dataAvailable = true;
do {
boost::shared_ptr<NetPacket> tmpPacket;
// This is necessary, because we use TCP.
// Packets may be received in multiple chunks or
// several packets may be received at once.
if (recvBufUsed) {
try {
// This call will also handle the memmove stuff, i.e.
// buffering for partial packets.
tmpPacket = NetPacket::Create(recvBuf, recvBufUsed);
} catch (const exception &e) {
// Reset buffer on error.
recvBufUsed = 0;
LOG_ERROR(e.what());
}
}
if (tmpPacket) {
//cerr << "IN:" << endl << tmpPacket->ToString() << endl;
if (asn_check_constraints(&asn_DEF_PokerTHMessage, tmpPacket->GetMsg(), NULL, NULL) == 0)
receivedPackets.push_back(tmpPacket);
else
LOG_ERROR("Invalid packet: " << endl << tmpPacket->ToString());
} else
dataAvailable = false;
} while(dataAvailable);
}
void
ReceiveBuffer::ProcessPackets(boost::shared_ptr<SessionData> session)
{
while (!receivedPackets.empty()) {
boost::shared_ptr<NetPacket> p = receivedPackets.front();
receivedPackets.pop_front();
// We need to catch specific exceptions, so that they do not affect the server.
try {
session->HandlePacket(p);
} catch (const PokerTHException &e) {
LOG_ERROR("Session " << session->GetId() << " - Read handler exception: " << e.what());
// TODO add error handling, close session.
}
}
if (recvBufUsed >= RECV_BUF_SIZE) {
LOG_ERROR("Session " << session->GetId() << " - Receive buf full: " << recvBufUsed);
recvBufUsed = 0;
}
}
-70
View File
@@ -1,70 +0,0 @@
/***************************************************************************
* Copyright (C) 2007 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 *
* the Free Software Foundation; either version 2 of the License, or *
* (at your option) any later version. *
* *
* This program is distributed in the hope that it will be useful, *
* but WITHOUT ANY WARRANTY; without even the implied warranty of *
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the *
* GNU General Public License for more details. *
* *
* You should have received a copy of the GNU General Public License *
* along with this program; if not, write to the *
* Free Software Foundation, Inc., *
* 59 Temple Place - Suite 330, Boston, MA 02111-1307, USA. *
***************************************************************************/
#ifdef _WIN32
#define FD_SETSIZE 512
#endif
#include <net/receiverhelper.h>
#include <net/socket_msg.h>
#include <net/netexception.h>
#include <core/loghelper.h>
using namespace std;
ReceiverHelper::ReceiverHelper()
{
}
ReceiverHelper::~ReceiverHelper()
{
}
void
ReceiverHelper::ScanPackets(ReceiveBuffer &buf)
{
bool dataAvailable = true;
do {
boost::shared_ptr<NetPacket> tmpPacket;
// This is necessary, because we use TCP.
// Packets may be received in multiple chunks or
// several packets may be received at once.
if (buf.recvBufUsed) {
try {
// This call will also handle the memmove stuff, i.e.
// buffering for partial packets.
tmpPacket = NetPacket::Create(buf.recvBuf, buf.recvBufUsed);
} catch (const exception &e) {
// Reset buffer on error.
buf.recvBufUsed = 0;
LOG_ERROR(e.what());
}
}
if (tmpPacket) {
//cerr << "IN:" << endl << tmpPacket->ToString() << endl;
if (asn_check_constraints(&asn_DEF_PokerTHMessage, tmpPacket->GetMsg(), NULL, NULL) == 0)
buf.receivedPackets.push_back(tmpPacket);
else
LOG_ERROR("Invalid packet: " << endl << tmpPacket->ToString());
} else
dataAvailable = false;
} while(dataAvailable);
}
+3 -3
View File
@@ -18,8 +18,8 @@
***************************************************************************/
#include <net/senderhelper.h>
#include <net/sessiondata.h>
#include <net/sendbuffer.h>
#include <net/sendercallback.h>
#include <net/socket_helper.h>
#include <net/socket_msg.h>
#include <core/loghelper.h>
@@ -28,8 +28,8 @@
using namespace std;
SenderHelper::SenderHelper(SenderCallback &cb, boost::shared_ptr<boost::asio::io_service> ioService)
: m_callback(cb), m_ioService(ioService)
SenderHelper::SenderHelper(boost::shared_ptr<boost::asio::io_service> ioService)
: m_ioService(ioService)
{
}
-10
View File
@@ -26,7 +26,6 @@
#include <net/serverlobbythread.h>
#include <net/serverexception.h>
#include <net/senderhelper.h>
#include <net/receiverhelper.h>
#include <net/socket_msg.h>
#include <core/loghelper.h>
#include <db/serverdbinterface.h>
@@ -57,8 +56,6 @@ ServerGame::ServerGame(boost::shared_ptr<ServerLobbyThread> lobbyThread, u_int32
m_stateTimer1(lobbyThread->GetIOService()), m_stateTimer2(lobbyThread->GetIOService())
{
LOG_VERBOSE("Game object " << GetId() << " created.");
m_receiver.reset(new ReceiverHelper);
}
ServerGame::~ServerGame()
@@ -918,13 +915,6 @@ ServerGame::GetStateTimer2()
return m_stateTimer2;
}
ReceiverHelper &
ServerGame::GetReceiver()
{
assert(m_receiver.get());
return *m_receiver;
}
unsigned
ServerGame::GetSmallDelaySec() const
{
-1
View File
@@ -20,7 +20,6 @@
#include <net/servergamestate.h>
#include <net/servergame.h>
#include <net/serverlobbythread.h>
#include <net/receiverhelper.h>
#include <net/senderhelper.h>
#include <net/netpacket.h>
#include <net/socket_msg.h>
+43 -97
View File
@@ -21,10 +21,9 @@
#include <net/servergame.h>
#include <net/serverbanmanager.h>
#include <net/serverexception.h>
#include <net/receivebuffer.h>
#include <net/senderhelper.h>
#include <net/sendercallback.h>
#include <net/serverircbotcallback.h>
#include <net/receiverhelper.h>
#include <net/socket_msg.h>
#include <net/chatcleanermanager.h>
#include <net/net_helper.h>
@@ -83,18 +82,18 @@ using namespace std;
using boost::asio::ip::tcp;
class InternalServerCallback : public SenderCallback, public SessionDataCallback, public ChatCleanerCallback, public ServerDBCallback
class InternalServerCallback : public SessionDataCallback, public ChatCleanerCallback, public ServerDBCallback
{
public:
InternalServerCallback(ServerLobbyThread &server) : m_server(server) {}
virtual ~InternalServerCallback() {}
virtual void SignalNetError(SessionId /*session*/, int /*errorID*/, int /*osErrorID*/) {
// We just ignore send errors for now, on server side.
// A serious send error should trigger a read error or a read
// returning 0 afterwards, and we will handle this error.
virtual void CloseSession(boost::shared_ptr<SessionData> session) {
m_server.CloseSession(session->GetId());
}
virtual void SignalSessionTerminated(unsigned /*session*/) {
virtual void HandlePacket(boost::shared_ptr<SessionData> session, boost::shared_ptr<NetPacket> packet) {
m_server.DispatchPacket(session, packet);
}
virtual void SignalChatBotMessage(const string &msg) {
@@ -181,8 +180,7 @@ ServerLobbyThread::ServerLobbyThread(GuiInterface &gui, ServerMode mode, ServerI
m_startTime(boost::posix_time::second_clock::local_time())
{
m_internalServerCallback.reset(new InternalServerCallback(*this));
m_sender.reset(new SenderHelper(*m_internalServerCallback, m_ioService));
m_receiver.reset(new ReceiverHelper);
m_sender.reset(new SenderHelper(m_ioService));
m_banManager.reset(new ServerBanManager(m_ioService));
m_chatCleanerManager.reset(new ChatCleanerManager(*m_internalServerCallback, m_ioService));
DBFactory dbFactory;
@@ -276,15 +274,7 @@ ServerLobbyThread::AddConnection(boost::shared_ptr<tcp::socket> sock)
netAnnounce->numPlayersOnServer = m_statData.numberOfPlayersOnServer;
}
GetSender().Send(sessionData, packet);
sock->async_read_some(
boost::asio::buffer(sessionData->GetReceiveBuffer().recvBuf, RECV_BUF_SIZE),
boost::bind(
&ServerLobbyThread::HandleRead,
this,
boost::asio::placeholders::error,
sessionData->GetId(),
boost::asio::placeholders::bytes_transferred));
sessionData->GetReceiveBuffer().StartAsyncRead(sessionData);
}
}
if (!hasClientIp) {
@@ -359,7 +349,22 @@ ServerLobbyThread::RemoveSessionFromGame(SessionWrapper session)
{
// Just remove the session. Only for fatal errors.
CloseSession(session);
session.sessionData->SetGameId(0);
}
void
ServerLobbyThread::CloseSession(SessionId sessionId)
{
SessionWrapper session = m_sessionManager.GetSessionById(sessionId);
if (!session.sessionData)
session = m_gameSessionManager.GetSessionById(sessionId);
if (session.sessionData) {
GameMap::iterator pos = m_gameMap.find(session.sessionData->GetGameId());
if (pos != m_gameMap.end()) {
pos->second->ErrorRemoveSession(session);
} else {
CloseSession(session);
}
}
}
void
@@ -376,6 +381,7 @@ ServerLobbyThread::CloseSession(SessionWrapper session)
NotifyPlayerLeftLobby(session.playerData->GetUniqueId());
// Update stats (if needed).
UpdateStatisticsNumberOfPlayers();
session.sessionData->SetGameId(0);
}
}
@@ -838,78 +844,25 @@ ServerLobbyThread::InitChatCleaner()
}
void
ServerLobbyThread::HandleRead(const boost::system::error_code &ec, SessionId sessionId, size_t bytesRead)
ServerLobbyThread::DispatchPacket(boost::shared_ptr<SessionData> s, boost::shared_ptr<NetPacket> packet)
{
if (ec != boost::asio::error::operation_aborted) {
try {
// Find the session.
SessionWrapper session = m_sessionManager.GetSessionById(sessionId);
if (!session.sessionData)
session = m_gameSessionManager.GetSessionById(sessionId);
if (session.sessionData) {
ReceiveBuffer &buf = session.sessionData->GetReceiveBuffer();
if (!ec) {
if (buf.recvBufUsed + bytesRead > RECV_BUF_SIZE)
LOG_ERROR("Session " << session.sessionData->GetId() << " - Internal error: Receive buffer overflow!");
buf.recvBufUsed += bytesRead;
GetReceiver().ScanPackets(buf);
bool errorFlag = false;
while (!buf.receivedPackets.empty()) {
boost::shared_ptr<NetPacket> packet = buf.receivedPackets.front();
buf.receivedPackets.pop_front();
// Retrieve current game, if applicable.
boost::shared_ptr<ServerGame> game = InternalGetGameFromId(session.sessionData->GetGameId());
if (game) {
// We need to catch game-specific exceptions, so that they do not affect the server.
try {
game->HandlePacket(session, packet);
} catch (const PokerTHException &e) {
LOG_ERROR("Game " << game->GetId() << " - Read handler exception: " << e.what());
game->RemoveAllSessions();
errorFlag = true;
break;
}
} else
HandlePacket(session, packet);
}
if (buf.recvBufUsed >= RECV_BUF_SIZE) {
LOG_ERROR("Session " << session.sessionData->GetId() << " - Full receive buf but no valid packet.");
buf.recvBufUsed = 0;
}
if (!errorFlag) {
session.sessionData->GetAsioSocket()->async_read_some(
boost::asio::buffer(buf.recvBuf + buf.recvBufUsed, RECV_BUF_SIZE - buf.recvBufUsed),
boost::bind(
&ServerLobbyThread::HandleRead,
this,
boost::asio::placeholders::error,
sessionId,
boost::asio::placeholders::bytes_transferred));
}
} else if (ec == boost::asio::error::interrupted || ec == boost::asio::error::try_again) {
LOG_ERROR("Session " << sessionId << " - recv interrupted: " << ec);
session.sessionData->GetAsioSocket()->async_read_some(
boost::asio::buffer(buf.recvBuf + buf.recvBufUsed, RECV_BUF_SIZE - buf.recvBufUsed),
boost::bind(
&ServerLobbyThread::HandleRead,
this,
boost::asio::placeholders::error,
sessionId,
boost::asio::placeholders::bytes_transferred));
} else {
LOG_ERROR("Session " << sessionId << " - Connection closed: " << ec);
// On error: Close this session.
boost::shared_ptr<ServerGame> game = InternalGetGameFromId(session.sessionData->GetGameId());
if (game)
game->ErrorRemoveSession(session);
else
CloseSession(session);
}
// Find the session.
SessionWrapper session = m_sessionManager.GetSessionById(s->GetId());
if (!session.sessionData)
session = m_gameSessionManager.GetSessionById(s->GetId());
if (session.sessionData) {
// Retrieve current game, if applicable.
boost::shared_ptr<ServerGame> game = InternalGetGameFromId(session.sessionData->GetGameId());
if (game) {
// We need to catch game-specific exceptions, so that they do not affect the server.
try {
game->HandlePacket(session, packet);
} catch (const PokerTHException &e) {
LOG_ERROR("Game " << game->GetId() << " - Read handler exception: " << e.what());
game->RemoveAllSessions();
}
} catch (const exception &e) {
LOG_ERROR("Session " << sessionId << " - unhandled exception in HandleRead: " << e.what());
}
} else
HandlePacket(session, packet);
}
}
@@ -2106,13 +2059,6 @@ ServerLobbyThread::GetIrcBotCallback()
return m_ircBotCb;
}
ReceiverHelper &
ServerLobbyThread::GetReceiver()
{
assert(m_receiver.get());
return *m_receiver;
}
InternalServerCallback &
ServerLobbyThread::GetSenderCallback()
{
+2 -1
View File
@@ -18,6 +18,7 @@
***************************************************************************/
#include <net/sessiondata.h>
#include <net/receivebuffer.h>
#include <net/sendbuffer.h>
#include <gsasl.h>
@@ -30,12 +31,12 @@ SessionData::SessionData(boost::shared_ptr<boost::asio::ip::tcp::socket> sock, S
m_autoDisconnectTimer(boost::posix_time::time_duration(0, 0, 0), boost::timers::portable::microsec_timer::auto_start),
m_callback(cb), m_authSession(NULL), m_curAuthStep(0)
{
m_receiveBuffer.reset(new ReceiveBuffer);
m_sendBuffer.reset(new SendBuffer);
}
SessionData::~SessionData()
{
m_callback.SignalSessionTerminated(m_id);
InternalClearAuthSession();
}