Files
smartservice_native/hdssd/huslib/MultiplexSocket.h
T
2026-05-20 03:08:08 +09:00

487 lines
12 KiB
C++

#ifndef MULTIPLEX_SOCKET__H_
#define MULTIPLEX_SOCKET__H_
#ifdef _WIN32
#pragma once
#pragma warning(disable:4244)
#pragma warning(disable:4267)
#endif
#include <list>
class HMpsDispatcher;
class HMultiplexSocket : public HSocket
{
public:
enum
{
MPS_DISABLE = 1,
MPS_READY,
MPS_CONNECTING,
MPS_LISTENING,
MPS_CLOSING,
MPS_DELETE
};
enum
{
MPS_MAX_SOCKET_BUFFER = 16384
};
HMultiplexSocket();
HMultiplexSocket(SOCKET sock, struct sockaddr_in * psaiRemote, HMpsDispatcher * pDispatcher);
~HMultiplexSocket();
inline int GetState() { return m_nState; }
inline void SetState(int nState) { m_nState = nState; }
inline unsigned int GetLastReceivedTime() { return m_nLastReceivedTime; }
inline bool IsSendRetry() { return m_bSendRetry; }
HMpsDispatcher * GetDispatcher() { return m_pDispatcher; }
void SetDispatcher(HMpsDispatcher * pDispatcher) { m_pDispatcher = pDispatcher; }
virtual bool Close();
virtual bool Connect(unsigned nConnectTime = 60000); // msec
bool IsConnectTimeout(DWORD nCurrTime);
bool Bind(struct sockaddr_in * psaiLocal);
bool Bind(const char * pszAddress, u_short nPort);
struct sockaddr_in * GetRemoteAddress() { return &m_saiRemote; }
bool GetRemoteAddress(char * pBuf, u_short * pnPort);
bool SetRemoteAddress(struct sockaddr_in * psaiLocal);
bool SetRemoteAddress(const char * pszAddress, u_short nPort);
bool SetRemoteAddressHostByName(const char * pszHostName, u_short nPort);
virtual int Send(const char * pData, int nSize) = 0;
virtual int Receive(char * pBuffer, int nSize) = 0;
virtual void FinishTcp() {}
protected:
virtual void OnClose() {}
virtual void OnConnect(bool bIsSuccess); // only called for client socket
virtual void OnReceive();
virtual int OnSendRetry();
virtual void OnAccepted() {} // only called for child socket
public:
virtual HMultiplexSocket * OnAccept(SOCKET /* sock */, struct sockaddr_in * /* psaiRemote */) { return NULL; }
protected:
int m_nState;
bool m_bSendRetry;
unsigned int m_nLastReceivedTime;
struct sockaddr_in m_saiRemote;
HMpsDispatcher * m_pDispatcher;
private:
unsigned long m_nConnectTime;
unsigned long m_nLastTryConnectTime;
friend class HMpsDispatcher;
template<class _Ty> friend class HMultiplexListenThread;
};
class HMultiplexTcpSocket : public HMultiplexSocket
{
public:
HMultiplexTcpSocket() {}
HMultiplexTcpSocket(SOCKET sock, struct sockaddr_in * psaiRemote, HMpsDispatcher * pDispatcher) : HMultiplexSocket(sock, psaiRemote, pDispatcher) {}
~HMultiplexTcpSocket() {}
bool Create(bool bIsTcpNoDelay = false);
virtual bool IsSslSocket() { return false; }
virtual int Send(const char * pData, int nSize);
virtual int Receive(char * pBuffer, int nSize);
virtual void FinishTcp(); // remain send and shutdown
};
class HMultiplexUdpSocket : public HMultiplexSocket
{
public:
HMultiplexUdpSocket() {}
~HMultiplexUdpSocket() {}
virtual bool Create();
virtual bool SetRemoteAddress(struct sockaddr_in * psaiLocal);
virtual bool SetRemoteAddress(const char * pszAddress, u_short nPort);
virtual bool Bind(struct sockaddr_in * psaiLocal);
virtual inline int Send(const char * pData, int nSize)
{
return sendto(m_fdSocket, pData, nSize, 0, (struct sockaddr *)&m_saiRemote, sizeof(struct sockaddr));
}
virtual inline int Receive(char * pBuffer, int nSize)
{
return recvfrom(m_fdSocket, pBuffer, nSize, 0, NULL, 0);
}
virtual inline int Receive(char * pBuffer, int nSize, struct sockaddr_in * psaiUDPFrom)
{
socklen_t len = sizeof(struct sockaddr);
return recvfrom(m_fdSocket, pBuffer, nSize, 0, (struct sockaddr *)psaiUDPFrom, &len);
}
};
class HMultiplexListenSocket : public HMultiplexTcpSocket
{
public:
HMultiplexListenSocket() {}
~HMultiplexListenSocket() {}
bool Listen(const char * pszAddress, u_short nPort, int nBackLog = 10)
{
if (Create() == false)
return false;
#ifdef REUSEADDR
int nOptVal = 1;
SetSockOpt(SOL_SOCKET, SO_REUSEADDR, (const char *)&nOptVal, sizeof(nOptVal)); // for multiple ip in UNIX
#endif
if (Bind(pszAddress, nPort) == false)
{
HLOGF(HLOG_ERROR, "[ERROR], MpLstnSock, Failed to bind, (%s:%d), error = %d\n", pszAddress, nPort, SockErrorNo());
return false;
}
if (listen(m_fdSocket, nBackLog) == SOCKET_ERROR)
{
HLOGF(HLOG_ERROR, "[ERROR], MpLstnSock, Failed to listen, (%s:%d), error = %d\n", pszAddress, nPort, SockErrorNo());
return false;
}
m_nState = MPS_LISTENING;
HLOGF(HLOG_INFO1, "INFO, MpLstnSock, Succeed to listen TCP, (%s:%d), socket = %d\n", pszAddress, nPort, m_fdSocket);
return true;
}
private:
virtual HMultiplexTcpSocket * OnAccept(SOCKET sock, struct sockaddr_in * psaiRemote) = 0;
};
typedef std::list<HMultiplexSocket *> LIST_MULTIPLEX_SOCKET;
class HMpsDispatcher
{
public:
HMpsDispatcher();
virtual ~HMpsDispatcher();
void Initialize(unsigned long nLoopInterval) { m_nLoopInterval = nLoopInterval; } // msec, for connect timeout
void SetTerminateFlag() { m_bTerminate = true; }
LIST_MULTIPLEX_SOCKET * GetSocketList() { return &m_listSocket; }
bool Attach(HMultiplexSocket * pSocket);
bool AttachWithLock(HMultiplexSocket * pSocket);
bool Detach(HMultiplexSocket * pSocket);
void RemoveAllSocket();
#ifndef SINGLE_THREADED_MODEL
inline bool IsProcessing() { return m_bProecessing; }
#endif
bool Dispatch();
//private:
public:
LIST_MULTIPLEX_SOCKET m_listSocket;
unsigned long m_nLoopInterval;
volatile bool m_bTerminate;
public: //...TempCode
#ifndef SINGLE_THREADED_MODEL
HSyncObject m_syncObject;
volatile bool m_bProecessing;
#endif
};
class HMpsDispatcherList : public std::list<HMpsDispatcher *>
{
public:
HMpsDispatcherList() { m_nLazyIndex = 0; }
~HMpsDispatcherList() {}
HMpsDispatcher * GetLazyDispatcher()
{
iterator iter;
int i = 0;
for (iter = begin(); iter != end(); iter++)
{
if ((*iter)->IsProcessing() == false)
i++;
}
//HLOGF(HLOG_WARNING, "[TRACE], DispatcherList_0x%x, lazy count = %d\n", this, i);
i = 0;
for (iter = begin(); iter != end(); iter++)
{
if (i == m_nLazyIndex)
{
m_nLazyIndex++;
if (m_nLazyIndex >= (int)size())
m_nLazyIndex = 0;
return *iter;
}
i++;
}
return NULL;
}
public:
volatile int m_nLazyIndex;
};
///////////////////////////////////////
class HMultiplexSocketDispatchThread : public HThread
{
public:
HMultiplexSocketDispatchThread()
: HThread(true, 1024 * 1024 * 4) // stack size : 4MB
{
m_pMpsDispatcher = new HMpsDispatcher;
}
~HMultiplexSocketDispatchThread()
{
delete m_pMpsDispatcher;
}
inline HMpsDispatcher * GetDispatcher() { return m_pMpsDispatcher; }
inline void DispatchLock() { m_pMpsDispatcher->m_syncObject.Lock(); }
inline void DispatchUnlock() { m_pMpsDispatcher->m_syncObject.Unlock(); }
inline unsigned long BeginDispatch(unsigned long nLoopInterval = 50)
{
m_pMpsDispatcher->Initialize(nLoopInterval);
return HThread::Begin();
}
void Terminate()
{
m_pMpsDispatcher->SetTerminateFlag();
}
protected:
void Main() { m_pMpsDispatcher->Dispatch(); }
private:
inline unsigned long Begin() { return 0; } // do not use
private:
HMpsDispatcher * m_pMpsDispatcher;
};
///////////////////////////////////////
template<class _Ty>
class HMultiplexListenThread : public HThread
{
public:
HMultiplexListenThread()
: HThread(true, 0)
{
m_sockListen = INVALID_SOCKET;
m_bTerminate = false;
}
~HMultiplexListenThread() {}
inline void SetDispatcherList(HMpsDispatcherList * pList) { m_pListDispatcher = pList; }
bool Listen(const char * pszAddress, u_short nPort, int nBackLog = 10)
{
if (m_sockListen != INVALID_SOCKET)
{
HLOGF(HLOG_ERROR, "[ERROR], MpLstnThr, This socket aleady created. socket = %d\n", m_sockListen);
return false;
}
if ((m_sockListen = socket(AF_INET, SOCK_STREAM, 0)) == INVALID_SOCKET)
{
HLOGF(HLOG_ERROR, "[ERROR], MpLstnThr, Don't create socket. Error = %d\n", SockErrorNo());
return false;
}
// Set non blocking mode
#ifdef _WIN32
ULONG nNonBlock = 1;
if (ioctlsocket(m_sockListen, FIONBIO, &nNonBlock) == SOCKET_ERROR)
{
m_sockListen = INVALID_SOCKET;
HLOGF(HLOG_ERROR, "[ERROR], MpLstnThrd, Failed to ioctlsocket, socket = %d, error = %d\n", m_sockListen, SockErrorNo());
return false;
}
#else
int nFlags = fcntl(m_sockListen, F_GETFL, 0);
fcntl(m_sockListen, F_SETFL, nFlags | O_NONBLOCK);
#endif
struct sockaddr_in sai;
memset(&sai, 0, sizeof(struct sockaddr_in));
if (pszAddress == NULL || pszAddress[0] == '\0')
sai.sin_addr.s_addr = htonl(INADDR_ANY);
else
sai.sin_addr.s_addr = inet_addr(pszAddress);
sai.sin_family = AF_INET;
sai.sin_port = htons(nPort);
#ifdef REUSEADDR
int nOptVal = 1;
setsockopt(m_sockListen, SOL_SOCKET, SO_REUSEADDR, (const char *)&nOptVal, sizeof(nOptVal)); // for multiple ip in UNIX
#endif
if (bind(m_sockListen, (struct sockaddr *)&sai, sizeof(struct sockaddr)) < 0)
{
HLOGF(HLOG_ERROR, "[ERROR], MpLstnThrd, Failed to bind, (%s:%d), socket = %d, error = %d\n", pszAddress, nPort, m_sockListen, SockErrorNo());
return false;
}
if (listen(m_sockListen, nBackLog) == SOCKET_ERROR)
{
HLOGF(HLOG_ERROR, "[ERROR], MpLstnThrd, Failed to listen, (%s:%d), socket = %d, error = %d\n", pszAddress, nPort, m_sockListen, SockErrorNo());
return false;
}
HLOGF(HLOG_INFO1, "INFO, MpLstnThrd, Succeed to listen TCP, (%s:%d), socket = %d\n", pszAddress, nPort, m_sockListen);
m_strListenIp = pszAddress;
m_nListenPort = nPort;
Begin();
return true;
}
void Terminate()
{
closesocket(m_sockListen);
for (HMpsDispatcherList::iterator iter = m_pListDispatcher->begin(); iter != m_pListDispatcher->end(); iter++)
{
(*iter)->SetTerminateFlag(); // auto delete
}
m_bTerminate = true;
}
protected:
void Main()
{
HLOGF(HLOG_INFO5, "INFO, MpLstnThrd, listen thread begins, (%s:%d)\n", m_strListenIp.psz(), m_nListenPort);
fd_set fdSetRead;//, fdSetError;
int nSelectResult;
SOCKET sockClient = INVALID_SOCKET;
socklen_t len = sizeof(struct sockaddr);
struct sockaddr_in sai;
FD_ZERO(&fdSetRead);
//FD_ZERO(&fdSetError);
while (m_bTerminate == false)
{
FD_SET(m_sockListen, &fdSetRead);
//FD_SET(m_sockListen, &fdSetError);
if ((nSelectResult = select(m_sockListen + 1, &fdSetRead, NULL, NULL, NULL)) == SOCKET_ERROR)
{
Sleep(10);
continue;
}
if ((sockClient = accept(m_sockListen, (struct sockaddr *)&sai, &len)) == INVALID_SOCKET)
{
HLOGF(HLOG_WARNING, "[WARN], MpLstnThrd, Failed to accept, (%s:%d), socket = %d, error = %d\n", m_strListenIp.psz(), m_nListenPort, m_sockListen, SockErrorNo());
Sleep(50);
}
else
{
HLOGF(HLOG_INFO4, "INFO, MpLstnThrd, Succeeded to accept, (%s:%d), (%s:%d), socket = %d\n", m_strListenIp.psz(), m_nListenPort, inet_ntoa(sai.sin_addr), ntohs(sai.sin_port), sockClient);
// Set non blocking mode
#ifdef _WIN32
ULONG nNonBlock = 1;
if (ioctlsocket(sockClient, FIONBIO, &nNonBlock) == SOCKET_ERROR)
HLOGF(HLOG_ERROR, "[ERROR], MpLstnThrd, Failed to ioctlsocket, socket = %d, error = %d\n", sockClient, SockErrorNo());
#else
int nFlags = fcntl(sockClient, F_GETFL, 0);
fcntl(sockClient, F_SETFL, nFlags | O_NONBLOCK);
#endif
HMultiplexTcpSocket * pAcceptedSocket = OnAccept(sockClient, &sai);
if (pAcceptedSocket)
{
pAcceptedSocket->m_nLastReceivedTime = HTIMER()->GetTickCount();
pAcceptedSocket->OnAccepted();
}
else
closesocket(sockClient);
}
}
}
private:
virtual HMultiplexTcpSocket * OnAccept(SOCKET sock, struct sockaddr_in * psaiRemote)
{
#if 1
HMpsDispatcher * pDispatcher = m_pListDispatcher->GetLazyDispatcher();
HMultiplexTcpSocket * pMpSocket = new _Ty(sock, psaiRemote, pDispatcher);
pDispatcher->Attach(pMpSocket);
return pMpSocket;
#else
HMpsDispatcherList::iterator iter;
int i = 0;
for (iter = m_pListDispatcher->begin(); iter != m_pListDispatcher->end(); iter++)
{
if ((*iter)->IsProcessing() == false)
i++;
}
//HLOGF(HLOG_WARNING, "[TRACE], MplxLstnThr_0x%x, lazy dispatcher count = %d\n", this, i);
for (;;)
{
for (iter = m_pListDispatcher->begin(); iter != m_pListDispatcher->end(); iter++)
{
if ((*iter)->IsProcessing() == false)
{
HMultiplexTcpSocket * pMpSocket = new _Ty(sock, psaiRemote, *iter);
(*iter)->Attach(pMpSocket);
return pMpSocket;
}
}
Sleep(10);
}
#endif
}
protected:
HMpsDispatcherList * m_pListDispatcher;
SOCKET m_sockListen;
HString m_strListenIp;
u_short m_nListenPort;
volatile bool m_bTerminate;
};
#endif // MULTIPLEX_SOCKET__H_