Fixed sender again, as it did not work when sending errors.

This commit is contained in:
lotodore
2007-10-28 15:06:57 +00:00
parent e951bda1db
commit cdea81b2ba
12 changed files with 104 additions and 148 deletions
+49 -58
View File
@@ -28,7 +28,7 @@ using namespace std;
SenderThread::SenderThread(SenderCallback &cb)
: m_curSession(INVALID_SESSION), m_tmpOutBufSize(0), m_callback(cb)
: m_tmpOutBufSize(0), m_callback(cb)
{
}
@@ -37,9 +37,9 @@ SenderThread::~SenderThread()
}
void
SenderThread::Send(SessionId session, boost::shared_ptr<NetPacket> packet)
SenderThread::Send(boost::shared_ptr<SessionData> session, boost::shared_ptr<NetPacket> packet)
{
if (packet.get() && session != INVALID_SESSION)
if (packet.get() && session.get())
{
boost::mutex::scoped_lock lock(m_outBufMutex);
InternalStore(m_outBuf, SEND_QUEUE_SIZE, session, packet);
@@ -47,9 +47,9 @@ SenderThread::Send(SessionId session, boost::shared_ptr<NetPacket> packet)
}
void
SenderThread::Send(SessionId session, const NetPacketList &packetList)
SenderThread::Send(boost::shared_ptr<SessionData> session, const NetPacketList &packetList)
{
if (!packetList.empty() && session != INVALID_SESSION)
if (!packetList.empty() && session.get())
{
boost::mutex::scoped_lock lock(m_outBufMutex);
InternalStore(m_outBuf, SEND_QUEUE_SIZE, session, packetList);
@@ -57,9 +57,9 @@ SenderThread::Send(SessionId session, const NetPacketList &packetList)
}
void
SenderThread::SendLowPrio(SessionId session, boost::shared_ptr<NetPacket> packet)
SenderThread::SendLowPrio(boost::shared_ptr<SessionData> session, boost::shared_ptr<NetPacket> packet)
{
if (packet.get() && session != INVALID_SESSION)
if (packet.get() && session.get())
{
boost::mutex::scoped_lock lock(m_lowPrioOutBufMutex);
InternalStore(m_lowPrioOutBuf, SEND_LOW_PRIO_QUEUE_SIZE, session, packet);
@@ -67,9 +67,9 @@ SenderThread::SendLowPrio(SessionId session, boost::shared_ptr<NetPacket> packet
}
void
SenderThread::SendLowPrio(SessionId session, const NetPacketList &packetList)
SenderThread::SendLowPrio(boost::shared_ptr<SessionData> session, const NetPacketList &packetList)
{
if (!packetList.empty() && session != INVALID_SESSION)
if (!packetList.empty() && session.get())
{
boost::mutex::scoped_lock lock(m_lowPrioOutBufMutex);
InternalStore(m_lowPrioOutBuf, SEND_LOW_PRIO_QUEUE_SIZE, session, packetList);
@@ -77,7 +77,7 @@ SenderThread::SendLowPrio(SessionId session, const NetPacketList &packetList)
}
void
SenderThread::InternalStore(SendDataDeque &sendQueue, unsigned maxQueueSize, SessionId session, boost::shared_ptr<NetPacket> packet)
SenderThread::InternalStore(SendDataDeque &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(std::make_pair(packet, session));
@@ -85,7 +85,7 @@ SenderThread::InternalStore(SendDataDeque &sendQueue, unsigned maxQueueSize, Ses
}
void
SenderThread::InternalStore(SendDataDeque &sendQueue, unsigned maxQueueSize, SessionId session, const NetPacketList &packetList)
SenderThread::InternalStore(SendDataDeque &sendQueue, unsigned maxQueueSize, boost::shared_ptr<SessionData> session, const NetPacketList &packetList)
{
if (sendQueue.size() + packetList.size() < maxQueueSize)
{
@@ -134,7 +134,7 @@ SenderThread::Main()
if (tmpData.first.get())
{
if (tmpData.second != INVALID_SESSION)
if (tmpData.second.get())
m_curSession = tmpData.second;
u_int16_t tmpLen = tmpData.first->GetLen();
@@ -147,25 +147,37 @@ SenderThread::Main()
}
if (m_tmpOutBufSize)
{
SOCKET tmpSocket;
if (!m_callback.GetSocketForSession(m_curSession, tmpSocket))
SOCKET tmpSocket = m_curSession->GetSocket();
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))
{
// Invalid session - skip.
m_tmpOutBufSize = 0;
m_curSession = INVALID_SESSION;
// 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);
m_tmpOutBufSize = 0;
m_curSession.reset();
}
Msleep(SEND_TIMEOUT_MSEC);
}
else
if (selectResult > 0) // send is possible
{
fd_set writeSet;
struct timeval timeout;
// send next chunk of data
int bytesSent = send(tmpSocket, m_tmpOutBuf, m_tmpOutBufSize, 0);
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))
if (!IS_VALID_SEND(bytesSent))
{
// Never assume that this is a fatal error.
int errCode = SOCKET_ERRNO();
@@ -174,42 +186,21 @@ 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(m_curSession, ERR_SOCK_SELECT_FAILED, errCode);
m_callback.SignalNetError(m_curSession->GetId(), ERR_SOCK_SEND_FAILED, errCode);
m_tmpOutBufSize = 0;
m_curSession = INVALID_SESSION;
m_curSession.reset();
}
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, 0);
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, ERR_SOCK_SEND_FAILED, errCode);
m_tmpOutBufSize = 0;
m_curSession = INVALID_SESSION;
}
Msleep(SEND_TIMEOUT_MSEC);
}
else if ((unsigned)bytesSent < m_tmpOutBufSize)
{
m_tmpOutBufSize -= (unsigned)bytesSent;
memmove(m_tmpOutBuf, m_tmpOutBuf + bytesSent, m_tmpOutBufSize);
}
else
{
m_tmpOutBufSize = 0;
m_curSession = INVALID_SESSION;
}
m_tmpOutBufSize -= (unsigned)bytesSent;
memmove(m_tmpOutBuf, m_tmpOutBuf + bytesSent, m_tmpOutBufSize);
}
else
{
m_tmpOutBufSize = 0;
m_curSession.reset();
}
}
}