UdpNode.cpp 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757
  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. #if defined(__GNUC__) && __GNUC__ >= 11
  24. #pragma GCC diagnostic push
  25. #pragma GCC diagnostic ignored "-Warray-bounds"
  26. #pragma GCC diagnostic ignored "-Wstringop-overflow"
  27. #endif
  28. #include "UdpNode.h"
  29. #ifdef _UDP_SUPPORT
  30. BOOL CUdpNode::Start(LPCTSTR lpszBindAddress, USHORT usPort, EnCastMode enCastMode, LPCTSTR lpszCastAddress)
  31. {
  32. m_enCastMode = enCastMode;
  33. if(!CheckParams() || !CheckStarting())
  34. return FALSE;
  35. PrepareStart();
  36. HP_SOCKADDR bindAddr(AF_UNSPEC, TRUE);
  37. if(ParseBindAddr(lpszBindAddress, usPort, lpszCastAddress, bindAddr))
  38. if(CreateListenSocket(bindAddr))
  39. if(CreateWorkerThreads())
  40. if(StartAccept())
  41. {
  42. m_enState = SS_STARTED;
  43. return TRUE;
  44. }
  45. EXECUTE_RESTORE_ERROR(Stop());
  46. return FALSE;
  47. }
  48. BOOL CUdpNode::CheckParams()
  49. {
  50. if (((int)m_dwFreeBufferPoolSize >= 0) &&
  51. ((int)m_dwFreeBufferPoolHold >= 0) &&
  52. ((int)m_dwPostReceiveCount > 0) &&
  53. ((int)m_dwWorkerThreadCount > 0 && m_dwWorkerThreadCount <= MAX_WORKER_THREAD_COUNT) &&
  54. (m_enCastMode >= CM_UNICAST && m_enCastMode <= CM_BROADCAST) &&
  55. (m_iMCTtl >= 0 && m_iMCTtl <= 255) &&
  56. (m_bMCLoop == TRUE || m_bMCLoop == FALSE) &&
  57. ((int)m_dwMaxDatagramSize > 0 && m_dwMaxDatagramSize <= MAXIMUM_UDP_MAX_DATAGRAM_SIZE) )
  58. return TRUE;
  59. SetLastError(SE_INVALID_PARAM, __FUNCTION__, ERROR_INVALID_PARAMETER);
  60. return FALSE;
  61. }
  62. BOOL CUdpNode::CheckStarting()
  63. {
  64. CReentrantWriteLock locallock(m_lcState);
  65. if(m_enState == SS_STOPPED)
  66. m_enState = SS_STARTING;
  67. else
  68. {
  69. SetLastError(SE_ILLEGAL_STATE, __FUNCTION__, ERROR_INVALID_STATE);
  70. return FALSE;
  71. }
  72. return TRUE;
  73. }
  74. void CUdpNode::PrepareStart()
  75. {
  76. m_bfObjPool.SetItemCapacity(m_dwMaxDatagramSize);
  77. m_bfObjPool.SetPoolSize(m_dwFreeBufferPoolSize);
  78. m_bfObjPool.SetPoolHold(m_dwFreeBufferPoolHold);
  79. m_bfObjPool.Prepare();
  80. TNodeBufferObjList* pBufferObjList = (TNodeBufferObjList*)malloc(m_dwWorkerThreadCount * sizeof(TNodeBufferObjList));
  81. for(int i = 0; i < (int)m_dwWorkerThreadCount; i++)
  82. new (pBufferObjList + i) TNodeBufferObjList(m_bfObjPool);
  83. m_sndBuffs.reset(pBufferObjList);
  84. m_csSends = make_unique<CCriSec[]>(m_dwWorkerThreadCount);
  85. m_rcBuffers = make_unique<CBufferPtr[]>(m_dwWorkerThreadCount);
  86. for_each(m_rcBuffers.get(), m_rcBuffers.get() + m_dwWorkerThreadCount, [this](CBufferPtr& buff) {buff.Malloc(m_dwMaxDatagramSize);});
  87. m_soListens = make_unique<SOCKET[]>(m_dwWorkerThreadCount);
  88. for_each(m_soListens.get(), m_soListens.get() + m_dwWorkerThreadCount, [](SOCKET& sock) {sock = INVALID_FD;});
  89. }
  90. BOOL CUdpNode::ParseBindAddr(LPCTSTR lpszBindAddress, USHORT usPort, LPCTSTR lpszCastAddress, HP_SOCKADDR& bindAddr)
  91. {
  92. if(::IsStrEmpty(lpszCastAddress))
  93. {
  94. if(m_enCastMode == CM_BROADCAST)
  95. lpszCastAddress = DEFAULT_IPV4_BROAD_CAST_ADDRESS;
  96. else if(m_enCastMode == CM_MULTICAST)
  97. {
  98. SetLastError(SE_SOCKET_CREATE, __FUNCTION__, ERROR_ADDRNOTAVAIL);
  99. return FALSE;
  100. }
  101. }
  102. if(m_enCastMode != CM_UNICAST && !::sockaddr_A_2_IN(lpszCastAddress, usPort, m_castAddr))
  103. {
  104. SetLastError(SE_SOCKET_CREATE, __FUNCTION__, ::WSAGetLastError());
  105. return FALSE;
  106. }
  107. if(::IsStrEmpty(lpszBindAddress))
  108. {
  109. bindAddr.family = (m_enCastMode != CM_UNICAST) ? m_castAddr.family : AF_INET;
  110. bindAddr.SetPort(usPort);
  111. }
  112. else
  113. {
  114. if(!::sockaddr_A_2_IN(lpszBindAddress, usPort, bindAddr))
  115. {
  116. SetLastError(SE_SOCKET_CREATE, __FUNCTION__, ::WSAGetLastError());
  117. return FALSE;
  118. }
  119. }
  120. if(m_enCastMode == CM_BROADCAST && bindAddr.IsIPv6())
  121. {
  122. SetLastError(SE_SOCKET_CREATE, __FUNCTION__, ERROR_PFNOSUPPORT);
  123. return FALSE;
  124. }
  125. if(m_enCastMode != CM_UNICAST && m_castAddr.family != bindAddr.family)
  126. {
  127. SetLastError(SE_SOCKET_CREATE, __FUNCTION__, ERROR_AFNOSUPPORT);
  128. return FALSE;
  129. }
  130. return TRUE;
  131. }
  132. BOOL CUdpNode::CreateListenSocket(const HP_SOCKADDR& bindAddr)
  133. {
  134. for(DWORD i = 0; i < m_dwWorkerThreadCount; i++)
  135. {
  136. m_soListens[i] = socket(bindAddr.family, SOCK_DGRAM, IPPROTO_UDP);
  137. SOCKET soListen = m_soListens[i];
  138. if(IS_INVALID_FD(soListen))
  139. {
  140. SetLastError(SE_SOCKET_CREATE, __FUNCTION__, ::WSAGetLastError());
  141. return FALSE;
  142. }
  143. ::fcntl_SETFL(soListen, O_NOATIME | O_NONBLOCK | O_CLOEXEC);
  144. VERIFY(IS_NO_ERROR(::SSO_ReuseAddress(soListen, m_enReusePolicy)));
  145. if(IS_HAS_ERROR(::bind(soListen, bindAddr.Addr(), bindAddr.AddrSize())))
  146. {
  147. SetLastError(SE_SOCKET_BIND, __FUNCTION__, ::WSAGetLastError());
  148. return FALSE;
  149. }
  150. if(i == 0)
  151. {
  152. socklen_t dwAddrLen = (socklen_t)bindAddr.AddrSize();
  153. ENSURE(IS_NO_ERROR(::getsockname(soListen, m_localAddr.Addr(), &dwAddrLen)));
  154. }
  155. if(m_enCastMode == CM_MULTICAST)
  156. {
  157. if(!::SetMultiCastSocketOptions(soListen, bindAddr, m_castAddr, m_iMCTtl, m_bMCLoop))
  158. {
  159. SetLastError(SE_CONNECT_SERVER, __FUNCTION__, ::WSAGetLastError());
  160. return FALSE;
  161. }
  162. }
  163. else if(m_enCastMode == CM_BROADCAST)
  164. {
  165. ASSERT(m_castAddr.IsIPv4());
  166. BOOL bSet = TRUE;
  167. if(IS_HAS_ERROR(::SSO_SetSocketOption(soListen, SOL_SOCKET, SO_BROADCAST, &bSet, sizeof(BOOL))))
  168. {
  169. SetLastError(SE_CONNECT_SERVER, __FUNCTION__, ::WSAGetLastError());
  170. return FALSE;
  171. }
  172. }
  173. if(TRIGGER(FirePrepareListen(soListen)) == HR_ERROR)
  174. {
  175. SetLastError(SE_SOCKET_PREPARE, __FUNCTION__, ENSURE_ERROR_CANCELLED);
  176. return FALSE;
  177. }
  178. }
  179. return TRUE;
  180. }
  181. BOOL CUdpNode::CreateWorkerThreads()
  182. {
  183. return m_ioDispatcher.Start(this, m_dwPostReceiveCount, m_dwWorkerThreadCount);
  184. }
  185. BOOL CUdpNode::StartAccept()
  186. {
  187. for(int i = 0; i < (int)m_dwWorkerThreadCount; i++)
  188. {
  189. SOCKET& soListen = m_soListens[i];
  190. if(!m_ioDispatcher.AddFD(i, soListen, EPOLLIN | EPOLLOUT | EPOLLET, TO_PVOID(&soListen)))
  191. return FALSE;
  192. }
  193. return TRUE;
  194. }
  195. BOOL CUdpNode::Stop()
  196. {
  197. if(!CheckStoping())
  198. return FALSE;
  199. CloseListenSocket();
  200. WaitForWorkerThreadEnd();
  201. FireShutdown();
  202. ReleaseFreeBuffer();
  203. Reset();
  204. return TRUE;
  205. }
  206. BOOL CUdpNode::CheckStoping()
  207. {
  208. if(m_enState != SS_STOPPED)
  209. {
  210. CReentrantWriteLock locallock(m_lcState);
  211. if(HasStarted())
  212. {
  213. m_enState = SS_STOPPING;
  214. return TRUE;
  215. }
  216. }
  217. SetLastError(SE_ILLEGAL_STATE, __FUNCTION__, ERROR_INVALID_STATE);
  218. return FALSE;
  219. }
  220. void CUdpNode::CloseListenSocket()
  221. {
  222. if(m_soListens)
  223. {
  224. for_each(m_soListens.get(), m_soListens.get() + m_dwWorkerThreadCount, [](SOCKET& sock)
  225. {
  226. if(sock != INVALID_FD)
  227. {
  228. ::ManualCloseSocket(sock);
  229. sock = INVALID_FD;
  230. }
  231. });
  232. ::WaitFor(100);
  233. }
  234. }
  235. void CUdpNode::WaitForWorkerThreadEnd()
  236. {
  237. m_ioDispatcher.Stop();
  238. }
  239. void CUdpNode::ReleaseFreeBuffer()
  240. {
  241. for_each(m_sndBuffs.get(), m_sndBuffs.get() + m_dwWorkerThreadCount, [](TNodeBufferObjList& sndBuff)
  242. {
  243. sndBuff.Clear();
  244. sndBuff.~TNodeBufferObjList();
  245. });
  246. free(m_sndBuffs.release());
  247. m_csSends = nullptr;
  248. m_bfObjPool.Clear();
  249. }
  250. void CUdpNode::Reset()
  251. {
  252. m_castAddr.Reset();
  253. m_localAddr.Reset();
  254. m_soListens = nullptr;
  255. m_rcBuffers = nullptr;
  256. m_iSending = 0;
  257. m_enState = SS_STOPPED;
  258. m_evWait.SyncNotifyAll();
  259. }
  260. int CUdpNode::GenerateBufferIndex(const HP_SOCKADDR& addrRemote)
  261. {
  262. return (int)(addrRemote.Hash() % m_dwWorkerThreadCount);
  263. }
  264. BOOL CUdpNode::Send(LPCTSTR lpszRemoteAddress, USHORT usRemotePort, const BYTE* pBuffer, int iLength, int iOffset)
  265. {
  266. HP_SOCKADDR addrRemote;
  267. if(!::GetSockAddrByHostName(lpszRemoteAddress, usRemotePort, addrRemote))
  268. return FALSE;
  269. return DoSend(addrRemote, pBuffer, iLength, iOffset);
  270. }
  271. BOOL CUdpNode::SendPackets(LPCTSTR lpszRemoteAddress, USHORT usRemotePort, const WSABUF pBuffers[], int iCount)
  272. {
  273. HP_SOCKADDR addrRemote;
  274. if(!::GetSockAddrByHostName(lpszRemoteAddress, usRemotePort, addrRemote))
  275. return FALSE;
  276. return DoSendPackets(addrRemote, pBuffers, iCount);
  277. }
  278. BOOL CUdpNode::SendCast(const BYTE* pBuffer, int iLength, int iOffset)
  279. {
  280. if(m_enCastMode == CM_UNICAST)
  281. {
  282. ::SetLastError(ERROR_INVALID_OPERATION);
  283. return FALSE;
  284. }
  285. return DoSend(m_castAddr, pBuffer, iLength, iOffset);
  286. }
  287. BOOL CUdpNode::SendCastPackets(const WSABUF pBuffers[], int iCount)
  288. {
  289. if(m_enCastMode == CM_UNICAST)
  290. {
  291. ::SetLastError(ERROR_INVALID_OPERATION);
  292. return FALSE;
  293. }
  294. return DoSendPackets(m_castAddr, pBuffers, iCount);
  295. }
  296. BOOL CUdpNode::DoSend(const HP_SOCKADDR& addrRemote, const BYTE* pBuffer, int iLength, int iOffset)
  297. {
  298. ASSERT(pBuffer && iLength >= 0 && iLength <= (int)m_dwMaxDatagramSize);
  299. int result = NO_ERROR;
  300. if(IsValid())
  301. {
  302. if(addrRemote.family == m_localAddr.family)
  303. {
  304. if(pBuffer && iLength >= 0 && iLength <= (int)m_dwMaxDatagramSize)
  305. {
  306. if(iOffset != 0) pBuffer += iOffset;
  307. TNodeBufferObjPtr bufPtr(m_bfObjPool, m_bfObjPool.PickFreeItem());
  308. bufPtr->Cat(pBuffer, iLength);
  309. result = SendInternal(addrRemote, bufPtr);
  310. }
  311. else
  312. result = ERROR_INVALID_PARAMETER;
  313. }
  314. else
  315. result = ERROR_AFNOSUPPORT;
  316. }
  317. else
  318. result = ERROR_INVALID_STATE;
  319. if(result != NO_ERROR)
  320. ::SetLastError(result);
  321. return (result == NO_ERROR);
  322. }
  323. BOOL CUdpNode::DoSendPackets(const HP_SOCKADDR& addrRemote, const WSABUF pBuffers[], int iCount)
  324. {
  325. ASSERT(pBuffers && iCount > 0);
  326. if(!pBuffers || iCount <= 0)
  327. return ERROR_INVALID_PARAMETER;
  328. if(!IsValid())
  329. {
  330. ::SetLastError(ERROR_INVALID_STATE);
  331. return FALSE;
  332. }
  333. if(addrRemote.family != m_localAddr.family)
  334. {
  335. ::SetLastError(ERROR_AFNOSUPPORT);
  336. return FALSE;
  337. }
  338. int result = NO_ERROR;
  339. int iLength = 0;
  340. int iMaxLen = (int)m_dwMaxDatagramSize;
  341. TNodeBufferObjPtr bufPtr(m_bfObjPool, m_bfObjPool.PickFreeItem());
  342. for(int i = 0; i < iCount; i++)
  343. {
  344. int iBufLen = pBuffers[i].len;
  345. if(iBufLen > 0)
  346. {
  347. BYTE* pBuffer = (BYTE*)pBuffers[i].buf;
  348. ASSERT(pBuffer);
  349. iLength += iBufLen;
  350. if(iLength <= iMaxLen)
  351. bufPtr->Cat(pBuffer, iBufLen);
  352. else
  353. break;
  354. }
  355. }
  356. if(iLength > 0 && iLength <= iMaxLen)
  357. result = SendInternal(addrRemote, bufPtr);
  358. else
  359. result = ERROR_INCORRECT_SIZE;
  360. if(result != NO_ERROR)
  361. ::SetLastError(result);
  362. return (result == NO_ERROR);
  363. }
  364. int CUdpNode::SendInternal(const HP_SOCKADDR& addrRemote, TNodeBufferObjPtr& bufPtr)
  365. {
  366. BOOL bPending;
  367. int iBufferSize = bufPtr->Size();
  368. int idx = GenerateBufferIndex(addrRemote);
  369. addrRemote.Copy(bufPtr->remoteAddr);
  370. {
  371. CReentrantReadLock locallock(m_lcState);
  372. if(!IsValid())
  373. return ERROR_INVALID_STATE;
  374. TNodeBufferObjList& sndBuff = m_sndBuffs[idx];
  375. CCriSecLock locallock2(m_csSends[idx]);
  376. bPending = IsPending(idx);
  377. sndBuff.PushBack(bufPtr.Detach());
  378. if(iBufferSize == 0) sndBuff.IncreaseLength(1);
  379. ASSERT(sndBuff.Length() > 0);
  380. }
  381. if(!bPending && IsPending(idx))
  382. VERIFY(m_ioDispatcher.SendCommandByIndex(idx, DISP_CMD_SEND));
  383. return NO_ERROR;
  384. }
  385. BOOL CUdpNode::OnBeforeProcessIo(const TDispContext* pContext, PVOID pv, UINT events)
  386. {
  387. ASSERT(pv == &m_soListens[pContext->GetIndex()]);
  388. return TRUE;
  389. }
  390. VOID CUdpNode::OnAfterProcessIo(const TDispContext* pContext, PVOID pv, UINT events, BOOL rs)
  391. {
  392. }
  393. VOID CUdpNode::OnCommand(const TDispContext* pContext, TDispCommand* pCmd)
  394. {
  395. int idx = pContext->GetIndex();
  396. int flag = (int)(pCmd->wParam);
  397. switch(pCmd->type)
  398. {
  399. case DISP_CMD_SEND:
  400. HandleCmdSend(idx, flag);
  401. break;
  402. }
  403. }
  404. BOOL CUdpNode::OnReadyRead(const TDispContext* pContext, PVOID pv, UINT events)
  405. {
  406. return HandleReceive(pContext, RETRIVE_EVENT_FLAG_H(events));
  407. }
  408. BOOL CUdpNode::OnReadyWrite(const TDispContext* pContext, PVOID pv, UINT events)
  409. {
  410. return HandleSend(pContext, RETRIVE_EVENT_FLAG_H(events), RETRIVE_EVENT_FLAG_R(events));
  411. }
  412. BOOL CUdpNode::OnHungUp(const TDispContext* pContext, PVOID pv, UINT events)
  413. {
  414. return HandleClose(pContext->GetIndex(), nullptr, SO_CLOSE, 0);
  415. }
  416. BOOL CUdpNode::OnError(const TDispContext* pContext, PVOID pv, UINT events)
  417. {
  418. return HandleClose(pContext->GetIndex(), nullptr, SO_CLOSE, -1);
  419. }
  420. VOID CUdpNode::OnDispatchThreadStart(THR_ID tid)
  421. {
  422. OnWorkerThreadStart(tid);
  423. }
  424. VOID CUdpNode::OnDispatchThreadEnd(THR_ID tid)
  425. {
  426. OnWorkerThreadEnd(tid);
  427. }
  428. BOOL CUdpNode::HandleClose(int idx, TNodeBufferObj* pBufferObj, EnSocketOperation enOperation, int iErrorCode)
  429. {
  430. if(!HasStarted())
  431. return FALSE;
  432. if(iErrorCode == -1)
  433. iErrorCode = ::SSO_GetError(m_soListens[idx]);
  434. if(pBufferObj != nullptr)
  435. TRIGGER(FireError(&pBufferObj->remoteAddr, pBufferObj->Ptr(), pBufferObj->Size(), enOperation, iErrorCode));
  436. else
  437. TRIGGER(FireError(nullptr, nullptr, 0, enOperation, iErrorCode));
  438. return TRUE;
  439. }
  440. BOOL CUdpNode::HandleReceive(const TDispContext* pContext, int flag)
  441. {
  442. int idx = pContext->GetIndex();
  443. CBufferPtr& buffer = m_rcBuffers[idx];
  444. int iBufferLen = (int)buffer.Size();
  445. while(TRUE)
  446. {
  447. HP_SOCKADDR addr;
  448. socklen_t dwAddrLen = (socklen_t)addr.AddrSize();
  449. int rc = (int)recvfrom(m_soListens[idx], buffer.Ptr(), iBufferLen, MSG_TRUNC, addr.Addr(), &dwAddrLen);
  450. if(rc >= 0)
  451. {
  452. if(rc > iBufferLen)
  453. {
  454. TRIGGER(FireError(&addr, buffer.Ptr(), iBufferLen, SO_RECEIVE, ERROR_BAD_LENGTH));
  455. continue;
  456. }
  457. TRIGGER(FireReceive(&addr, buffer.Ptr(), rc));
  458. }
  459. else if(rc == SOCKET_ERROR)
  460. {
  461. int code = ::WSAGetLastError();
  462. if(code == ERROR_WOULDBLOCK)
  463. break;
  464. else if(!HandleClose(idx, nullptr, SO_RECEIVE, code))
  465. return FALSE;
  466. }
  467. else
  468. {
  469. ASSERT(FALSE);
  470. }
  471. }
  472. return TRUE;
  473. }
  474. BOOL CUdpNode::HandleSend(const TDispContext* pContext, int flag, int rd)
  475. {
  476. HandleCmdSend(pContext->GetIndex(), flag);
  477. return TRUE;
  478. }
  479. VOID CUdpNode::HandleCmdSend(int idx, int flag)
  480. {
  481. BOOL bBlocked = FALSE;
  482. TNodeBufferObjList& sndBuff = m_sndBuffs[idx];
  483. TNodeBufferObjPtr bufPtr(m_bfObjPool);
  484. while(IsPending(idx))
  485. {
  486. {
  487. CCriSecLock locallock(m_csSends[idx]);
  488. bufPtr = sndBuff.PopFront();
  489. }
  490. if(!bufPtr.IsValid())
  491. break;
  492. if(!SendItem(idx, sndBuff, bufPtr, bBlocked))
  493. return;
  494. if(bBlocked)
  495. {
  496. {
  497. CCriSecLock locallock(m_csSends[idx]);
  498. sndBuff.PushFront(bufPtr.Detach());
  499. }
  500. break;
  501. }
  502. }
  503. if(!bBlocked && IsPending(idx))
  504. VERIFY(m_ioDispatcher.SendCommandByIndex(idx, DISP_CMD_SEND));
  505. }
  506. BOOL CUdpNode::SendItem(int idx, TNodeBufferObjList& sndBuff, TNodeBufferObj* pBufferObj, BOOL& bBlocked)
  507. {
  508. int rc = (int)sendto(m_soListens[idx], pBufferObj->Ptr(), pBufferObj->Size(), 0, pBufferObj->remoteAddr.Addr(), pBufferObj->remoteAddr.AddrSize());
  509. if(rc >= 0)
  510. {
  511. ASSERT(rc == pBufferObj->Size());
  512. if(rc == 0)
  513. {
  514. CCriSecLock locallock(m_csSends[idx]);
  515. sndBuff.ReduceLength(1);
  516. }
  517. TRIGGER(FireSend(pBufferObj));
  518. }
  519. else if(rc == SOCKET_ERROR)
  520. {
  521. int code = ::WSAGetLastError();
  522. if(code == ERROR_WOULDBLOCK)
  523. bBlocked = TRUE;
  524. else if(!HandleClose(idx, pBufferObj, SO_SEND, code))
  525. return FALSE;
  526. }
  527. else
  528. {
  529. ASSERT(FALSE);
  530. }
  531. return TRUE;
  532. }
  533. BOOL CUdpNode::GetLocalAddress(TCHAR lpszAddress[], int& iAddressLen, USHORT& usPort)
  534. {
  535. ADDRESS_FAMILY usFamily;
  536. return ::sockaddr_IN_2_A(m_localAddr, usFamily, lpszAddress, iAddressLen, usPort);
  537. }
  538. BOOL CUdpNode::GetCastAddress(TCHAR lpszAddress[], int& iAddressLen, USHORT& usPort)
  539. {
  540. ADDRESS_FAMILY usFamily;
  541. return ::sockaddr_IN_2_A(m_castAddr, usFamily, lpszAddress, iAddressLen, usPort);
  542. }
  543. void CUdpNode::SetLastError(EnSocketError code, LPCSTR func, int ec)
  544. {
  545. m_enLastError = code;
  546. ::SetLastError(ec);
  547. }
  548. BOOL CUdpNode::GetPendingDataLength(int& iPending)
  549. {
  550. iPending = 0;
  551. {
  552. CReentrantReadLock locallock(m_lcState);
  553. if(!IsValid())
  554. return FALSE;
  555. for_each(m_sndBuffs.get(), m_sndBuffs.get() + m_dwWorkerThreadCount, [&iPending](TNodeBufferObjList& sndBuff) { iPending += sndBuff.Length(); });
  556. }
  557. return TRUE;
  558. }
  559. EnHandleResult CUdpNode::FireSend(TNodeBufferObj* pBufferObj)
  560. {
  561. TCHAR szAddress[60];
  562. int iAddressLen = ARRAY_SIZE(szAddress);
  563. ADDRESS_FAMILY usFamily;
  564. USHORT usPort;
  565. ::sockaddr_IN_2_A(pBufferObj->remoteAddr, usFamily, szAddress, iAddressLen, usPort);
  566. return m_pListener->OnSend(this, szAddress, usPort, pBufferObj->Ptr(), pBufferObj->Size());
  567. }
  568. EnHandleResult CUdpNode::FireReceive(const HP_SOCKADDR* pRemoteAddr, const BYTE* pData, int iLength)
  569. {
  570. TCHAR szAddress[60];
  571. int iAddressLen = ARRAY_SIZE(szAddress);
  572. ADDRESS_FAMILY usFamily;
  573. USHORT usPort;
  574. ::sockaddr_IN_2_A(*pRemoteAddr, usFamily, szAddress, iAddressLen, usPort);
  575. return m_pListener->OnReceive(this, szAddress, usPort, pData, iLength);
  576. }
  577. EnHandleResult CUdpNode::FireError(const HP_SOCKADDR* pRemoteAddr, const BYTE* pData, int iLength, EnSocketOperation enOperation, int iErrorCode)
  578. {
  579. TCHAR szAddress[60];
  580. int iAddressLen = ARRAY_SIZE(szAddress);
  581. ADDRESS_FAMILY usFamily;
  582. USHORT usPort;
  583. if(pRemoteAddr == nullptr)
  584. {
  585. ::sockaddr_IN_2_A(m_localAddr, usFamily, szAddress, iAddressLen, usPort);
  586. return m_pListener->OnError(this, enOperation, iErrorCode, szAddress, usPort, nullptr, 0);
  587. }
  588. ::sockaddr_IN_2_A(*pRemoteAddr, usFamily, szAddress, iAddressLen, usPort);
  589. return m_pListener->OnError(this, enOperation, iErrorCode, szAddress, usPort, pData, iLength);
  590. }
  591. #endif
  592. #if defined(__GNUC__) && __GNUC__ >= 11
  593. #pragma GCC diagnostic pop
  594. #endif