Fixed a really nasty receive bug (TCP only). Using SCTP would be so much better...
This commit is contained in:
@@ -376,7 +376,8 @@ AbstractClientStateReceiving::Process(ClientThread &client)
|
||||
int retVal = MSG_SOCK_INTERNAL_PENDING;
|
||||
|
||||
// delegate to receiver helper class
|
||||
boost::shared_ptr<NetPacket> tmpPacket = client.GetReceiver().Recv(client.GetContext().GetSocket());
|
||||
boost::shared_ptr<NetPacket> tmpPacket =
|
||||
client.GetReceiver().Recv(client.GetContext().GetSocket(), client.GetContext().GetReceiveBuffer());
|
||||
|
||||
if (tmpPacket.get())
|
||||
{
|
||||
|
||||
@@ -26,7 +26,6 @@ using namespace std;
|
||||
|
||||
|
||||
ReceiverHelper::ReceiverHelper()
|
||||
: m_socket(INVALID_SOCKET), m_tmpInBufSize(0)
|
||||
{
|
||||
}
|
||||
|
||||
@@ -34,25 +33,14 @@ ReceiverHelper::~ReceiverHelper()
|
||||
{
|
||||
}
|
||||
|
||||
void
|
||||
ReceiverHelper::Init(SOCKET socket)
|
||||
{
|
||||
if (!IS_VALID_SOCKET(socket))
|
||||
return; // TODO: throw exception
|
||||
|
||||
m_socket = socket;
|
||||
}
|
||||
|
||||
boost::shared_ptr<NetPacket>
|
||||
ReceiverHelper::Recv(SOCKET sock)
|
||||
ReceiverHelper::Recv(SOCKET sock, ReceiveBuffer &buf)
|
||||
{
|
||||
boost::shared_ptr<NetPacket> tmpPacket(InternalGetPacket());
|
||||
|
||||
if (!tmpPacket.get())
|
||||
if (buf.receivedPackets.empty())
|
||||
{
|
||||
unsigned bufSize = RECV_BUF_SIZE - m_tmpInBufSize;
|
||||
int bufSize = RECV_BUF_SIZE - buf.recvBufUsed;
|
||||
|
||||
if (bufSize) // check if there is room in the input buffer
|
||||
if (bufSize > 0) // check if there is room in the input buffer
|
||||
{
|
||||
fd_set readSet;
|
||||
struct timeval timeout;
|
||||
@@ -69,7 +57,7 @@ ReceiverHelper::Recv(SOCKET sock)
|
||||
}
|
||||
if (selectResult > 0) // recv is possible
|
||||
{
|
||||
int bytesRecvd = recv(sock, m_tmpInBuf + m_tmpInBufSize, bufSize, 0);
|
||||
int bytesRecvd = recv(sock, buf.recvBuf + buf.recvBufUsed, bufSize, 0);
|
||||
|
||||
if (!IS_VALID_RECV(bytesRecvd))
|
||||
{
|
||||
@@ -81,37 +69,49 @@ ReceiverHelper::Recv(SOCKET sock)
|
||||
}
|
||||
else
|
||||
{
|
||||
m_tmpInBufSize += bytesRecvd;
|
||||
tmpPacket = InternalGetPacket();
|
||||
buf.recvBufUsed += bytesRecvd;
|
||||
InternalGetPackets(buf);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
return tmpPacket;
|
||||
}
|
||||
|
||||
boost::shared_ptr<NetPacket>
|
||||
ReceiverHelper::InternalGetPacket()
|
||||
{
|
||||
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 (m_tmpInBufSize >= MIN_PACKET_SIZE)
|
||||
if (!buf.receivedPackets.empty())
|
||||
{
|
||||
try
|
||||
{
|
||||
// This call will also handle the memmove stuff, i.e.
|
||||
// buffering for partial packets.
|
||||
tmpPacket = NetPacket::Create(m_tmpInBuf, m_tmpInBufSize);
|
||||
} catch (const NetException &)
|
||||
{
|
||||
// Reset buffer on error.
|
||||
m_tmpInBufSize = 0;
|
||||
// TODO: log error/increase error counter.
|
||||
}
|
||||
tmpPacket = buf.receivedPackets.front();
|
||||
buf.receivedPackets.pop_front();
|
||||
}
|
||||
return tmpPacket;
|
||||
}
|
||||
|
||||
void
|
||||
ReceiverHelper::InternalGetPackets(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 >= MIN_PACKET_SIZE)
|
||||
{
|
||||
try
|
||||
{
|
||||
// This call will also handle the memmove stuff, i.e.
|
||||
// buffering for partial packets.
|
||||
tmpPacket = NetPacket::Create(buf.recvBuf, buf.recvBufUsed);
|
||||
} catch (const NetException &)
|
||||
{
|
||||
// Reset buffer on error.
|
||||
buf.recvBufUsed = 0;
|
||||
// TODO: log error/increase error counter.
|
||||
}
|
||||
}
|
||||
if (tmpPacket.get())
|
||||
buf.receivedPackets.push_back(tmpPacket);
|
||||
else
|
||||
dataAvailable = false;
|
||||
} while(dataAvailable);
|
||||
}
|
||||
|
||||
|
||||
@@ -139,7 +139,7 @@ AbstractServerGameStateReceiving::Process(ServerGameThread &server)
|
||||
try
|
||||
{
|
||||
// Receive the packet.
|
||||
packet = server.GetReceiver().Recv(session.sessionData->GetSocket());
|
||||
packet = server.GetReceiver().Recv(session.sessionData->GetSocket(), session.sessionData->GetReceiveBuffer());
|
||||
} catch (const NetException &)
|
||||
{
|
||||
server.CloseSessionDelayed(session);
|
||||
|
||||
@@ -193,7 +193,7 @@ ServerLobbyThread::ProcessLoop()
|
||||
try
|
||||
{
|
||||
// Receive the next packet.
|
||||
packet = GetReceiver().Recv(session.sessionData->GetSocket());
|
||||
packet = GetReceiver().Recv(session.sessionData->GetSocket(), session.sessionData->GetReceiveBuffer());
|
||||
} catch (const NetException &)
|
||||
{
|
||||
// On error: Close this session.
|
||||
|
||||
@@ -95,44 +95,55 @@ SessionManager::Select(unsigned timeoutMsec)
|
||||
|
||||
while (i != end)
|
||||
{
|
||||
// Collect all sockets.
|
||||
SOCKET tmpSock = i->first;
|
||||
FD_SET(tmpSock, &rdset);
|
||||
if (tmpSock > maxSock || maxSock == INVALID_SOCKET)
|
||||
maxSock = tmpSock;
|
||||
|
||||
// Check if a packet is available.
|
||||
if (!i->second.sessionData->GetReceiveBuffer().receivedPackets.empty())
|
||||
{
|
||||
retSession = i->second;
|
||||
break;
|
||||
}
|
||||
++i;
|
||||
}
|
||||
}
|
||||
|
||||
if (maxSock == INVALID_SOCKET)
|
||||
if (!retSession.sessionData.get())
|
||||
{
|
||||
Thread::Msleep(timeoutMsec); // just sleep if there is no session
|
||||
}
|
||||
else
|
||||
{
|
||||
// wait for data
|
||||
struct timeval timeout;
|
||||
timeout.tv_sec = timeoutMsec / 1000;
|
||||
timeout.tv_usec = (timeoutMsec % 1000) * 1000;
|
||||
int selectResult = select(maxSock + 1, &rdset, NULL, NULL, &timeout);
|
||||
if (!IS_VALID_SELECT(selectResult))
|
||||
if (maxSock == INVALID_SOCKET)
|
||||
{
|
||||
throw ServerException(ERR_SOCK_SELECT_FAILED, SOCKET_ERRNO());
|
||||
Thread::Msleep(timeoutMsec); // just sleep if there is no session
|
||||
}
|
||||
if (selectResult > 0) // one (or more) of the sockets is readable
|
||||
else
|
||||
{
|
||||
// Check which socket is readable, return the first.
|
||||
boost::mutex::scoped_lock lock(m_sessionMapMutex);
|
||||
SessionMap::iterator i = m_sessionMap.begin();
|
||||
SessionMap::iterator end = m_sessionMap.end();
|
||||
|
||||
while (i != end)
|
||||
// wait for data
|
||||
struct timeval timeout;
|
||||
timeout.tv_sec = timeoutMsec / 1000;
|
||||
timeout.tv_usec = (timeoutMsec % 1000) * 1000;
|
||||
int selectResult = select(maxSock + 1, &rdset, NULL, NULL, &timeout);
|
||||
if (!IS_VALID_SELECT(selectResult))
|
||||
{
|
||||
if (FD_ISSET(i->first, &rdset))
|
||||
throw ServerException(ERR_SOCK_SELECT_FAILED, SOCKET_ERRNO());
|
||||
}
|
||||
if (selectResult > 0) // one (or more) of the sockets is readable
|
||||
{
|
||||
// Check which socket is readable, return the first.
|
||||
boost::mutex::scoped_lock lock(m_sessionMapMutex);
|
||||
SessionMap::iterator i = m_sessionMap.begin();
|
||||
SessionMap::iterator end = m_sessionMap.end();
|
||||
|
||||
while (i != end)
|
||||
{
|
||||
retSession = i->second;
|
||||
break;
|
||||
if (FD_ISSET(i->first, &rdset))
|
||||
{
|
||||
retSession = i->second;
|
||||
break;
|
||||
}
|
||||
++i;
|
||||
}
|
||||
++i;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user