HPThreadPool.cpp 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496
  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 "HPThreadPool.h"
  24. #include "common/FuncHelper.h"
  25. #include <pthread.h>
  26. LPTSocketTask CreateSocketTaskObj( Fn_SocketTaskProc fnTaskProc,
  27. PVOID pSender, CONNID dwConnID,
  28. LPCBYTE pBuffer, INT iBuffLen, EnTaskBufferType enBuffType,
  29. WPARAM wParam, LPARAM lParam)
  30. {
  31. ASSERT(fnTaskProc != nullptr);
  32. ASSERT(iBuffLen >= 0);
  33. LPTSocketTask pTask = new TSocketTask;
  34. pTask->fn = fnTaskProc;
  35. pTask->sender = pSender;
  36. pTask->connID = dwConnID;
  37. pTask->bufLen = iBuffLen;
  38. pTask->bufType = enBuffType;
  39. pTask->wparam = wParam;
  40. pTask->lparam = lParam;
  41. if(enBuffType != TBT_COPY || !pBuffer)
  42. pTask->buf = pBuffer;
  43. else
  44. {
  45. pTask->buf = MALLOC(BYTE, iBuffLen);
  46. ::CopyMemory((LPBYTE)pTask->buf, pBuffer, iBuffLen);
  47. }
  48. return pTask;
  49. }
  50. void DestroySocketTaskObj(LPTSocketTask pTask)
  51. {
  52. if(pTask)
  53. {
  54. if(pTask->bufType != TBT_REFER && pTask->buf)
  55. FREE(pTask->buf);
  56. delete pTask;
  57. }
  58. }
  59. volatile UINT CHPThreadPool::sm_uiNum = MAXUINT;
  60. LPCTSTR CHPThreadPool::POOLED_THREAD_PREFIX = _T("hp-pool-");
  61. BOOL CHPThreadPool::Start(DWORD dwThreadCount, DWORD dwMaxQueueSize, EnRejectedPolicy enRejectedPolicy, DWORD dwStackSize)
  62. {
  63. if(!CheckStarting())
  64. return FALSE;
  65. m_dwStackSize = dwStackSize;
  66. m_dwMaxQueueSize = dwMaxQueueSize;
  67. m_enRejectedPolicy = enRejectedPolicy;
  68. FireStartup();
  69. if(!InternalAdjustThreadCount(dwThreadCount))
  70. {
  71. EXECUTE_RESTORE_ERROR(Stop());
  72. return FALSE;
  73. }
  74. m_enState = SS_STARTED;
  75. return TRUE;
  76. }
  77. BOOL CHPThreadPool::Stop(DWORD dwMaxWait)
  78. {
  79. if(!CheckStoping())
  80. return FALSE;
  81. ::WaitFor(15);
  82. Shutdown(dwMaxWait);
  83. FireShutdown();
  84. Reset();
  85. return TRUE;
  86. }
  87. BOOL CHPThreadPool::Shutdown(DWORD dwMaxWait)
  88. {
  89. BOOL isOK = TRUE;
  90. BOOL bLimited = (m_dwMaxQueueSize != 0);
  91. BOOL bInfinite = (dwMaxWait == (DWORD)INFINITE || dwMaxWait == 0);
  92. auto prdShutdown = [this]() {return m_stThreads.empty();};
  93. if(m_enRejectedPolicy == TRP_WAIT_FOR && bLimited)
  94. m_evQueue.SyncNotifyAll();
  95. VERIFY(DoAdjustThreadCount(0));
  96. if(bInfinite)
  97. m_evShutdown.Wait(prdShutdown);
  98. else
  99. m_evShutdown.WaitFor(dwMaxWait, prdShutdown);
  100. ASSERT(m_lsTasks.Size() == 0);
  101. ASSERT(m_stThreads.size() == 0);
  102. if(!m_lsTasks.IsEmpty())
  103. {
  104. TTask* pTask = nullptr;
  105. while(m_lsTasks.PopFront(&pTask))
  106. {
  107. if(pTask->freeArg)
  108. ::DestroySocketTaskObj((LPTSocketTask)pTask->arg);
  109. TTask::Destruct(pTask);
  110. }
  111. ::SetLastError(ERROR_CANCELLED);
  112. isOK = FALSE;
  113. }
  114. if(!m_stThreads.empty())
  115. {
  116. CCriSecLock lock(m_csThread);
  117. if(!m_stThreads.empty())
  118. {
  119. #if !defined(__ANDROID__)
  120. for(auto it = m_stThreads.begin(), end = m_stThreads.end(); it != end; ++it)
  121. pthread_cancel(*it);
  122. #endif
  123. m_stThreads.clear();
  124. ::SetLastError(ERROR_CANCELLED);
  125. isOK = FALSE;
  126. }
  127. }
  128. return isOK;
  129. }
  130. BOOL CHPThreadPool::Submit(Fn_TaskProc fnTaskProc, PVOID pvArg, DWORD dwMaxWait)
  131. {
  132. return DoSubmit(fnTaskProc, pvArg, FALSE, dwMaxWait);
  133. }
  134. BOOL CHPThreadPool::Submit(LPTSocketTask pTask, DWORD dwMaxWait)
  135. {
  136. return DoSubmit((Fn_TaskProc)pTask->fn, (PVOID)pTask, TRUE, dwMaxWait);
  137. }
  138. BOOL CHPThreadPool::DoSubmit(Fn_TaskProc fnTaskProc, PVOID pvArg, BOOL bFreeArg, DWORD dwMaxWait)
  139. {
  140. EnSubmitResult sr = DirectSubmit(fnTaskProc, pvArg, bFreeArg);
  141. if(sr != SUBMIT_FULL)
  142. return (sr == SUBMIT_OK);
  143. if(m_enRejectedPolicy == TRP_CALL_FAIL)
  144. {
  145. ::SetLastError(ERROR_DESTINATION_ELEMENT_FULL);
  146. return FALSE;
  147. }
  148. else if(m_enRejectedPolicy == TRP_WAIT_FOR)
  149. {
  150. return CycleWaitSubmit(fnTaskProc, pvArg, dwMaxWait, bFreeArg);
  151. }
  152. else if(m_enRejectedPolicy == TRP_CALLER_RUN)
  153. {
  154. DoRunTaskProc(fnTaskProc, pvArg, bFreeArg);
  155. }
  156. else
  157. {
  158. ASSERT(FALSE);
  159. ::SetLastError(ERROR_INVALID_PARAMETER);
  160. return FALSE;
  161. }
  162. return TRUE;
  163. }
  164. void CHPThreadPool::DoRunTaskProc(Fn_TaskProc fnTaskProc, PVOID pvArg, BOOL bFreeArg)
  165. {
  166. ::InterlockedIncrement(&m_dwTaskCount);
  167. fnTaskProc(pvArg);
  168. ::InterlockedDecrement(&m_dwTaskCount);
  169. if(bFreeArg)
  170. ::DestroySocketTaskObj((LPTSocketTask)pvArg);
  171. }
  172. CHPThreadPool::EnSubmitResult CHPThreadPool::DirectSubmit(Fn_TaskProc fnTaskProc, PVOID pvArg, BOOL bFreeArg)
  173. {
  174. if(!CheckStarted())
  175. return SUBMIT_ERROR;
  176. BOOL bLimited = (m_dwMaxQueueSize != 0);
  177. if(bLimited && m_lsTasks.Size() >= m_dwMaxQueueSize)
  178. return SUBMIT_FULL;
  179. else
  180. {
  181. TTask* pTask = TTask::Construct(fnTaskProc, pvArg, bFreeArg);
  182. m_lsTasks.PushBack(pTask);
  183. m_evTask.SyncNotifyOne();
  184. }
  185. return SUBMIT_OK;
  186. }
  187. BOOL CHPThreadPool::CycleWaitSubmit(Fn_TaskProc fnTaskProc, PVOID pvArg, DWORD dwMaxWait, BOOL bFreeArg)
  188. {
  189. ASSERT(m_dwMaxQueueSize != 0);
  190. DWORD dwTime = ::TimeGetTime();
  191. BOOL bInfinite = (dwMaxWait == (DWORD)INFINITE || dwMaxWait == 0);
  192. auto prdQueue = [this]() {return (m_lsTasks.Size() < m_dwMaxQueueSize);};
  193. while(CheckStarted())
  194. {
  195. EnSubmitResult sr = DirectSubmit(fnTaskProc, pvArg, bFreeArg);
  196. if(sr == SUBMIT_OK)
  197. return TRUE;
  198. if(sr == SUBMIT_ERROR)
  199. return FALSE;
  200. if(bInfinite)
  201. m_evQueue.Wait(prdQueue);
  202. else
  203. {
  204. DWORD dwNow = ::GetTimeGap32(dwTime);
  205. if(dwNow > dwMaxWait || !m_evQueue.WaitFor(chrono::milliseconds(dwMaxWait - dwNow), prdQueue))
  206. {
  207. ::SetLastError(ERROR_TIMEOUT);
  208. break;
  209. }
  210. }
  211. }
  212. return FALSE;
  213. }
  214. BOOL CHPThreadPool::AdjustThreadCount(DWORD dwNewThreadCount)
  215. {
  216. if(!CheckStarted())
  217. return FALSE;
  218. return InternalAdjustThreadCount(dwNewThreadCount);
  219. }
  220. BOOL CHPThreadPool::InternalAdjustThreadCount(DWORD dwNewThreadCount)
  221. {
  222. int iNewThreadCount = (int)dwNewThreadCount;
  223. if(iNewThreadCount == 0)
  224. iNewThreadCount = ::GetDefaultWorkerThreadCount();
  225. else if(iNewThreadCount < 0)
  226. iNewThreadCount = PROCESSOR_COUNT * (-iNewThreadCount);
  227. return DoAdjustThreadCount((DWORD)iNewThreadCount);
  228. }
  229. BOOL CHPThreadPool::DoAdjustThreadCount(DWORD dwNewThreadCount)
  230. {
  231. ASSERT((int)dwNewThreadCount >= 0);
  232. BOOL bRemove = FALSE;
  233. DWORD dwThreadCount = 0;
  234. if(dwNewThreadCount > m_dwThreadCount)
  235. {
  236. dwThreadCount = dwNewThreadCount - m_dwThreadCount;
  237. return CreateWorkerThreads(dwThreadCount);
  238. }
  239. else if(dwNewThreadCount < m_dwThreadCount)
  240. {
  241. bRemove = TRUE;
  242. dwThreadCount = m_dwThreadCount - dwNewThreadCount;
  243. ::InterlockedSub(&m_dwThreadCount, dwThreadCount);
  244. }
  245. if(bRemove)
  246. {
  247. for(DWORD i = 0; i < dwThreadCount; i++)
  248. m_evTask.SyncNotifyOne();
  249. }
  250. return TRUE;
  251. }
  252. BOOL CHPThreadPool::CreateWorkerThreads(DWORD dwThreadCount)
  253. {
  254. unique_ptr<pthread_attr_t> pThreadAttr;
  255. if(m_dwStackSize != 0)
  256. {
  257. pThreadAttr = make_unique<pthread_attr_t>();
  258. VERIFY_IS_NO_ERROR(pthread_attr_init(pThreadAttr.get()));
  259. int rs = pthread_attr_setstacksize(pThreadAttr.get(), m_dwStackSize);
  260. if(!IS_NO_ERROR(rs))
  261. {
  262. pthread_attr_destroy(pThreadAttr.get());
  263. ::SetLastError(rs);
  264. return FALSE;
  265. }
  266. }
  267. BOOL isOK = TRUE;
  268. for(DWORD i = 0; i < dwThreadCount; i++)
  269. {
  270. THR_ID dwThreadID;
  271. int rs = pthread_create(&dwThreadID, pThreadAttr.get(), ThreadProc, (PVOID)this);
  272. if(!IS_NO_ERROR(rs))
  273. {
  274. ::SetLastError(rs);
  275. isOK = FALSE;
  276. break;
  277. }
  278. ::InterlockedIncrement(&m_dwThreadCount);
  279. CCriSecLock lock(m_csThread);
  280. m_stThreads.emplace(dwThreadID);
  281. }
  282. if(pThreadAttr != nullptr)
  283. pthread_attr_destroy(pThreadAttr.get());
  284. return isOK;
  285. }
  286. PVOID CHPThreadPool::ThreadProc(LPVOID pv)
  287. {
  288. CHPThreadPool* pThis = (CHPThreadPool*)pv;
  289. ::SetSequenceThreadName(SELF_THREAD_ID, pThis->m_strPrefix, pThis->m_uiSeq);
  290. pThis->FireWorkerThreadStart();
  291. PVOID rs = (PVOID)(UINT_PTR)(pThis->WorkerProc());
  292. pThis->FireWorkerThreadEnd();
  293. return rs;
  294. }
  295. int CHPThreadPool::WorkerProc()
  296. {
  297. BOOL bLimited = (m_dwMaxQueueSize != 0);
  298. TTask* pTask = nullptr;
  299. auto prdTask = [this]() {return (!m_lsTasks.IsEmpty()) || (m_dwThreadCount < m_stThreads.size());};
  300. while(TRUE)
  301. {
  302. pTask = nullptr;
  303. while(m_lsTasks.PopFront(&pTask))
  304. {
  305. if(m_enRejectedPolicy == TRP_WAIT_FOR && bLimited)
  306. m_evQueue.SyncNotifyOne();
  307. DoRunTaskProc(pTask->fn, pTask->arg, pTask->freeArg);
  308. TTask::Destruct(pTask);
  309. }
  310. if(CheckWorkerThreadExit())
  311. break;
  312. m_evTask.Wait(prdTask);
  313. }
  314. return 0;
  315. }
  316. BOOL CHPThreadPool::CheckWorkerThreadExit()
  317. {
  318. BOOL bExit = FALSE;
  319. BOOL bShutdown = FALSE;
  320. if(m_dwThreadCount < m_stThreads.size())
  321. {
  322. CCriSecLock lock(m_csThread);
  323. if(m_dwThreadCount < m_stThreads.size())
  324. {
  325. VERIFY(m_stThreads.erase(SELF_THREAD_ID) == 1);
  326. bExit = TRUE;
  327. bShutdown = m_stThreads.empty();
  328. }
  329. }
  330. if(bExit)
  331. {
  332. pthread_detach(SELF_THREAD_ID);
  333. if(bShutdown)
  334. m_evShutdown.SyncNotifyOne();
  335. }
  336. return bExit;
  337. }
  338. BOOL CHPThreadPool::CheckStarting()
  339. {
  340. if(::InterlockedCompareExchange(&m_enState, SS_STARTING, SS_STOPPED) != SS_STOPPED)
  341. {
  342. ::SetLastError(ERROR_INVALID_STATE);
  343. return FALSE;
  344. }
  345. return TRUE;
  346. }
  347. BOOL CHPThreadPool::CheckStarted()
  348. {
  349. if(m_enState != SS_STARTED)
  350. {
  351. ::SetLastError(ERROR_INVALID_STATE);
  352. return FALSE;
  353. }
  354. return TRUE;
  355. }
  356. BOOL CHPThreadPool::CheckStoping()
  357. {
  358. if( ::InterlockedCompareExchange(&m_enState, SS_STOPPING, SS_STARTED) != SS_STARTED &&
  359. ::InterlockedCompareExchange(&m_enState, SS_STOPPING, SS_STARTING) != SS_STARTING)
  360. {
  361. while(m_enState != SS_STOPPED)
  362. ::WaitFor(5);
  363. ::SetLastError(ERROR_INVALID_STATE);
  364. return FALSE;
  365. }
  366. return TRUE;
  367. }
  368. void CHPThreadPool::Reset(BOOL bSetWaitEvent)
  369. {
  370. m_uiSeq = MAXUINT;
  371. m_dwStackSize = 0;
  372. m_dwTaskCount = 0;
  373. m_dwThreadCount = 0;
  374. m_dwMaxQueueSize = 0;
  375. m_enRejectedPolicy = TRP_CALL_FAIL;
  376. m_enState = SS_STOPPED;
  377. if(bSetWaitEvent)
  378. m_evWait.SyncNotifyAll();
  379. }
  380. void CHPThreadPool::MakePrefix()
  381. {
  382. UINT uiNumber = ::InterlockedIncrement(&sm_uiNum);
  383. m_strPrefix.Format(_T("%s%u-"), POOLED_THREAD_PREFIX, uiNumber);
  384. }