TcpServer.cpp 25 KB

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