UdpNode.h 8.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198
  1. /*
  2. * Copyright: JessMA Open Source (ldcsaa@gmail.com)
  3. *
  4. * Author : Bruce Liang
  5. * Website : https://github.com/ldcsaa
  6. * Project : https://github.com/ldcsaa/HP-Socket
  7. * Blog : http://www.cnblogs.com/ldcsaa
  8. * Wiki : http://www.oschina.net/p/hp-socket
  9. * QQ Group : 44636872, 75375912
  10. *
  11. * Licensed under the Apache License, Version 2.0 (the "License");
  12. * you may not use this file except in compliance with the License.
  13. * You may obtain a copy of the License at
  14. *
  15. * http://www.apache.org/licenses/LICENSE-2.0
  16. *
  17. * Unless required by applicable law or agreed to in writing, software
  18. * distributed under the License is distributed on an "AS IS" BASIS,
  19. * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
  20. * See the License for the specific language governing permissions and
  21. * limitations under the License.
  22. */
  23. #pragma once
  24. #include "SocketHelper.h"
  25. #include "./common/GeneralHelper.h"
  26. #include "./common/IODispatcher.h"
  27. #ifdef _UDP_SUPPORT
  28. class CUdpNode : public IUdpNode, private CIOHandler
  29. {
  30. public:
  31. virtual BOOL Start (LPCTSTR lpszBindAddress = nullptr, USHORT usPort = 0, EnCastMode enCastMode = CM_UNICAST, LPCTSTR lpszCastAddress = nullptr);
  32. virtual BOOL Stop ();
  33. virtual BOOL Send (LPCTSTR lpszRemoteAddress, USHORT usRemotePort, const BYTE* pBuffer, int iLength, int iOffset = 0);
  34. virtual BOOL SendPackets (LPCTSTR lpszRemoteAddress, USHORT usRemotePort, const WSABUF pBuffers[], int iCount);
  35. virtual BOOL SendCast (const BYTE* pBuffer, int iLength, int iOffset = 0);
  36. virtual BOOL SendCastPackets(const WSABUF pBuffers[], int iCount);
  37. virtual BOOL Wait (DWORD dwMilliseconds = INFINITE) {return m_evWait.WaitFor(dwMilliseconds, WAIT_FOR_STOP_PREDICATE);}
  38. virtual BOOL HasStarted () {return m_enState == SS_STARTED || m_enState == SS_STARTING;}
  39. virtual EnServiceState GetState () {return m_enState;}
  40. virtual BOOL GetLocalAddress (TCHAR lpszAddress[], int& iAddressLen, USHORT& usPort);
  41. virtual BOOL GetCastAddress (TCHAR lpszAddress[], int& iAddressLen, USHORT& usPort);
  42. virtual EnSocketError GetLastError () {return m_enLastError;}
  43. virtual LPCTSTR GetLastErrorDesc () {return ::GetSocketErrorDesc(m_enLastError);}
  44. virtual BOOL GetPendingDataLength (int& iPending);
  45. private:
  46. virtual BOOL OnBeforeProcessIo(const TDispContext* pContext, PVOID pv, UINT events) override;
  47. virtual VOID OnAfterProcessIo(const TDispContext* pContext, PVOID pv, UINT events, BOOL rs) override;
  48. virtual VOID OnCommand(const TDispContext* pContext, TDispCommand* pCmd) override;
  49. virtual BOOL OnReadyRead(const TDispContext* pContext, PVOID pv, UINT events) override;
  50. virtual BOOL OnReadyWrite(const TDispContext* pContext, PVOID pv, UINT events) override;
  51. virtual BOOL OnHungUp(const TDispContext* pContext, PVOID pv, UINT events) override;
  52. virtual BOOL OnError(const TDispContext* pContext, PVOID pv, UINT events) override;
  53. virtual VOID OnDispatchThreadStart(THR_ID tid) override;
  54. virtual VOID OnDispatchThreadEnd(THR_ID tid) override;
  55. public:
  56. virtual void SetReuseAddressPolicy (EnReuseAddressPolicy enReusePolicy){ENSURE_HAS_STOPPED(); ASSERT(m_enReusePolicy == enReusePolicy);}
  57. virtual void SetWorkerThreadCount (DWORD dwWorkerThreadCount) {ENSURE_HAS_STOPPED(); m_dwWorkerThreadCount = dwWorkerThreadCount;}
  58. virtual void SetFreeBufferPoolSize (DWORD dwFreeBufferPoolSize) {ENSURE_HAS_STOPPED(); m_dwFreeBufferPoolSize = dwFreeBufferPoolSize;}
  59. virtual void SetFreeBufferPoolHold (DWORD dwFreeBufferPoolHold) {ENSURE_HAS_STOPPED(); m_dwFreeBufferPoolHold = dwFreeBufferPoolHold;}
  60. virtual void SetPostReceiveCount (DWORD dwPostReceiveCount) {ENSURE_HAS_STOPPED(); m_dwPostReceiveCount = dwPostReceiveCount;}
  61. virtual void SetMaxDatagramSize (DWORD dwMaxDatagramSize) {ENSURE_HAS_STOPPED(); m_dwMaxDatagramSize = dwMaxDatagramSize;}
  62. virtual void SetMultiCastTtl (int iMCTtl) {ENSURE_HAS_STOPPED(); m_iMCTtl = iMCTtl;}
  63. virtual void SetMultiCastLoop (BOOL bMCLoop) {ENSURE_HAS_STOPPED(); m_bMCLoop = bMCLoop;}
  64. virtual void SetExtra (PVOID pExtra) {m_pExtra = pExtra;}
  65. virtual EnReuseAddressPolicy GetReuseAddressPolicy () {return m_enReusePolicy;}
  66. virtual DWORD GetWorkerThreadCount () {return m_dwWorkerThreadCount;}
  67. virtual DWORD GetFreeBufferPoolSize () {return m_dwFreeBufferPoolSize;}
  68. virtual DWORD GetFreeBufferPoolHold () {return m_dwFreeBufferPoolHold;}
  69. virtual DWORD GetPostReceiveCount () {return m_dwPostReceiveCount;}
  70. virtual DWORD GetMaxDatagramSize () {return m_dwMaxDatagramSize;}
  71. virtual EnCastMode GetCastMode () {return m_enCastMode;}
  72. virtual int GetMultiCastTtl () {return m_iMCTtl;}
  73. virtual BOOL IsMultiCastLoop () {return m_bMCLoop;}
  74. virtual PVOID GetExtra () {return m_pExtra;}
  75. protected:
  76. EnHandleResult FirePrepareListen(SOCKET soListen)
  77. {return m_pListener->OnPrepareListen(this, soListen);}
  78. EnHandleResult FireShutdown()
  79. {return m_pListener->OnShutdown(this);}
  80. EnHandleResult FireSend(TNodeBufferObj* pBufferObj);
  81. EnHandleResult FireReceive(const HP_SOCKADDR* pRemoteAddr, const BYTE* pData, int iLength);
  82. EnHandleResult FireError(const HP_SOCKADDR* pRemoteAddr, const BYTE* pData, int iLength, EnSocketOperation enOperation, int iErrorCode);
  83. void SetLastError(EnSocketError code, LPCSTR func, int ec);
  84. virtual BOOL CheckParams();
  85. virtual void PrepareStart();
  86. virtual void Reset();
  87. virtual void OnWorkerThreadStart(THR_ID dwThreadID) {}
  88. virtual void OnWorkerThreadEnd(THR_ID dwThreadID) {}
  89. BOOL DoSend(const HP_SOCKADDR& addrRemote, const BYTE* pBuffer, int iLength, int iOffset = 0);
  90. BOOL DoSendPackets(const HP_SOCKADDR& addrRemote, const WSABUF pBuffers[], int iCount);
  91. int SendInternal(const HP_SOCKADDR& addrRemote, TNodeBufferObjPtr& bufPtr);
  92. private:
  93. BOOL CheckStarting();
  94. BOOL CheckStoping();
  95. BOOL ParseBindAddr(LPCTSTR lpszBindAddress, USHORT usPort, LPCTSTR lpszCastAddress, HP_SOCKADDR& bindAddr);
  96. BOOL CreateListenSocket(const HP_SOCKADDR& bindAddr);
  97. BOOL CreateWorkerThreads();
  98. BOOL StartAccept();
  99. void CloseListenSocket();
  100. void WaitForWorkerThreadEnd();
  101. void ReleaseFreeBuffer();
  102. int GenerateBufferIndex(const HP_SOCKADDR& addrRemote);
  103. private:
  104. BOOL HandleReceive(const TDispContext* pContext, int flag = 0);
  105. BOOL HandleSend(const TDispContext* pContext, int flag = 0, int rd = 0);
  106. BOOL HandleClose(int idx, TNodeBufferObj* pBufferObj, EnSocketOperation enOperation, int iErrorCode);
  107. VOID HandleCmdSend(int idx, int flag);
  108. BOOL SendItem(int idx, TNodeBufferObjList& sndBuff, TNodeBufferObj* pBufferObj, BOOL& bBlocked);
  109. private:
  110. BOOL IsValid () {return m_enState == SS_STARTED;}
  111. BOOL IsPending (int idx) {return m_sndBuffs[idx].Length() > 0;}
  112. public:
  113. CUdpNode(IUdpNodeListener* pListener)
  114. : m_pListener (pListener)
  115. , m_iSending (0)
  116. , m_enLastError (SE_OK)
  117. , m_enState (SS_STOPPED)
  118. , m_enReusePolicy (RAP_ADDR_AND_PORT)
  119. , m_dwWorkerThreadCount (DEFAULT_WORKER_THREAD_COUNT)
  120. , m_dwFreeBufferPoolSize (DEFAULT_FREE_BUFFEROBJ_POOL)
  121. , m_dwFreeBufferPoolHold (DEFAULT_FREE_BUFFEROBJ_HOLD)
  122. , m_dwPostReceiveCount (DEFAULT_UDP_POST_RECEIVE_COUNT)
  123. , m_dwMaxDatagramSize (DEFAULT_UDP_MAX_DATAGRAM_SIZE)
  124. , m_pExtra (nullptr)
  125. , m_iMCTtl (1)
  126. , m_bMCLoop (FALSE)
  127. , m_enCastMode (CM_UNICAST)
  128. , m_castAddr (AF_UNSPEC, TRUE)
  129. , m_localAddr (AF_UNSPEC, TRUE)
  130. {
  131. ASSERT(m_pListener);
  132. }
  133. virtual ~CUdpNode()
  134. {
  135. ENSURE_STOP();
  136. }
  137. private:
  138. EnReuseAddressPolicy m_enReusePolicy;
  139. DWORD m_dwWorkerThreadCount;
  140. DWORD m_dwFreeBufferPoolSize;
  141. DWORD m_dwFreeBufferPoolHold;
  142. DWORD m_dwPostReceiveCount;
  143. DWORD m_dwMaxDatagramSize;
  144. PVOID m_pExtra;
  145. int m_iMCTtl;
  146. BOOL m_bMCLoop;
  147. EnCastMode m_enCastMode;
  148. private:
  149. CSEM m_evWait;
  150. HP_SOCKADDR m_castAddr;
  151. HP_SOCKADDR m_localAddr;
  152. CNodeBufferObjPool m_bfObjPool;
  153. CNodeCriSecs m_csSends;
  154. TNodeBufferObjLists m_sndBuffs;
  155. IUdpNodeListener* m_pListener;
  156. ListenSocketsPtr m_soListens;
  157. EnServiceState m_enState;
  158. EnSocketError m_enLastError;
  159. CReceiveBuffersPtr m_rcBuffers;
  160. CRWLock m_lcState;
  161. volatile long m_iSending;
  162. CIODispatcher m_ioDispatcher;
  163. };
  164. #endif