TcpServer.cpp 27 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187
  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. #include "TcpServer.h"
  24. #include "./common/FileHelper.h"
  25. BOOL CTcpServer::Start(LPCTSTR lpszBindAddress, USHORT usPort)
  26. {
  27. if(!CheckParams() || !CheckStarting())
  28. return FALSE;
  29. PrepareStart();
  30. if(CreateListenSocket(lpszBindAddress, usPort))
  31. if(CreateWorkerThreads())
  32. if(StartAccept())
  33. {
  34. m_enState = SS_STARTED;
  35. return TRUE;
  36. }
  37. EXECUTE_RESTORE_ERROR(Stop());
  38. return FALSE;
  39. }
  40. void CTcpServer::SetLastError(EnSocketError code, LPCSTR func, int ec)
  41. {
  42. m_enLastError = code;
  43. ::SetLastError(ec);
  44. }
  45. BOOL CTcpServer::CheckParams()
  46. {
  47. if ((m_enSendPolicy >= SP_PACK && m_enSendPolicy <= SP_DIRECT) &&
  48. (m_enOnSendSyncPolicy >= OSSP_NONE && m_enOnSendSyncPolicy <= OSSP_RECEIVE) &&
  49. ((int)m_dwMaxConnectionCount > 0 && m_dwMaxConnectionCount <= MAX_CONNECTION_COUNT) &&
  50. ((int)m_dwWorkerThreadCount > 0 && m_dwWorkerThreadCount <= MAX_WORKER_THREAD_COUNT) &&
  51. ((int)m_dwAcceptSocketCount > 0) &&
  52. ((int)m_dwSocketBufferSize >= MIN_SOCKET_BUFFER_SIZE) &&
  53. ((int)m_dwSocketListenQueue > 0) &&
  54. ((int)m_dwFreeSocketObjLockTime >= 1000) &&
  55. ((int)m_dwFreeSocketObjPool >= 0) &&
  56. ((int)m_dwFreeBufferObjPool >= 0) &&
  57. ((int)m_dwFreeSocketObjHold >= 0) &&
  58. ((int)m_dwFreeBufferObjHold >= 0) &&
  59. ((int)m_dwKeepAliveTime >= 1000 || m_dwKeepAliveTime == 0) &&
  60. ((int)m_dwKeepAliveInterval >= 1000 || m_dwKeepAliveInterval == 0) )
  61. return TRUE;
  62. SetLastError(SE_INVALID_PARAM, __FUNCTION__, ERROR_INVALID_PARAMETER);
  63. return FALSE;
  64. }
  65. void CTcpServer::PrepareStart()
  66. {
  67. m_bfActiveSockets.Reset(m_dwMaxConnectionCount);
  68. m_lsFreeSocket.Reset(m_dwFreeSocketObjPool);
  69. m_bfObjPool.SetItemCapacity(m_dwSocketBufferSize);
  70. m_bfObjPool.SetPoolSize(m_dwFreeBufferObjPool);
  71. m_bfObjPool.SetPoolHold(m_dwFreeBufferObjHold);
  72. m_bfObjPool.Prepare();
  73. m_rcBuffers = make_unique<CBufferPtr[]>(m_dwWorkerThreadCount);
  74. for_each(m_rcBuffers.get(), m_rcBuffers.get() + m_dwWorkerThreadCount, [this](CBufferPtr& buff) {buff.Malloc(m_dwSocketBufferSize);});
  75. m_soListens = make_unique<SOCKET[]>(m_dwWorkerThreadCount);
  76. for_each(m_soListens.get(), m_soListens.get() + m_dwWorkerThreadCount, [](SOCKET& sock) {sock = INVALID_FD;});
  77. }
  78. BOOL CTcpServer::CheckStarting()
  79. {
  80. CSpinLock locallock(m_csState);
  81. if(m_enState == SS_STOPPED)
  82. m_enState = SS_STARTING;
  83. else
  84. {
  85. SetLastError(SE_ILLEGAL_STATE, __FUNCTION__, ERROR_INVALID_STATE);
  86. return FALSE;
  87. }
  88. return TRUE;
  89. }
  90. BOOL CTcpServer::CheckStoping()
  91. {
  92. if(m_enState != SS_STOPPED)
  93. {
  94. CSpinLock locallock(m_csState);
  95. if(HasStarted())
  96. {
  97. m_enState = SS_STOPPING;
  98. return TRUE;
  99. }
  100. }
  101. SetLastError(SE_ILLEGAL_STATE, __FUNCTION__, ERROR_INVALID_STATE);
  102. return FALSE;
  103. }
  104. BOOL CTcpServer::CreateListenSocket(LPCTSTR lpszBindAddress, USHORT usPort)
  105. {
  106. if(::IsStrEmpty(lpszBindAddress))
  107. lpszBindAddress = DEFAULT_IPV4_BIND_ADDRESS;
  108. HP_SOCKADDR addr;
  109. if(!::sockaddr_A_2_IN(lpszBindAddress, usPort, addr))
  110. {
  111. SetLastError(SE_SOCKET_CREATE, __FUNCTION__, ::WSAGetLastError());
  112. return FALSE;
  113. }
  114. for(DWORD i = 0; i < m_dwWorkerThreadCount; i++)
  115. {
  116. m_soListens[i] = socket(addr.family, SOCK_STREAM, IPPROTO_TCP);
  117. SOCKET soListen = m_soListens[i];
  118. if(IS_INVALID_FD(soListen))
  119. {
  120. SetLastError(SE_SOCKET_CREATE, __FUNCTION__, ::WSAGetLastError());
  121. return FALSE;
  122. }
  123. ::fcntl_SETFL(soListen, O_NOATIME | O_NONBLOCK | O_CLOEXEC);
  124. BOOL bOnOff = (m_dwKeepAliveTime > 0 && m_dwKeepAliveInterval > 0);
  125. VERIFY(IS_NO_ERROR(::SSO_KeepAliveVals(soListen, bOnOff, m_dwKeepAliveTime, m_dwKeepAliveInterval)));
  126. VERIFY(IS_NO_ERROR(::SSO_ReuseAddress(soListen, m_enReusePolicy)));
  127. VERIFY(IS_NO_ERROR(::SSO_NoDelay(soListen, m_bNoDelay)));
  128. if(IS_HAS_ERROR(::bind(soListen, addr.Addr(), addr.AddrSize())))
  129. {
  130. SetLastError(SE_SOCKET_BIND, __FUNCTION__, ::WSAGetLastError());
  131. return FALSE;
  132. }
  133. if(TRIGGER(FirePrepareListen(soListen)) == HR_ERROR)
  134. {
  135. SetLastError(SE_SOCKET_PREPARE, __FUNCTION__, ENSURE_ERROR_CANCELLED);
  136. return FALSE;
  137. }
  138. if(IS_HAS_ERROR(::listen(soListen, m_dwSocketListenQueue)))
  139. {
  140. SetLastError(SE_SOCKET_LISTEN, __FUNCTION__, ::WSAGetLastError());
  141. return FALSE;
  142. }
  143. }
  144. return TRUE;
  145. }
  146. BOOL CTcpServer::CreateWorkerThreads()
  147. {
  148. DWORD dwWorkerThreadCount = m_dwWorkerThreadCount
  149. #ifdef USE_EXTERNAL_GC
  150. + 1
  151. #endif
  152. ;
  153. if(!m_ioDispatcher.Start(this, m_dwAcceptSocketCount, dwWorkerThreadCount))
  154. {
  155. SetLastError(SE_WORKER_THREAD_CREATE, __FUNCTION__, ::WSAGetLastError());
  156. return FALSE;
  157. }
  158. return TRUE;
  159. }
  160. BOOL CTcpServer::StartAccept()
  161. {
  162. for(int i = 0; i < (int)m_dwWorkerThreadCount; i++)
  163. {
  164. SOCKET& soListen = m_soListens[i];
  165. if(!m_ioDispatcher.AddFD(i, soListen, EPOLLIN | EPOLLET, TO_PVOID(&soListen)))
  166. {
  167. SetLastError(SE_SOCKE_ATTACH_TO_CP, __FUNCTION__, ::WSAGetLastError());
  168. return FALSE;
  169. }
  170. }
  171. #ifdef USE_EXTERNAL_GC
  172. m_fdGCTimer = m_ioDispatcher.AddTimer(m_dwWorkerThreadCount, GC_CHECK_INTERVAL, this);
  173. if(IS_INVALID_FD(m_fdGCTimer))
  174. {
  175. SetLastError(SE_GC_START, __FUNCTION__, ::WSAGetLastError());
  176. return FALSE;
  177. }
  178. #endif
  179. return TRUE;
  180. }
  181. BOOL CTcpServer::Stop()
  182. {
  183. if(!CheckStoping())
  184. return FALSE;
  185. CloseListenSocket();
  186. DisconnectClientSocket();
  187. WaitForClientSocketClose();
  188. WaitForWorkerThreadEnd();
  189. ReleaseClientSocket();
  190. FireShutdown();
  191. ReleaseFreeSocket();
  192. Reset();
  193. return TRUE;
  194. }
  195. void CTcpServer::CloseListenSocket()
  196. {
  197. if(m_soListens)
  198. {
  199. for_each(m_soListens.get(), m_soListens.get() + m_dwWorkerThreadCount, [](SOCKET& sock)
  200. {
  201. if(sock != INVALID_FD)
  202. {
  203. ::ManualCloseSocket(sock);
  204. sock = INVALID_FD;
  205. }
  206. });
  207. ::WaitFor(100);
  208. }
  209. }
  210. void CTcpServer::DisconnectClientSocket()
  211. {
  212. if(m_bfActiveSockets.Elements() == 0)
  213. return;
  214. TSocketObjPtrPool::IndexSet indexes;
  215. m_bfActiveSockets.CopyIndexes(indexes);
  216. for(auto it = indexes.begin(), end = indexes.end(); it != end; ++it)
  217. Disconnect(*it);
  218. }
  219. void CTcpServer::WaitForClientSocketClose()
  220. {
  221. while(m_bfActiveSockets.Elements() > 0)
  222. ::WaitFor(50);
  223. }
  224. void CTcpServer::WaitForWorkerThreadEnd()
  225. {
  226. m_ioDispatcher.Stop();
  227. }
  228. void CTcpServer::ReleaseClientSocket()
  229. {
  230. VERIFY(m_bfActiveSockets.IsEmpty());
  231. m_bfActiveSockets.Reset();
  232. }
  233. void CTcpServer::ReleaseFreeSocket()
  234. {
  235. m_lsFreeSocket.Clear();
  236. #ifdef USE_EXTERNAL_GC
  237. if(IS_VALID_FD(m_fdGCTimer))
  238. {
  239. close(m_fdGCTimer);
  240. m_fdGCTimer = INVALID_FD;
  241. }
  242. #endif
  243. ReleaseGCSocketObj(TRUE);
  244. VERIFY(m_lsGCSocket.IsEmpty());
  245. }
  246. void CTcpServer::Reset()
  247. {
  248. m_phSocket.Reset();
  249. m_bfObjPool.Clear();
  250. m_rcBuffers = nullptr;
  251. m_soListens = nullptr;
  252. m_enState = SS_STOPPED;
  253. m_evWait.SyncNotifyAll();
  254. }
  255. TSocketObj* CTcpServer::GetFreeSocketObj(CONNID dwConnID, SOCKET soClient)
  256. {
  257. DWORD dwIndex;
  258. TSocketObj* pSocketObj = nullptr;
  259. if(m_lsFreeSocket.TryLock(&pSocketObj, dwIndex))
  260. {
  261. if(::GetTimeGap32(pSocketObj->freeTime) >= m_dwFreeSocketObjLockTime)
  262. VERIFY(m_lsFreeSocket.ReleaseLock(nullptr, dwIndex));
  263. else
  264. {
  265. VERIFY(m_lsFreeSocket.ReleaseLock(pSocketObj, dwIndex));
  266. pSocketObj = nullptr;
  267. }
  268. }
  269. if(!pSocketObj) pSocketObj = CreateSocketObj();
  270. pSocketObj->Reset(dwConnID, soClient);
  271. return pSocketObj;
  272. }
  273. TSocketObj* CTcpServer::CreateSocketObj()
  274. {
  275. return TSocketObj::Construct(m_phSocket, m_bfObjPool);
  276. }
  277. void CTcpServer::DeleteSocketObj(TSocketObj* pSocketObj)
  278. {
  279. TSocketObj::Destruct(pSocketObj);
  280. }
  281. void CTcpServer::AddFreeSocketObj(TSocketObj* pSocketObj, EnSocketCloseFlag enFlag, EnSocketOperation enOperation, int iErrorCode)
  282. {
  283. if(!InvalidSocketObj(pSocketObj))
  284. return;
  285. CloseClientSocketObj(pSocketObj, enFlag, enOperation, iErrorCode);
  286. m_bfActiveSockets.Remove(pSocketObj->connID);
  287. TSocketObj::Release(pSocketObj);
  288. #ifndef USE_EXTERNAL_GC
  289. ReleaseGCSocketObj();
  290. #endif
  291. if(!m_lsFreeSocket.TryPut(pSocketObj))
  292. m_lsGCSocket.PushBack(pSocketObj);
  293. }
  294. void CTcpServer::ReleaseGCSocketObj(BOOL bForce)
  295. {
  296. ::ReleaseGCObj(m_lsGCSocket, m_dwFreeSocketObjLockTime, bForce);
  297. }
  298. BOOL CTcpServer::InvalidSocketObj(TSocketObj* pSocketObj)
  299. {
  300. return TSocketObj::InvalidSocketObj(pSocketObj);
  301. }
  302. void CTcpServer::AddClientSocketObj(CONNID dwConnID, TSocketObj* pSocketObj, const HP_SOCKADDR& remoteAddr)
  303. {
  304. ASSERT(FindSocketObj(dwConnID) == nullptr);
  305. pSocketObj->connTime = ::TimeGetTime();
  306. pSocketObj->activeTime = pSocketObj->connTime;
  307. remoteAddr.Copy(pSocketObj->remoteAddr);
  308. pSocketObj->SetConnected();
  309. VERIFY(m_bfActiveSockets.ReleaseLock(dwConnID, pSocketObj));
  310. }
  311. TSocketObj* CTcpServer::FindSocketObj(CONNID dwConnID)
  312. {
  313. TSocketObj* pSocketObj = nullptr;
  314. if(m_bfActiveSockets.Get(dwConnID, &pSocketObj) != TSocketObjPtrPool::GR_VALID)
  315. pSocketObj = nullptr;
  316. return pSocketObj;
  317. }
  318. void CTcpServer::CloseClientSocketObj(TSocketObj* pSocketObj, EnSocketCloseFlag enFlag, EnSocketOperation enOperation, int iErrorCode, int iShutdownFlag)
  319. {
  320. ASSERT(TSocketObj::IsExist(pSocketObj));
  321. if(enFlag == SCF_CLOSE)
  322. FireClose(pSocketObj, SO_CLOSE, SE_OK);
  323. else if(enFlag == SCF_ERROR)
  324. FireClose(pSocketObj, enOperation, iErrorCode);
  325. SOCKET socket = pSocketObj->socket;
  326. pSocketObj->socket = INVALID_SOCKET;
  327. ::ManualCloseSocket(socket, iShutdownFlag);
  328. }
  329. BOOL CTcpServer::GetListenAddress(TCHAR lpszAddress[], int& iAddressLen, USHORT& usPort)
  330. {
  331. ASSERT(lpszAddress != nullptr && iAddressLen > 0);
  332. if(!HasStarted())
  333. {
  334. ::SetLastError(ERROR_INVALID_STATE);
  335. return FALSE;
  336. }
  337. return ::GetSocketLocalAddress(m_soListens[0], lpszAddress, iAddressLen, usPort);
  338. }
  339. BOOL CTcpServer::GetLocalAddress(CONNID dwConnID, TCHAR lpszAddress[], int& iAddressLen, USHORT& usPort)
  340. {
  341. ASSERT(lpszAddress != nullptr && iAddressLen > 0);
  342. TSocketObj* pSocketObj = FindSocketObj(dwConnID);
  343. if(!TSocketObj::IsValid(pSocketObj))
  344. {
  345. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  346. return FALSE;
  347. }
  348. return ::GetSocketLocalAddress(pSocketObj->socket, lpszAddress, iAddressLen, usPort);
  349. }
  350. BOOL CTcpServer::GetRemoteAddress(CONNID dwConnID, TCHAR lpszAddress[], int& iAddressLen, USHORT& usPort)
  351. {
  352. ASSERT(lpszAddress != nullptr && iAddressLen > 0);
  353. TSocketObj* pSocketObj = FindSocketObj(dwConnID);
  354. if(!TSocketObj::IsExist(pSocketObj))
  355. {
  356. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  357. return FALSE;
  358. }
  359. ADDRESS_FAMILY usFamily;
  360. return ::sockaddr_IN_2_A(pSocketObj->remoteAddr, usFamily, lpszAddress, iAddressLen, usPort);
  361. }
  362. BOOL CTcpServer::SetConnectionExtra(CONNID dwConnID, PVOID pExtra)
  363. {
  364. TSocketObj* pSocketObj = FindSocketObj(dwConnID);
  365. return SetConnectionExtra(pSocketObj, pExtra);
  366. }
  367. BOOL CTcpServer::SetConnectionExtra(TSocketObj* pSocketObj, PVOID pExtra)
  368. {
  369. if(!TSocketObj::IsExist(pSocketObj))
  370. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  371. else
  372. {
  373. pSocketObj->extra = pExtra;
  374. return TRUE;
  375. }
  376. return FALSE;
  377. }
  378. BOOL CTcpServer::GetConnectionExtra(CONNID dwConnID, PVOID* ppExtra)
  379. {
  380. TSocketObj* pSocketObj = FindSocketObj(dwConnID);
  381. return GetConnectionExtra(pSocketObj, ppExtra);
  382. }
  383. BOOL CTcpServer::GetConnectionExtra(TSocketObj* pSocketObj, PVOID* ppExtra)
  384. {
  385. ASSERT(ppExtra != nullptr);
  386. if(!TSocketObj::IsExist(pSocketObj))
  387. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  388. else
  389. {
  390. *ppExtra = pSocketObj->extra;
  391. return TRUE;
  392. }
  393. return FALSE;
  394. }
  395. BOOL CTcpServer::SetConnectionReserved(CONNID dwConnID, PVOID pReserved)
  396. {
  397. TSocketObj* pSocketObj = FindSocketObj(dwConnID);
  398. return SetConnectionReserved(pSocketObj, pReserved);
  399. }
  400. BOOL CTcpServer::SetConnectionReserved(TSocketObj* pSocketObj, PVOID pReserved)
  401. {
  402. if(!TSocketObj::IsExist(pSocketObj))
  403. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  404. else
  405. {
  406. pSocketObj->reserved = pReserved;
  407. return TRUE;
  408. }
  409. return FALSE;
  410. }
  411. BOOL CTcpServer::GetConnectionReserved(CONNID dwConnID, PVOID* ppReserved)
  412. {
  413. TSocketObj* pSocketObj = FindSocketObj(dwConnID);
  414. return GetConnectionReserved(pSocketObj, ppReserved);
  415. }
  416. BOOL CTcpServer::GetConnectionReserved(TSocketObj* pSocketObj, PVOID* ppReserved)
  417. {
  418. ASSERT(ppReserved != nullptr);
  419. if(!TSocketObj::IsExist(pSocketObj))
  420. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  421. else
  422. {
  423. *ppReserved = pSocketObj->reserved;
  424. return TRUE;
  425. }
  426. return FALSE;
  427. }
  428. BOOL CTcpServer::SetConnectionReserved2(CONNID dwConnID, PVOID pReserved2)
  429. {
  430. TSocketObj* pSocketObj = FindSocketObj(dwConnID);
  431. return SetConnectionReserved2(pSocketObj, pReserved2);
  432. }
  433. BOOL CTcpServer::SetConnectionReserved2(TSocketObj* pSocketObj, PVOID pReserved2)
  434. {
  435. if(!TSocketObj::IsExist(pSocketObj))
  436. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  437. else
  438. {
  439. pSocketObj->reserved2 = pReserved2;
  440. return TRUE;
  441. }
  442. return FALSE;
  443. }
  444. BOOL CTcpServer::GetConnectionReserved2(CONNID dwConnID, PVOID* ppReserved2)
  445. {
  446. TSocketObj* pSocketObj = FindSocketObj(dwConnID);
  447. return GetConnectionReserved2(pSocketObj, ppReserved2);
  448. }
  449. BOOL CTcpServer::GetConnectionReserved2(TSocketObj* pSocketObj, PVOID* ppReserved2)
  450. {
  451. ASSERT(ppReserved2 != nullptr);
  452. if(!TSocketObj::IsExist(pSocketObj))
  453. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  454. else
  455. {
  456. *ppReserved2 = pSocketObj->reserved2;
  457. return TRUE;
  458. }
  459. return FALSE;
  460. }
  461. BOOL CTcpServer::IsPauseReceive(CONNID dwConnID, BOOL& bPaused)
  462. {
  463. TSocketObj* pSocketObj = FindSocketObj(dwConnID);
  464. if(!TSocketObj::IsValid(pSocketObj))
  465. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  466. else
  467. {
  468. bPaused = pSocketObj->paused;
  469. return TRUE;
  470. }
  471. return FALSE;
  472. }
  473. BOOL CTcpServer::IsConnected(CONNID dwConnID)
  474. {
  475. TSocketObj* pSocketObj = FindSocketObj(dwConnID);
  476. if(!TSocketObj::IsValid(pSocketObj))
  477. {
  478. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  479. return FALSE;
  480. }
  481. return pSocketObj->HasConnected();
  482. }
  483. BOOL CTcpServer::GetPendingDataLength(CONNID dwConnID, int& iPending)
  484. {
  485. TSocketObj* pSocketObj = FindSocketObj(dwConnID);
  486. if(!TSocketObj::IsValid(pSocketObj))
  487. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  488. else
  489. {
  490. iPending = pSocketObj->Pending();
  491. return TRUE;
  492. }
  493. return FALSE;
  494. }
  495. DWORD CTcpServer::GetConnectionCount()
  496. {
  497. return m_bfActiveSockets.Elements();
  498. }
  499. BOOL CTcpServer::GetAllConnectionIDs(CONNID pIDs[], DWORD& dwCount)
  500. {
  501. return m_bfActiveSockets.GetAllElementIndexes(pIDs, dwCount);
  502. }
  503. BOOL CTcpServer::GetConnectPeriod(CONNID dwConnID, DWORD& dwPeriod)
  504. {
  505. BOOL isOK = TRUE;
  506. TSocketObj* pSocketObj = FindSocketObj(dwConnID);
  507. if(TSocketObj::IsValid(pSocketObj))
  508. dwPeriod = ::GetTimeGap32(pSocketObj->connTime);
  509. else
  510. {
  511. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  512. isOK = FALSE;
  513. }
  514. return isOK;
  515. }
  516. BOOL CTcpServer::GetSilencePeriod(CONNID dwConnID, DWORD& dwPeriod)
  517. {
  518. if(!m_bMarkSilence)
  519. return FALSE;
  520. BOOL isOK = TRUE;
  521. TSocketObj* pSocketObj = FindSocketObj(dwConnID);
  522. if(TSocketObj::IsValid(pSocketObj))
  523. dwPeriod = ::GetTimeGap32(pSocketObj->activeTime);
  524. else
  525. {
  526. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  527. isOK = FALSE;
  528. }
  529. return isOK;
  530. }
  531. BOOL CTcpServer::Disconnect(CONNID dwConnID, BOOL bForce)
  532. {
  533. TSocketObj* pSocketObj = FindSocketObj(dwConnID);
  534. if(!TSocketObj::IsValid(pSocketObj))
  535. {
  536. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  537. return FALSE;
  538. }
  539. return m_ioDispatcher.SendCommandByFD(pSocketObj->socket, DISP_CMD_DISCONNECT, dwConnID, bForce);
  540. }
  541. BOOL CTcpServer::DisconnectLongConnections(DWORD dwPeriod, BOOL bForce)
  542. {
  543. if(dwPeriod > MAX_CONNECTION_PERIOD)
  544. return FALSE;
  545. if(m_bfActiveSockets.Elements() == 0)
  546. return TRUE;
  547. DWORD now = ::TimeGetTime();
  548. TSocketObjPtrPool::IndexSet indexes;
  549. m_bfActiveSockets.CopyIndexes(indexes);
  550. for(auto it = indexes.begin(), end = indexes.end(); it != end; ++it)
  551. {
  552. CONNID connID = *it;
  553. TSocketObj* pSocketObj = FindSocketObj(connID);
  554. if(TSocketObj::IsValid(pSocketObj) && (int)(now - pSocketObj->connTime) >= (int)dwPeriod)
  555. Disconnect(connID, bForce);
  556. }
  557. return TRUE;
  558. }
  559. BOOL CTcpServer::DisconnectSilenceConnections(DWORD dwPeriod, BOOL bForce)
  560. {
  561. if(!m_bMarkSilence)
  562. return FALSE;
  563. if(dwPeriod > MAX_CONNECTION_PERIOD)
  564. return FALSE;
  565. if(m_bfActiveSockets.Elements() == 0)
  566. return TRUE;
  567. DWORD now = ::TimeGetTime();
  568. TSocketObjPtrPool::IndexSet indexes;
  569. m_bfActiveSockets.CopyIndexes(indexes);
  570. for(auto it = indexes.begin(), end = indexes.end(); it != end; ++it)
  571. {
  572. CONNID connID = *it;
  573. TSocketObj* pSocketObj = FindSocketObj(connID);
  574. if(TSocketObj::IsValid(pSocketObj) && (int)(now - pSocketObj->activeTime) >= (int)dwPeriod)
  575. Disconnect(connID, bForce);
  576. }
  577. return TRUE;
  578. }
  579. BOOL CTcpServer::PauseReceive(CONNID dwConnID, BOOL bPause)
  580. {
  581. TSocketObj* pSocketObj = FindSocketObj(dwConnID);
  582. if(!TSocketObj::IsValid(pSocketObj))
  583. {
  584. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  585. return FALSE;
  586. }
  587. if(pSocketObj->paused == bPause)
  588. return TRUE;
  589. pSocketObj->paused = bPause;
  590. if(!bPause)
  591. return m_ioDispatcher.SendCommandByFD(pSocketObj->socket, DISP_CMD_UNPAUSE, pSocketObj->connID);
  592. return TRUE;
  593. }
  594. BOOL CTcpServer::OnBeforeProcessIo(const TDispContext* pContext, PVOID pv, UINT events)
  595. {
  596. if(pv == &m_soListens[pContext->GetIndex()])
  597. {
  598. HandleAccept(pContext, events);
  599. return FALSE;
  600. }
  601. else if(pv == this)
  602. {
  603. ReleaseGCSocketObj(FALSE);
  604. ::ReadTimer(m_fdGCTimer);
  605. return FALSE;
  606. }
  607. TSocketObj* pSocketObj = (TSocketObj*)(pv);
  608. if(!TSocketObj::IsValid(pSocketObj))
  609. return FALSE;
  610. if(events & _EPOLL_ALL_ERROR_EVENTS)
  611. pSocketObj->SetConnected(FALSE);
  612. pSocketObj->Increment();
  613. if(!TSocketObj::IsValid(pSocketObj))
  614. {
  615. pSocketObj->Decrement();
  616. return FALSE;
  617. }
  618. return TRUE;
  619. }
  620. VOID CTcpServer::OnAfterProcessIo(const TDispContext* pContext, PVOID pv, UINT events, BOOL rs)
  621. {
  622. TSocketObj* pSocketObj = (TSocketObj*)(pv);
  623. if(TSocketObj::IsValid(pSocketObj))
  624. {
  625. ASSERT(rs && !(events & (EPOLLERR | EPOLLHUP | EPOLLRDHUP)));
  626. UINT evts = (pSocketObj->IsPending() ? EPOLLOUT : 0) | (pSocketObj->IsPaused() ? 0 : EPOLLIN);
  627. m_ioDispatcher.ModFD(pSocketObj->socket, evts | EPOLLRDHUP, pSocketObj);
  628. }
  629. pSocketObj->Decrement();
  630. }
  631. VOID CTcpServer::OnCommand(const TDispContext* pContext, TDispCommand* pCmd)
  632. {
  633. switch(pCmd->type)
  634. {
  635. case DISP_CMD_SEND:
  636. HandleCmdSend(pContext, (CONNID)(pCmd->wParam));
  637. break;
  638. case DISP_CMD_UNPAUSE:
  639. HandleCmdUnpause(pContext, (CONNID)(pCmd->wParam));
  640. break;
  641. case DISP_CMD_DISCONNECT:
  642. HandleCmdDisconnect(pContext, (CONNID)(pCmd->wParam), (BOOL)pCmd->lParam);
  643. break;
  644. }
  645. }
  646. VOID CTcpServer::HandleCmdSend(const TDispContext* pContext, CONNID dwConnID)
  647. {
  648. TSocketObj* pSocketObj = FindSocketObj(dwConnID);
  649. if(TSocketObj::IsValid(pSocketObj) && pSocketObj->IsPending())
  650. m_ioDispatcher.ProcessIo(pContext, pSocketObj, EPOLLOUT);
  651. }
  652. VOID CTcpServer::HandleCmdUnpause(const TDispContext* pContext, CONNID dwConnID)
  653. {
  654. TSocketObj* pSocketObj = FindSocketObj(dwConnID);
  655. if(!TSocketObj::IsValid(pSocketObj))
  656. return;
  657. if(BeforeUnpause(pSocketObj))
  658. m_ioDispatcher.ProcessIo(pContext, pSocketObj, EPOLLIN);
  659. else
  660. AddFreeSocketObj(pSocketObj, SCF_ERROR, SO_RECEIVE, ENSURE_ERROR_CANCELLED);
  661. }
  662. VOID CTcpServer::HandleCmdDisconnect(const TDispContext* pContext, CONNID dwConnID, BOOL bForce)
  663. {
  664. TSocketObj* pSocketObj = FindSocketObj(dwConnID);
  665. if(TSocketObj::IsValid(pSocketObj))
  666. m_ioDispatcher.ProcessIo(pContext, pSocketObj, EPOLLHUP);
  667. }
  668. BOOL CTcpServer::OnReadyRead(const TDispContext* pContext, PVOID pv, UINT events)
  669. {
  670. return HandleReceive(pContext, (TSocketObj*)pv, RETRIVE_EVENT_FLAG_H(events));
  671. }
  672. BOOL CTcpServer::OnReadyWrite(const TDispContext* pContext, PVOID pv, UINT events)
  673. {
  674. return HandleSend(pContext, (TSocketObj*)pv, RETRIVE_EVENT_FLAG_H(events));
  675. }
  676. BOOL CTcpServer::OnHungUp(const TDispContext* pContext, PVOID pv, UINT events)
  677. {
  678. return HandleClose(pContext, (TSocketObj*)pv, SCF_CLOSE, events);
  679. }
  680. BOOL CTcpServer::OnError(const TDispContext* pContext, PVOID pv, UINT events)
  681. {
  682. return HandleClose(pContext, (TSocketObj*)pv, SCF_ERROR, events);
  683. }
  684. VOID CTcpServer::OnDispatchThreadStart(THR_ID tid)
  685. {
  686. OnWorkerThreadStart(tid);
  687. }
  688. VOID CTcpServer::OnDispatchThreadEnd(THR_ID tid)
  689. {
  690. OnWorkerThreadEnd(tid);
  691. }
  692. BOOL CTcpServer::HandleClose(const TDispContext* pContext, TSocketObj* pSocketObj, EnSocketCloseFlag enFlag, UINT events)
  693. {
  694. EnSocketOperation enOperation = SO_CLOSE;
  695. if(events & _EPOLL_HUNGUP_EVENTS)
  696. enOperation = SO_CLOSE;
  697. else if(events & EPOLLIN)
  698. enOperation = SO_RECEIVE;
  699. else if(events & EPOLLOUT)
  700. enOperation = SO_SEND;
  701. int iErrorCode = 0;
  702. if(enFlag == SCF_ERROR)
  703. iErrorCode = ::SSO_GetError(pSocketObj->socket);
  704. AddFreeSocketObj(pSocketObj, enFlag, enOperation, iErrorCode);
  705. return TRUE;
  706. }
  707. BOOL CTcpServer::HandleAccept(const TDispContext* pContext, UINT events)
  708. {
  709. if(events & _EPOLL_ALL_ERROR_EVENTS)
  710. {
  711. ASSERT(!HasStarted());
  712. return FALSE;
  713. }
  714. while(TRUE)
  715. {
  716. HP_SOCKADDR addr;
  717. socklen_t addrLen = (socklen_t)addr.AddrSize();
  718. SOCKET soClient = ::accept(m_soListens[pContext->GetIndex()], addr.Addr(), &addrLen);
  719. if(soClient == INVALID_SOCKET)
  720. {
  721. int code = ::WSAGetLastError();
  722. switch(code)
  723. {
  724. case ERROR_WOULDBLOCK : return TRUE;
  725. case ERROR_CONNABORTED : continue;
  726. case ERROR_HANDLES_CLOSED : return FALSE;
  727. default : ERROR_EXIT2(EXIT_CODE_SOFTWARE, code);
  728. }
  729. }
  730. if(!::fcntl_SETFL(soClient, O_NOATIME | O_NONBLOCK | O_CLOEXEC))
  731. {
  732. ::ManualCloseSocket(soClient, SHUT_RDWR);
  733. continue;
  734. }
  735. CONNID dwConnID = 0;
  736. if(!m_bfActiveSockets.AcquireLock(dwConnID))
  737. {
  738. ::ManualCloseSocket(soClient, SHUT_RDWR);
  739. continue;
  740. }
  741. TSocketObj* pSocketObj = GetFreeSocketObj(dwConnID, soClient);
  742. AddClientSocketObj(dwConnID, pSocketObj, addr);
  743. if(TRIGGER(FireAccept(pSocketObj)) == HR_ERROR)
  744. {
  745. AddFreeSocketObj(pSocketObj, SCF_NONE);
  746. continue;
  747. }
  748. UINT evts = (pSocketObj->IsPending() ? EPOLLOUT : 0) | (pSocketObj->IsPaused() ? 0 : EPOLLIN);
  749. if(!m_ioDispatcher.AddFD(pSocketObj->socket, evts | EPOLLRDHUP, pSocketObj))
  750. {
  751. AddFreeSocketObj(pSocketObj, SCF_ERROR, SO_ACCEPT, ::WSAGetLastError());
  752. continue;
  753. }
  754. }
  755. return TRUE;
  756. }
  757. BOOL CTcpServer::HandleReceive(const TDispContext* pContext, TSocketObj* pSocketObj, int flag)
  758. {
  759. ASSERT(TSocketObj::IsValid(pSocketObj));
  760. if(m_bMarkSilence) pSocketObj->activeTime = ::TimeGetTime();
  761. CBufferPtr& buffer = m_rcBuffers[pContext->GetIndex()];
  762. int reads = flag ? -1 : MAX_CONTINUE_READS;
  763. for(int i = 0; i < reads || reads < 0; i++)
  764. {
  765. if(pSocketObj->paused)
  766. break;
  767. int rc = (int)read(pSocketObj->socket, buffer.Ptr(), buffer.Size());
  768. if(rc > 0)
  769. {
  770. if(TRIGGER(FireReceive(pSocketObj, buffer.Ptr(), rc)) == HR_ERROR)
  771. {
  772. TRACE("<S-CNNID: %zu> OnReceive() event return 'HR_ERROR', connection will be closed !", pSocketObj->connID);
  773. AddFreeSocketObj(pSocketObj, SCF_ERROR, SO_RECEIVE, ENSURE_ERROR_CANCELLED);
  774. return FALSE;
  775. }
  776. }
  777. else if(rc == 0)
  778. {
  779. AddFreeSocketObj(pSocketObj, SCF_CLOSE, SO_RECEIVE, SE_OK);
  780. return FALSE;
  781. }
  782. else
  783. {
  784. ASSERT(rc == SOCKET_ERROR);
  785. int code = ::WSAGetLastError();
  786. if(code == ERROR_WOULDBLOCK)
  787. break;
  788. AddFreeSocketObj(pSocketObj, SCF_ERROR, SO_RECEIVE, code);
  789. return FALSE;
  790. }
  791. }
  792. return TRUE;
  793. }
  794. BOOL CTcpServer::HandleSend(const TDispContext* pContext, TSocketObj* pSocketObj, int flag)
  795. {
  796. ASSERT(TSocketObj::IsValid(pSocketObj));
  797. if(!pSocketObj->IsPending())
  798. return TRUE;
  799. BOOL bBlocked = FALSE;
  800. int writes = flag ? -1 : MAX_CONTINUE_WRITES;
  801. TBufferObjList& sndBuff = pSocketObj->sndBuff;
  802. TItemPtr itPtr(sndBuff);
  803. for(int i = 0; i < writes || writes < 0; i++)
  804. {
  805. {
  806. CReentrantCriSecLock locallock(pSocketObj->csSend);
  807. itPtr = sndBuff.PopFront();
  808. }
  809. if(!itPtr.IsValid())
  810. break;
  811. ASSERT(!itPtr->IsEmpty());
  812. if(!SendItem(pSocketObj, itPtr, bBlocked))
  813. return FALSE;
  814. if(bBlocked)
  815. {
  816. ASSERT(!itPtr->IsEmpty());
  817. CReentrantCriSecLock locallock(pSocketObj->csSend);
  818. sndBuff.PushFront(itPtr.Detach());
  819. break;
  820. }
  821. }
  822. return TRUE;
  823. }
  824. BOOL CTcpServer::SendItem(TSocketObj* pSocketObj, TItem* pItem, BOOL& bBlocked)
  825. {
  826. while(!pItem->IsEmpty())
  827. {
  828. int rc = (int)write(pSocketObj->socket, pItem->Ptr(), pItem->Size());
  829. if(rc > 0)
  830. {
  831. if(TRIGGER(FireSend(pSocketObj, pItem->Ptr(), rc)) == HR_ERROR)
  832. {
  833. TRACE("<S-CNNID: %zu> OnSend() event should not return 'HR_ERROR' !!", pSocketObj->connID);
  834. ASSERT(FALSE);
  835. }
  836. pItem->Reduce(rc);
  837. }
  838. else if(rc == SOCKET_ERROR)
  839. {
  840. int code = ::WSAGetLastError();
  841. if(code == ERROR_WOULDBLOCK)
  842. {
  843. bBlocked = TRUE;
  844. break;
  845. }
  846. else
  847. {
  848. AddFreeSocketObj(pSocketObj, SCF_ERROR, SO_SEND, code);
  849. return FALSE;
  850. }
  851. }
  852. else
  853. ASSERT(FALSE);
  854. }
  855. return TRUE;
  856. }
  857. BOOL CTcpServer::Send(CONNID dwConnID, const BYTE* pBuffer, int iLength, int iOffset)
  858. {
  859. ASSERT(pBuffer && iLength > 0);
  860. if(iOffset != 0) pBuffer += iOffset;
  861. WSABUF buffer;
  862. buffer.len = iLength;
  863. buffer.buf = (BYTE*)pBuffer;
  864. return SendPackets(dwConnID, &buffer, 1);
  865. }
  866. BOOL CTcpServer::DoSendPackets(CONNID dwConnID, const WSABUF pBuffers[], int iCount)
  867. {
  868. ASSERT(pBuffers && iCount > 0);
  869. TSocketObj* pSocketObj = FindSocketObj(dwConnID);
  870. if(!TSocketObj::IsValid(pSocketObj))
  871. {
  872. ::SetLastError(ERROR_OBJECT_NOT_FOUND);
  873. return FALSE;
  874. }
  875. return DoSendPackets(pSocketObj, pBuffers, iCount);
  876. }
  877. BOOL CTcpServer::DoSendPackets(TSocketObj* pSocketObj, const WSABUF pBuffers[], int iCount)
  878. {
  879. ASSERT(pSocketObj && pBuffers && iCount > 0);
  880. int result = NO_ERROR;
  881. if(pBuffers && iCount > 0)
  882. {
  883. CLocalSafeCounter localcounter(*pSocketObj);
  884. CReentrantCriSecLock locallock(pSocketObj->csSend);
  885. if(TSocketObj::IsValid(pSocketObj))
  886. result = SendInternal(pSocketObj, pBuffers, iCount);
  887. else
  888. result = ERROR_OBJECT_NOT_FOUND;
  889. }
  890. else
  891. result = ERROR_INVALID_PARAMETER;
  892. if(result != NO_ERROR)
  893. ::SetLastError(result);
  894. return (result == NO_ERROR);
  895. }
  896. int CTcpServer::SendInternal(TSocketObj* pSocketObj, const WSABUF pBuffers[], int iCount)
  897. {
  898. BOOL bPending = pSocketObj->IsPending();
  899. for(int i = 0; i < iCount; i++)
  900. {
  901. int iBufLen = pBuffers[i].len;
  902. if(iBufLen > 0)
  903. {
  904. BYTE* pBuffer = (BYTE*)pBuffers[i].buf;
  905. ASSERT(pBuffer);
  906. pSocketObj->sndBuff.Cat(pBuffer, iBufLen);
  907. ASSERT(pSocketObj->sndBuff.Length() > 0);
  908. }
  909. }
  910. if(!bPending && pSocketObj->IsPending())
  911. {
  912. if(!m_ioDispatcher.SendCommandByFD(pSocketObj->socket, DISP_CMD_SEND, pSocketObj->connID))
  913. return ::GetLastError();
  914. }
  915. return NO_ERROR;
  916. }
  917. BOOL CTcpServer::SendSmallFile(CONNID dwConnID, LPCTSTR lpszFileName, const LPWSABUF pHead, const LPWSABUF pTail)
  918. {
  919. CFile file;
  920. CFileMapping fmap;
  921. WSABUF szBuf[3];
  922. HRESULT hr = ::MakeSmallFilePackage(lpszFileName, file, fmap, szBuf, pHead, pTail);
  923. if(FAILED(hr))
  924. {
  925. ::SetLastError(hr);
  926. return FALSE;
  927. }
  928. return SendPackets(dwConnID, szBuf, 3);
  929. }