766 lines
19 KiB
C++
766 lines
19 KiB
C++
|
|
|
|
#ifndef _WIN32
|
|
|
|
#include <fcntl.h> // fcntl()
|
|
#include <netinet/tcp.h> // TCP_NODELAY
|
|
|
|
#endif
|
|
|
|
#include "stringex.h"
|
|
#include "Defines.h"
|
|
#include "Common.h"
|
|
#include "Object.h"
|
|
#include "String.h"
|
|
#include "Logger.h"
|
|
#include "Socket.h"
|
|
|
|
#ifndef SINGLE_THREADED_MODEL
|
|
#include "Thread.h"
|
|
#endif
|
|
#include "Timer.h"
|
|
|
|
#include "Global.h"
|
|
#include "MultiplexSocket.h"
|
|
|
|
|
|
#ifdef _WIN32
|
|
|
|
typedef int socklen_t;
|
|
|
|
#else
|
|
|
|
#define SD_BOTH SHUT_RDWR
|
|
#define FD_READ SHUT_RD
|
|
#define FD_WRITE SHUT_WR
|
|
#define closesocket close
|
|
|
|
#endif
|
|
|
|
|
|
|
|
HMultiplexSocket::HMultiplexSocket()
|
|
{
|
|
m_fdSocket = INVALID_SOCKET;
|
|
m_nState = MPS_DISABLE;
|
|
m_bSendRetry = false;
|
|
m_nLastReceivedTime = 0;
|
|
m_pDispatcher = NULL;
|
|
|
|
m_nConnectTime = g_nTcpConnectTimeout;
|
|
|
|
memset(&m_saiRemote, 0, sizeof(m_saiRemote));
|
|
}
|
|
|
|
HMultiplexSocket::HMultiplexSocket(SOCKET sock, struct sockaddr_in * psaiRemote, HMpsDispatcher * pDispatcher)
|
|
{
|
|
m_fdSocket = sock;
|
|
m_nState = MPS_READY;
|
|
m_bSendRetry = false;
|
|
m_nLastReceivedTime = 0;
|
|
m_pDispatcher = pDispatcher;
|
|
|
|
m_nConnectTime = g_nTcpConnectTimeout;
|
|
|
|
memcpy(&m_saiRemote, psaiRemote, sizeof(m_saiRemote));
|
|
}
|
|
|
|
HMultiplexSocket::~HMultiplexSocket()
|
|
{
|
|
Close();
|
|
}
|
|
|
|
bool HMultiplexSocket::Close()
|
|
{
|
|
if (m_nState == MPS_DISABLE ||
|
|
(m_nState == MPS_DELETE && m_fdSocket == INVALID_SOCKET))
|
|
//if (m_fdSocket == INVALID_SOCKET)
|
|
return false;
|
|
|
|
OnClose();
|
|
|
|
if (m_nState != MPS_LISTENING)
|
|
shutdown(m_fdSocket, SD_BOTH);
|
|
|
|
if (m_nState != MPS_DELETE)
|
|
m_nState = MPS_DISABLE;
|
|
|
|
closesocket(m_fdSocket);
|
|
|
|
HLOGF(HLOG_INFO5, "INFO, MpSock, socket closed, socket = %d\n", m_fdSocket);
|
|
|
|
m_fdSocket = INVALID_SOCKET;
|
|
|
|
return true;
|
|
}
|
|
|
|
bool HMultiplexSocket::Connect(unsigned nConnectTime)
|
|
{
|
|
if (m_fdSocket == INVALID_SOCKET)
|
|
{
|
|
HLOGF(HLOG_WARNING, "[WARN], MpSock, This socket did not created.\n");
|
|
return false;
|
|
}
|
|
|
|
m_nLastTryConnectTime = 0;
|
|
|
|
if (nConnectTime)
|
|
m_nConnectTime = nConnectTime;
|
|
|
|
int nResult = connect(m_fdSocket, (struct sockaddr *)&m_saiRemote, sizeof(struct sockaddr));
|
|
if (nResult == SOCKET_ERROR &&
|
|
#ifdef _WIN32
|
|
SockErrorNo() != WSAEWOULDBLOCK)
|
|
#else
|
|
SockErrorNo() != EWOULDBLOCK && SockErrorNo() != EINPROGRESS)
|
|
#endif
|
|
{
|
|
HLOGF(HLOG_WARNING, "[WARN], MpSock, Failed to connect. socket = %d, error = %d\n", m_fdSocket, SockErrorNo());
|
|
OnConnect(false);
|
|
m_nState = MPS_CLOSING;
|
|
return 0;
|
|
}
|
|
|
|
HLOGF(HLOG_INFO4, "INFO, MpSock, TCP connecting... to %s:%d (socket = %d)\n", inet_ntoa(m_saiRemote.sin_addr), ntohs(m_saiRemote.sin_port), m_fdSocket);
|
|
m_nState = MPS_CONNECTING;
|
|
|
|
return true;
|
|
}
|
|
|
|
bool HMultiplexSocket::IsConnectTimeout(DWORD nCurrTime)
|
|
{
|
|
if (m_nConnectTime <= 0)
|
|
return false;
|
|
|
|
if (m_nLastTryConnectTime == 0)
|
|
m_nLastTryConnectTime = nCurrTime;
|
|
|
|
if (m_nConnectTime <= TimeDiff(m_nLastTryConnectTime, nCurrTime))
|
|
return true;
|
|
|
|
return false;
|
|
}
|
|
|
|
bool HMultiplexSocket::Bind(struct sockaddr_in * psaiLocal)
|
|
{
|
|
if (m_fdSocket == INVALID_SOCKET)
|
|
{
|
|
HLOGF(HLOG_ERROR, "[ERROR], MpSock, This socket did not created.\n");
|
|
return false;
|
|
}
|
|
|
|
if (bind(m_fdSocket, (struct sockaddr *)psaiLocal, sizeof(struct sockaddr)) == SOCKET_ERROR)
|
|
{
|
|
HLOGF(HLOG_ERROR, "[ERROR], MpSock, Failed to bind, (%s:%d), error = %d\n", inet_ntoa(psaiLocal->sin_addr), ntohs(psaiLocal->sin_port), SockErrorNo());
|
|
return false;
|
|
}
|
|
|
|
HLOGF(HLOG_INFO4, "INFO, MpSock, Succeed to bind, (%s:%d), socket = %d\n", inet_ntoa(psaiLocal->sin_addr), ntohs(psaiLocal->sin_port), m_fdSocket);
|
|
|
|
return true;
|
|
}
|
|
|
|
bool HMultiplexSocket::Bind(const char * pszAddress, u_short nPort)
|
|
{
|
|
if (m_fdSocket == INVALID_SOCKET)
|
|
{
|
|
HLOGF(HLOG_ERROR, "[ERROR], MpSock, This socket did not created.\n");
|
|
return false;
|
|
}
|
|
|
|
struct sockaddr_in sai;
|
|
|
|
sai.sin_family = AF_INET;
|
|
sai.sin_addr.s_addr = inet_addr(pszAddress);
|
|
sai.sin_port = htons(nPort);
|
|
|
|
return HMultiplexSocket::Bind(&sai);
|
|
}
|
|
|
|
bool HMultiplexSocket::GetRemoteAddress(char * pBuf, u_short * pnPort)
|
|
{
|
|
if (pBuf)
|
|
strcpy(pBuf, inet_ntoa(m_saiRemote.sin_addr));
|
|
|
|
if (pnPort)
|
|
*pnPort = ntohs(m_saiRemote.sin_port);
|
|
|
|
return true;
|
|
}
|
|
|
|
bool HMultiplexSocket::SetRemoteAddress(struct sockaddr_in * psaiLocal)
|
|
{
|
|
memcpy(&m_saiRemote, psaiLocal, sizeof(struct sockaddr_in));
|
|
return true;
|
|
}
|
|
|
|
bool HMultiplexSocket::SetRemoteAddress(const char * pszAddress, u_short nPort)
|
|
{
|
|
m_saiRemote.sin_family = AF_INET;
|
|
m_saiRemote.sin_addr.s_addr = inet_addr(pszAddress);
|
|
m_saiRemote.sin_port = htons(nPort);
|
|
return true;
|
|
}
|
|
|
|
bool HMultiplexSocket::SetRemoteAddressHostByName(const char * pszHostName, u_short nPort)
|
|
{
|
|
struct hostent * pHE;
|
|
|
|
// resolve host address
|
|
if ((pHE = gethostbyname(pszHostName)) == NULL)
|
|
{
|
|
HLOGF(HLOG_ERROR, "[ERROR], MpSock, Failed to gethostbyname, socket = %d, error = %d\n", m_fdSocket, SockErrorNo());
|
|
#if _MSC_VER <= 1200
|
|
return false;
|
|
#else
|
|
char szPort[16];
|
|
struct addrinfo aiHints;
|
|
struct addrinfo *aiList = NULL;
|
|
struct addrinfo *aiPtr = NULL;
|
|
int retVal;
|
|
|
|
memset(&aiHints, 0, sizeof(aiHints));
|
|
aiHints.ai_family = AF_INET;
|
|
aiHints.ai_socktype = SOCK_STREAM;
|
|
aiHints.ai_protocol = IPPROTO_TCP;
|
|
|
|
sprintf(szPort, "%d", nPort);
|
|
|
|
if ((retVal = getaddrinfo(pszHostName, szPort, &aiHints, &aiList)) != 0)
|
|
{
|
|
HLOGF(HLOG_ERROR, "[ERROR], MpSock, Failed to getaddrinfo, socket = %d, error = %d, retVal = %d\n", m_fdSocket, SockErrorNo(), retVal);
|
|
return false;
|
|
}
|
|
|
|
for(aiPtr = aiList; aiPtr != NULL; aiPtr=aiPtr->ai_next)
|
|
{
|
|
if (aiPtr->ai_family == AF_INET)
|
|
{
|
|
memcpy(&m_saiRemote, aiPtr->ai_addr, sizeof(struct sockaddr_in));
|
|
freeaddrinfo(aiList);
|
|
return true;
|
|
}
|
|
}
|
|
|
|
return false;
|
|
#endif
|
|
}
|
|
|
|
memcpy(&m_saiRemote.sin_addr, pHE->h_addr_list[0], pHE->h_length);
|
|
m_saiRemote.sin_port = htons(nPort); // short, network byte order
|
|
m_saiRemote.sin_family = AF_INET; // host byte order
|
|
|
|
return true;
|
|
}
|
|
|
|
void HMultiplexSocket::OnConnect(bool /* bIsSuccess */)
|
|
{
|
|
HLOGF(HLOG_WARNING, "[WARN], MpSock, Called to OnConnect of parent class, socket = %d\n", m_fdSocket);
|
|
}
|
|
|
|
void HMultiplexSocket::OnReceive()
|
|
{
|
|
HLOGF(HLOG_WARNING, "[WARN], MpSock, Called to OnReceive of parent class, socket = %d\n", m_fdSocket);
|
|
}
|
|
|
|
int HMultiplexSocket::OnSendRetry()
|
|
{
|
|
HLOGF(HLOG_WARNING, "[WARN], MpSock, Called to OnReceive of parent class, socket = %d\n", m_fdSocket);
|
|
return -1;
|
|
}
|
|
|
|
//////////////////////////////////////////////////////////////////////////////
|
|
|
|
bool HMultiplexTcpSocket::Create(bool bIsTcpNoDelay)
|
|
{
|
|
if (m_fdSocket != INVALID_SOCKET)
|
|
{
|
|
HLOGF(HLOG_ERROR, "[ERROR], MpSock, This socket aleady created. socket = %d\n", m_fdSocket);
|
|
return false;
|
|
}
|
|
|
|
if ((m_fdSocket = socket(AF_INET, SOCK_STREAM, 0)) == INVALID_SOCKET)
|
|
{
|
|
HLOGF(HLOG_ERROR, "[ERROR], MpSock, Don't create socket. error = %d\n", SockErrorNo());
|
|
return false;
|
|
}
|
|
|
|
// Set non blocking mode
|
|
#ifdef _WIN32
|
|
ULONG nNonBlock = 1;
|
|
if (ioctlsocket(m_fdSocket, FIONBIO, &nNonBlock) == SOCKET_ERROR)
|
|
{
|
|
m_fdSocket = INVALID_SOCKET;
|
|
HLOGF(HLOG_ERROR, "[ERROR], MpSock, Failed to ioctlsocket. socket = %d, WSAGetLastError = %d\n", m_fdSocket, SockErrorNo());
|
|
return false;
|
|
}
|
|
#else
|
|
int nFlags = fcntl(m_fdSocket, F_GETFL, 0);
|
|
fcntl(m_fdSocket, F_SETFL, nFlags | O_NONBLOCK);
|
|
#endif
|
|
|
|
int nBuffSize = MPS_MAX_SOCKET_BUFFER;
|
|
if (setsockopt(m_fdSocket, SOL_SOCKET, SO_SNDBUF, (const char *)&nBuffSize, sizeof(nBuffSize)) == SOCKET_ERROR)
|
|
{
|
|
HLOGF(HLOG_ERROR, "[ERROR], MpSock, Failed to setsockopt SO_SNDBUF. socket = %d, error = %d\n", m_fdSocket, SockErrorNo());
|
|
}
|
|
|
|
if (setsockopt(m_fdSocket, SOL_SOCKET, SO_RCVBUF, (const char *)&nBuffSize, sizeof(nBuffSize)) == SOCKET_ERROR)
|
|
{
|
|
HLOGF(HLOG_ERROR, "[ERROR], MpSock, Failed to setsockopt SO_RCVBUF. socket = %d, error = %d\n", m_fdSocket, SockErrorNo());
|
|
}
|
|
|
|
if (bIsTcpNoDelay) //... linux error
|
|
{
|
|
bool bNoDelay = false;
|
|
if (setsockopt(m_fdSocket, IPPROTO_TCP, TCP_NODELAY, (const char *)&bNoDelay, sizeof(bNoDelay)) == SOCKET_ERROR)
|
|
{
|
|
HLOGF(HLOG_ERROR, "[ERROR], MpSock, Failed setsockopt TCP_NODELAY. socket = %d, error = %d\n", m_fdSocket, SockErrorNo());
|
|
}
|
|
}
|
|
|
|
int nKeepAlive = 1;
|
|
if (setsockopt(m_fdSocket, SOL_SOCKET, SO_KEEPALIVE, (const char *)&nKeepAlive, sizeof(nKeepAlive)) == SOCKET_ERROR)
|
|
{
|
|
HLOGF(HLOG_ERROR, "[ERROR], MpSock, Failed to setsockopt SO_KEEPALIVE. socket = %d, error = %d\n", m_fdSocket, SockErrorNo());
|
|
}
|
|
|
|
HLOGF(HLOG_INFO5, "INFO, MpSock, TCP socket created. socket = %d\n", m_fdSocket);
|
|
|
|
return true;
|
|
}
|
|
|
|
int HMultiplexTcpSocket::Send(const char * pData, int nSize)
|
|
{
|
|
int nSentSize = send(m_fdSocket, pData, nSize, 0);
|
|
if (nSentSize == SOCKET_ERROR)
|
|
{
|
|
int nError = SockErrorNo();
|
|
|
|
#ifdef _WIN32
|
|
if (nError != WSAEWOULDBLOCK)
|
|
#else
|
|
if (nError != EWOULDBLOCK && nError != EAGAIN)
|
|
#endif
|
|
{
|
|
HLOGF(HLOG_WARNING, "[WARN], MpSock, Failed to send, socket = %d, error = %d\n", m_fdSocket, nError);
|
|
m_nState = MPS_CLOSING;
|
|
return -1;
|
|
}
|
|
|
|
m_bSendRetry = true;
|
|
return 0;
|
|
}
|
|
else if (nSentSize < nSize)
|
|
{
|
|
m_bSendRetry = true;
|
|
return nSentSize;
|
|
}
|
|
|
|
if (m_bSendRetry)
|
|
m_bSendRetry = false;
|
|
|
|
return nSentSize;
|
|
}
|
|
|
|
int HMultiplexTcpSocket::Receive(char * pBuffer, int nSize)
|
|
{
|
|
int nReceivedSize = recv(m_fdSocket, pBuffer, nSize, 0);
|
|
if (nReceivedSize == SOCKET_ERROR)
|
|
{
|
|
HLOGF(HLOG_INFO4, "INFO, MpSock, The virtual circuit was reset by the remote side executing a \"hard\" or \"abortive\" close, socket = %d, error = %d\n", m_fdSocket, SockErrorNo());
|
|
m_nState = MPS_CLOSING;
|
|
return SOCKET_ERROR;
|
|
}
|
|
else if (nReceivedSize == 0)
|
|
{
|
|
HLOGF(HLOG_INFO4, "INFO, MpSock, The virtual circuit was reset by the remote side executing a \"gracefull\" shutdown, socket = %d\n", m_fdSocket);
|
|
m_nState = MPS_CLOSING;
|
|
return 0;
|
|
}
|
|
|
|
return nReceivedSize;
|
|
}
|
|
|
|
void HMultiplexTcpSocket::FinishTcp()
|
|
{
|
|
shutdown(m_fdSocket, FD_WRITE);
|
|
|
|
HLOGF(HLOG_INFO5, "INFO, MPSock, write shutdown, socket = %d\n", m_fdSocket);
|
|
m_nState = MPS_CLOSING;
|
|
}
|
|
|
|
//////////////////////////////////////////////////////////////////////////////
|
|
|
|
bool HMultiplexUdpSocket::Create()
|
|
{
|
|
if (m_fdSocket != INVALID_SOCKET)
|
|
{
|
|
HLOGF(HLOG_ERROR, "[ERROR], MpSock, This socket aleady created. socket = %d\n", m_fdSocket);
|
|
return false;
|
|
}
|
|
|
|
if ((m_fdSocket = socket(AF_INET, SOCK_DGRAM, 0)) == INVALID_SOCKET)
|
|
{
|
|
HLOGF(HLOG_ERROR, "[ERROR], MpSock, Don't create socket. error = %d\n", SockErrorNo());
|
|
return false;
|
|
}
|
|
|
|
// Set non blocking mode
|
|
#ifdef _WIN32
|
|
ULONG nNonBlock = 1;
|
|
if (ioctlsocket(m_fdSocket, FIONBIO, &nNonBlock) == SOCKET_ERROR)
|
|
{
|
|
m_fdSocket = INVALID_SOCKET;
|
|
HLOGF(HLOG_ERROR, "[ERROR], MpSock, Failed to ioctlsocket. socket = %d, WSAGetLastError = %d\n", m_fdSocket, SockErrorNo());
|
|
return false;
|
|
}
|
|
#else
|
|
int nFlags = fcntl(m_fdSocket, F_GETFL, 0);
|
|
fcntl(m_fdSocket, F_SETFL, nFlags | O_NONBLOCK);
|
|
#endif
|
|
|
|
int nBuffSize = MPS_MAX_SOCKET_BUFFER;
|
|
if (setsockopt(m_fdSocket, SOL_SOCKET, SO_SNDBUF, (const char *)&nBuffSize, sizeof(nBuffSize)) == SOCKET_ERROR)
|
|
{
|
|
HLOGF(HLOG_ERROR, "[ERROR], MpSock, Failed to setsockopt SO_SNDBUF. socket = %d, error = %d\n", m_fdSocket, SockErrorNo());
|
|
}
|
|
|
|
if (setsockopt(m_fdSocket, SOL_SOCKET, SO_RCVBUF, (const char *)&nBuffSize, sizeof(nBuffSize)) == SOCKET_ERROR)
|
|
{
|
|
HLOGF(HLOG_ERROR, "[ERROR], MpSock, Failed to setsockopt SO_RCVBUF. socket = %d, error = %d\n", m_fdSocket, SockErrorNo());
|
|
}
|
|
|
|
HLOGF(HLOG_INFO5, "INFO, MpSock, UDP socket created. socket = %d\n", m_fdSocket);
|
|
|
|
return true;
|
|
}
|
|
|
|
bool HMultiplexUdpSocket::Bind(struct sockaddr_in * psaiLocal)
|
|
{
|
|
if (HMultiplexSocket::Bind(psaiLocal) == false)
|
|
return false;
|
|
|
|
m_nState = MPS_READY;
|
|
|
|
return true;
|
|
}
|
|
|
|
bool HMultiplexUdpSocket::SetRemoteAddress(struct sockaddr_in * psaiLocal)
|
|
{
|
|
memcpy(&m_saiRemote, psaiLocal, sizeof(struct sockaddr_in));
|
|
m_nState = MPS_READY;
|
|
return true;
|
|
}
|
|
|
|
bool HMultiplexUdpSocket::SetRemoteAddress(const char * pszAddress, u_short nPort)
|
|
{
|
|
m_saiRemote.sin_family = AF_INET;
|
|
m_saiRemote.sin_addr.s_addr = inet_addr(pszAddress);
|
|
m_saiRemote.sin_port = htons(nPort);
|
|
m_nState = MPS_READY;
|
|
return true;
|
|
}
|
|
|
|
//////////////////////////////////////////////////////////////////////////////
|
|
|
|
HMpsDispatcher::HMpsDispatcher()
|
|
{
|
|
m_nLoopInterval = 0;
|
|
m_bTerminate = false;
|
|
|
|
#ifndef SINGLE_THREADED_MODEL
|
|
m_bProecessing = false;
|
|
#endif
|
|
}
|
|
|
|
HMpsDispatcher::~HMpsDispatcher()
|
|
{
|
|
RemoveAllSocket();
|
|
}
|
|
|
|
bool HMpsDispatcher::Attach(HMultiplexSocket * pSocket)
|
|
{
|
|
pSocket->SetDispatcher(this);
|
|
m_listSocket.push_front(pSocket);
|
|
return true;
|
|
}
|
|
|
|
bool HMpsDispatcher::AttachWithLock(HMultiplexSocket * pSocket)
|
|
{
|
|
pSocket->SetDispatcher(this);
|
|
|
|
#ifndef SINGLE_THREADED_MODEL
|
|
m_syncObject.Lock();
|
|
#endif
|
|
|
|
m_listSocket.push_front(pSocket);
|
|
|
|
#ifndef SINGLE_THREADED_MODEL
|
|
m_syncObject.Unlock();
|
|
#endif
|
|
|
|
return true;
|
|
}
|
|
|
|
bool HMpsDispatcher::Detach(HMultiplexSocket * pSocket)
|
|
{
|
|
pSocket->m_pDispatcher = NULL;
|
|
pSocket->SetState(HMultiplexSocket::MPS_DELETE);
|
|
|
|
return true;
|
|
}
|
|
|
|
void HMpsDispatcher::RemoveAllSocket()
|
|
{
|
|
if (m_listSocket.empty())
|
|
return;
|
|
|
|
#ifndef SINGLE_THREADED_MODEL
|
|
m_syncObject.Lock();
|
|
#endif
|
|
|
|
for (LIST_MULTIPLEX_SOCKET::iterator iter = m_listSocket.begin();
|
|
iter != m_listSocket.end(); iter++)
|
|
{
|
|
HMultiplexSocket * pSocket = *iter;
|
|
if (pSocket)
|
|
{
|
|
pSocket->Close();
|
|
delete pSocket;
|
|
}
|
|
}
|
|
m_listSocket.clear();
|
|
|
|
#ifndef SINGLE_THREADED_MODEL
|
|
m_syncObject.Unlock();
|
|
#endif
|
|
}
|
|
|
|
bool HMpsDispatcher::Dispatch()
|
|
{
|
|
HMultiplexSocket * pSocket;
|
|
fd_set fdSetRead, fdSetWrite;//, fdSetError;
|
|
int nSelectResult, nFDMax, nFDCount;
|
|
LIST_MULTIPLEX_SOCKET::iterator iter;
|
|
struct timeval tvWait;
|
|
|
|
long nTimeout = m_nLoopInterval * 1000; // msec
|
|
m_bTerminate = false;
|
|
|
|
while (m_bTerminate == false)
|
|
{
|
|
FD_ZERO(&fdSetRead);
|
|
FD_ZERO(&fdSetWrite);
|
|
//FD_ZERO(&fdSetError);
|
|
|
|
tvWait.tv_sec = 0;
|
|
tvWait.tv_usec = nTimeout;
|
|
|
|
nFDMax = 0;
|
|
nFDCount = 0;
|
|
|
|
#ifndef SINGLE_THREADED_MODEL
|
|
m_bProecessing = true;
|
|
m_syncObject.Lock();
|
|
#endif
|
|
for (iter = m_listSocket.begin(); iter != m_listSocket.end(); iter++)
|
|
{
|
|
pSocket = *iter;
|
|
|
|
if (pSocket->GetState() == HMultiplexSocket::MPS_CLOSING ||
|
|
pSocket->GetState() == HMultiplexSocket::MPS_DELETE)
|
|
{
|
|
pSocket->Close();
|
|
|
|
if (pSocket->GetState() == HMultiplexSocket::MPS_DELETE)
|
|
{
|
|
delete pSocket;
|
|
iter = m_listSocket.erase(iter);
|
|
|
|
if (m_listSocket.empty() || iter == m_listSocket.end())
|
|
break;
|
|
|
|
pSocket = *iter;
|
|
}
|
|
}
|
|
|
|
if (pSocket->GetState() == HMultiplexSocket::MPS_READY ||
|
|
pSocket->GetState() == HMultiplexSocket::MPS_LISTENING)
|
|
{
|
|
FD_SET(pSocket->GetSocket(), &fdSetRead);
|
|
//FD_SET(pSocket->GetSocket(), &fdSetError);
|
|
nFDCount++;
|
|
#ifndef _WIN32
|
|
if (nFDMax < pSocket->GetSocket())
|
|
nFDMax = pSocket->GetSocket();
|
|
#endif
|
|
}
|
|
else if (pSocket->GetState() == HMultiplexSocket::MPS_CONNECTING)
|
|
{
|
|
if (pSocket->IsConnectTimeout(HTIMER()->GetTickCount()))
|
|
{
|
|
HLOGF(HLOG_WARNING, "[WARN], MpsDisp, Failed to connect because of timeout. socket = %d\n", pSocket->GetSocket());
|
|
pSocket->OnConnect(false);
|
|
pSocket->SetState(HMultiplexSocket::MPS_CLOSING);
|
|
}
|
|
else
|
|
{
|
|
FD_SET(pSocket->GetSocket(), &fdSetRead);
|
|
FD_SET(pSocket->GetSocket(), &fdSetWrite);
|
|
//FD_SET(pSocket->GetSocket(), &fdSetError);
|
|
nFDCount++;
|
|
#ifndef _WIN32
|
|
if (nFDMax < pSocket->GetSocket())
|
|
nFDMax = pSocket->GetSocket();
|
|
#endif
|
|
}
|
|
}
|
|
|
|
if (pSocket->IsSendRetry())
|
|
{
|
|
FD_SET(pSocket->GetSocket(), &fdSetWrite);
|
|
nFDCount++;
|
|
#ifndef _WIN32
|
|
if (nFDMax < pSocket->GetSocket())
|
|
nFDMax = pSocket->GetSocket();
|
|
#endif
|
|
}
|
|
}
|
|
|
|
#ifndef SINGLE_THREADED_MODEL
|
|
m_bProecessing = false;
|
|
m_syncObject.Unlock();
|
|
#endif
|
|
|
|
if (m_listSocket.empty() || nFDCount == 0)
|
|
{
|
|
Sleep(10);
|
|
continue;
|
|
}
|
|
|
|
//if ((nSelectResult = select(nFDMax + 1, &fdSetRead, &fdSetWrite, &fdSetError, &tvWait)) == SOCKET_ERROR)
|
|
if ((nSelectResult = select(nFDMax + 1, &fdSetRead, &fdSetWrite, NULL, m_nLoopInterval ? &tvWait : NULL)) == SOCKET_ERROR)
|
|
{
|
|
Sleep(10);
|
|
continue;
|
|
}
|
|
//printf("nSelectResult = %d\n", nSelectResult);
|
|
|
|
if (nSelectResult <= 0)
|
|
continue;
|
|
|
|
#ifndef SINGLE_THREADED_MODEL
|
|
m_bProecessing = true;
|
|
m_syncObject.Lock();
|
|
#endif
|
|
for (iter = m_listSocket.begin(); iter != m_listSocket.end(); iter++)
|
|
{
|
|
pSocket = *iter;
|
|
|
|
if (pSocket->GetState() == HMultiplexSocket::MPS_DISABLE ||
|
|
pSocket->GetState() == HMultiplexSocket::MPS_CLOSING ||
|
|
pSocket->GetState() == HMultiplexSocket::MPS_DELETE)
|
|
{
|
|
continue;
|
|
}
|
|
|
|
if (FD_ISSET(pSocket->GetSocket(), &fdSetRead) != 0)
|
|
{
|
|
if (pSocket->GetState() == HMultiplexSocket::MPS_READY)
|
|
{
|
|
pSocket->OnReceive();
|
|
pSocket->m_nLastReceivedTime = HTIMER()->GetTickCount();
|
|
}
|
|
else if (pSocket->GetState() == HMultiplexSocket::MPS_LISTENING)
|
|
{
|
|
SOCKET sock = INVALID_SOCKET;
|
|
socklen_t len = sizeof(struct sockaddr);
|
|
struct sockaddr_in sai;
|
|
|
|
if ((sock = accept(pSocket->GetSocket(), (struct sockaddr *)&sai, &len)) == INVALID_SOCKET)
|
|
{
|
|
HLOGF(HLOG_WARNING, "[WARN], MpsDisp, Failed to accept, socket = %d, error = %d\n", pSocket->GetSocket(), SockErrorNo());
|
|
Sleep(50);
|
|
}
|
|
else
|
|
{
|
|
char szListenIp[32];
|
|
u_short nListenPort = 0;
|
|
pSocket->GetLocalAddress(szListenIp, &nListenPort);
|
|
|
|
HLOGF(HLOG_INFO4, "INFO, MpsDisp, Succeeded to accept, (%s:%d), (%s:%d), socket = %d\n", szListenIp, nListenPort, inet_ntoa(sai.sin_addr), ntohs(sai.sin_port), sock);
|
|
|
|
HMultiplexSocket * pAcceptedSocket = pSocket->OnAccept(sock, &sai);
|
|
if (pAcceptedSocket)
|
|
{
|
|
// Set non blocking mode
|
|
#ifdef _WIN32
|
|
ULONG nNonBlock = 1;
|
|
if (ioctlsocket(pAcceptedSocket->GetSocket(), FIONBIO, &nNonBlock) == SOCKET_ERROR)
|
|
HLOGF(HLOG_ERROR, "[ERROR], MpSock, Failed to ioctlsocket, socket = %d, WSAGetLastError = %d\n", pAcceptedSocket->GetSocket(), SockErrorNo());
|
|
#else
|
|
int nFlags = fcntl(pAcceptedSocket->GetSocket(), F_GETFL, 0);
|
|
fcntl(pAcceptedSocket->GetSocket(), F_SETFL, nFlags | O_NONBLOCK);
|
|
#endif
|
|
m_listSocket.push_front(pAcceptedSocket);
|
|
|
|
pAcceptedSocket->m_nLastReceivedTime = HTIMER()->GetTickCount();
|
|
pAcceptedSocket->OnAccepted();
|
|
}
|
|
else
|
|
closesocket(sock);
|
|
}
|
|
}
|
|
|
|
FD_CLR(pSocket->GetSocket(), &fdSetRead);
|
|
|
|
if (--nSelectResult <= 0)
|
|
break;
|
|
}
|
|
|
|
if (FD_ISSET(pSocket->GetSocket(), &fdSetWrite) != 0)
|
|
{
|
|
if (pSocket->GetState() == HMultiplexSocket::MPS_CONNECTING)
|
|
{
|
|
int nSocketError = pSocket->GetSockOptError();
|
|
if (nSocketError != 0)
|
|
{
|
|
HLOGF(HLOG_WARNING, "[WARN], MpsDisp, Failed to connect (%s:%d), socket = %d, error = %d\n",
|
|
inet_ntoa(pSocket->GetRemoteAddress()->sin_addr), ntohs(pSocket->GetRemoteAddress()->sin_port), pSocket->GetSocket(), nSocketError);
|
|
|
|
pSocket->OnConnect(false);
|
|
pSocket->SetState(HMultiplexSocket::MPS_CLOSING);
|
|
}
|
|
else
|
|
{
|
|
HLOGF(HLOG_INFO4, "INFO, MpsDisp, Succeeded to connect (%s:%d) socket = %d\n",
|
|
inet_ntoa(pSocket->GetRemoteAddress()->sin_addr), ntohs(pSocket->GetRemoteAddress()->sin_port), pSocket->GetSocket());
|
|
|
|
pSocket->SetState(HMultiplexSocket::MPS_READY);
|
|
pSocket->OnConnect(true);
|
|
pSocket->m_nLastReceivedTime = HTIMER()->GetTickCount();
|
|
}
|
|
}
|
|
else if (pSocket->GetState() == HMultiplexSocket::MPS_READY)
|
|
{
|
|
if (pSocket->OnSendRetry() < 0) // remain packet send
|
|
pSocket->FinishTcp();
|
|
}
|
|
|
|
FD_CLR(pSocket->GetSocket(), &fdSetWrite);
|
|
|
|
if (--nSelectResult <= 0)
|
|
break;
|
|
}
|
|
}
|
|
|
|
#ifndef SINGLE_THREADED_MODEL
|
|
m_bProecessing = false;
|
|
m_syncObject.Unlock();
|
|
#endif
|
|
}
|
|
|
|
return false;
|
|
}
|