Thread.h 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612
  1. /*
  2. * Copyright: JessMA Open Source (ldcsaa@gmail.com)
  3. *
  4. * Author : Bruce Liang
  5. * Website : https://github.com/ldcsaa
  6. * Project : https://github.com/ldcsaa/HP-Socket
  7. * Blog : http://www.cnblogs.com/ldcsaa
  8. * Wiki : http://www.oschina.net/p/hp-socket
  9. * QQ Group : 44636872, 75375912
  10. *
  11. * Licensed under the Apache License, Version 2.0 (the "License");
  12. * you may not use this file except in compliance with the License.
  13. * You may obtain a copy of the License at
  14. *
  15. * http://www.apache.org/licenses/LICENSE-2.0
  16. *
  17. * Unless required by applicable law or agreed to in writing, software
  18. * distributed under the License is distributed on an "AS IS" BASIS,
  19. * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
  20. * See the License for the specific language governing permissions and
  21. * limitations under the License.
  22. */
  23. #pragma once
  24. #include "hpsocket/GlobalDef.h"
  25. #include "hpsocket/GlobalErrno.h"
  26. #include "RWLock.h"
  27. #include "STLHelper.h"
  28. #include <pthread.h>
  29. #include <signal.h>
  30. #include <utility>
  31. using namespace std;
  32. /* Used to retry syscalls that can return EINTR. */
  33. #define NO_EINTR_EXCEPT_THR_INTR(exp) ({ \
  34. long int _rc; \
  35. do {_rc = (long int)(exp);} \
  36. while (IS_HAS_ERROR(_rc) && IS_INTR_ERROR() \
  37. && !::IsThreadInterrupted()); \
  38. _rc; })
  39. #define NO_EINTR_EXCEPT_THR_INTR_INT(exp) ((int)NO_EINTR_EXCEPT_THR_INTR(exp))
  40. class __CThread_Interrupt_
  41. {
  42. public:
  43. static const int SIG_NO_INTERRUPT = (_NSIG - 5);
  44. private:
  45. friend BOOL IsThreadInterrupted();
  46. template<typename T, typename P, typename R> friend class CThread;
  47. private:
  48. static BOOL IsInterrupted();
  49. static BOOL InitSigAction();
  50. static void SignalHandler(int sig);
  51. private:
  52. ~__CThread_Interrupt_();
  53. private:
  54. static BOOL sm_bInitFlag;
  55. };
  56. inline BOOL IsThreadInterrupted() {return __CThread_Interrupt_::IsInterrupted();}
  57. class __CFakeRunnerClass_ {};
  58. template<class T, class P = VOID, class R = UINT_PTR> class CThread
  59. {
  60. public:
  61. using F = R (T::*)(P*);
  62. using SF = R (*)(P*);
  63. struct TWorker
  64. {
  65. CThread* m_pThread;
  66. BOOL m_bDetach;
  67. T* m_pRunner;
  68. F m_pFunc;
  69. P* m_pArg;
  70. public:
  71. TWorker(CThread* pThread, BOOL bDetach = FALSE, T* pRunner = nullptr, F pFunc = nullptr, P* pArg = nullptr)
  72. : m_pThread(pThread)
  73. {
  74. Reset(bDetach, pRunner, pFunc, pArg);
  75. }
  76. void Reset(BOOL bDetach = FALSE, T* pRunner = nullptr, F pFunc = nullptr, P* pArg = nullptr)
  77. {
  78. m_bDetach = bDetach;
  79. m_pRunner = pRunner;
  80. m_pFunc = pFunc;
  81. m_pArg = pArg;
  82. }
  83. public:
  84. template<typename T_, typename R_, typename = enable_if_t<!is_same<T_, __CFakeRunnerClass_>::value && !is_void<R_>::value>>
  85. PVOID Run(T_*, R_*)
  86. {
  87. return (PVOID)(UINT_PTR)((m_pRunner->*m_pFunc)(m_pArg));
  88. }
  89. template<typename T_, typename = enable_if_t<!is_same<T_, __CFakeRunnerClass_>::value>>
  90. PVOID Run(T_*, PVOID)
  91. {
  92. (m_pRunner->*m_pFunc)(m_pArg);
  93. return nullptr;
  94. }
  95. template<typename R_, typename = enable_if_t<!is_void<R_>::value>>
  96. PVOID Run(__CFakeRunnerClass_*, R_*)
  97. {
  98. return (PVOID)(UINT_PTR)(*(SF*)&m_pFunc)(m_pArg);
  99. }
  100. PVOID Run(__CFakeRunnerClass_*, VOID*)
  101. {
  102. (*(SF*)&m_pFunc)(m_pArg);
  103. return nullptr;
  104. }
  105. };
  106. friend struct TWorker;
  107. public:
  108. BOOL Start(SF pFunc, P* pArg = nullptr, BOOL bDetach = FALSE, const pthread_attr_t* pAttr = nullptr)
  109. {
  110. return Start((__CFakeRunnerClass_*)nullptr, *(F*)&pFunc, pArg, bDetach, pAttr);
  111. }
  112. BOOL Start(T* pRunner, F pFunc, P* pArg = nullptr, BOOL bDetach = FALSE, const pthread_attr_t* pAttr = nullptr)
  113. {
  114. int rs = ERROR_INVALID_STATE;
  115. if(IsRunning())
  116. ::SetLastError(rs);
  117. else
  118. {
  119. m_Worker.Reset(bDetach, pRunner, pFunc, pArg);
  120. SetRunning(TRUE);
  121. rs = pthread_create(&m_ulThreadID, pAttr, ThreadProc, (PVOID)(&m_Worker));
  122. if(rs != NO_ERROR)
  123. {
  124. Reset();
  125. ::SetLastError(rs);
  126. }
  127. }
  128. return (rs == NO_ERROR);
  129. }
  130. #if !defined(__ANDROID__)
  131. BOOL Cancel()
  132. {
  133. int rs = NO_ERROR;
  134. if(!IsRunning() || ::IsSelfThread(m_ulThreadID))
  135. rs = ERROR_INVALID_STATE;
  136. else
  137. rs = pthread_cancel(m_ulThreadID);
  138. if(rs != NO_ERROR)
  139. ::SetLastError(rs);
  140. return (rs == NO_ERROR);
  141. }
  142. BOOL Join(R* pResult = nullptr, BOOL bWait = TRUE, LONG lWaitMillsec = INFINITE)
  143. {
  144. int rs = NO_ERROR;
  145. if(!IsRunning() || ::IsSelfThread(m_ulThreadID))
  146. rs = ERROR_INVALID_STATE;
  147. else
  148. {
  149. if(!bWait)
  150. rs = pthread_tryjoin_np(m_ulThreadID, (PVOID*)pResult);
  151. else if(IS_INFINITE(lWaitMillsec))
  152. rs = pthread_join(m_ulThreadID, (PVOID*)pResult);
  153. else
  154. {
  155. timespec ts;
  156. ::GetFutureTimespec(lWaitMillsec, ts, CLOCK_REALTIME);
  157. rs = pthread_timedjoin_np(m_ulThreadID, (PVOID*)pResult, &ts);
  158. }
  159. }
  160. if(rs == NO_ERROR)
  161. SetRunning(FALSE);
  162. else
  163. ::SetLastError(rs);
  164. return (rs == NO_ERROR);
  165. }
  166. #else
  167. BOOL Cancel()
  168. {
  169. SetLastError(ERROR_CALL_NOT_IMPLEMENTED);
  170. return FALSE;
  171. }
  172. BOOL Join(R* pResult = nullptr)
  173. {
  174. int rs = NO_ERROR;
  175. if(!IsRunning() || ::IsSelfThread(m_ulThreadID))
  176. rs = ERROR_INVALID_STATE;
  177. else
  178. rs = pthread_join(m_ulThreadID, (PVOID*)pResult);
  179. if(rs == NO_ERROR)
  180. SetRunning(FALSE);
  181. else
  182. ::SetLastError(rs);
  183. return (rs == NO_ERROR);
  184. }
  185. #endif
  186. BOOL Detach()
  187. {
  188. int rs = NO_ERROR;
  189. if(!IsRunning())
  190. rs = ERROR_INVALID_STATE;
  191. else
  192. rs = pthread_detach(m_ulThreadID);
  193. if(rs == NO_ERROR)
  194. Reset();
  195. else
  196. ::SetLastError(rs);
  197. return (rs == NO_ERROR);
  198. }
  199. void Reset()
  200. {
  201. SetRunning(FALSE);
  202. m_ulThreadID = 0;
  203. m_lNativeID = 0;
  204. m_Worker.Reset();
  205. }
  206. BOOL Interrupt()
  207. {
  208. if(!IsRunning())
  209. {
  210. SetLastError(ERROR_INVALID_STATE);
  211. return FALSE;
  212. }
  213. return (IS_NO_ERROR(pthread_kill(m_ulThreadID, __CThread_Interrupt_::SIG_NO_INTERRUPT)));
  214. }
  215. void SetRunning(BOOL bRunning) {m_bRunning = bRunning;}
  216. BOOL IsRunning () const {return m_bRunning;}
  217. T* GetRunner () const {return m_Worker.m_pRunner;}
  218. F GetFunc () const {return m_Worker.m_pFunc;}
  219. SF GetSFunc () const {return *(SF*)&m_Worker.m_pFunc;}
  220. P* GetArg () const {return m_Worker.m_pArg;}
  221. THR_ID GetThreadID () const {return m_ulThreadID;}
  222. NTHR_ID GetNativeID () const {return m_lNativeID;}
  223. BOOL IsInMyThread () const {return IsMyThreadID(SELF_THREAD_ID);}
  224. BOOL IsMyThreadID (THR_ID ulThreadID) const {return ::IsSameThread(ulThreadID, m_ulThreadID);}
  225. BOOL IsMyNativeThreadID (NTHR_ID lNativeID) const {return ::IsSameNativeThread(lNativeID, m_lNativeID);}
  226. private:
  227. static PVOID ThreadProc(LPVOID pv)
  228. {
  229. UnmaskInterruptSignal();
  230. __CThread_Interrupt_ tlsInterrupt;
  231. TWorker* pWorker = (TWorker*)pv;
  232. if(pWorker->m_bDetach)
  233. pWorker->m_pThread->Detach();
  234. else
  235. pWorker->m_pThread->m_lNativeID = SELF_NATIVE_THREAD_ID;
  236. PVOID pResult = pWorker->Run((T*)nullptr, (R*)nullptr);
  237. return pResult;
  238. }
  239. static void UnmaskInterruptSignal()
  240. {
  241. sigset_t ss;
  242. sigemptyset(&ss);
  243. sigaddset(&ss, __CThread_Interrupt_::SIG_NO_INTERRUPT);
  244. pthread_sigmask(SIG_UNBLOCK, &ss, nullptr);
  245. }
  246. public:
  247. CThread()
  248. : m_Worker(this)
  249. {
  250. Reset();
  251. }
  252. virtual ~CThread()
  253. {
  254. if(IsRunning())
  255. {
  256. Interrupt();
  257. Join(nullptr);
  258. }
  259. ASSERT(!IsRunning());
  260. }
  261. DECLARE_NO_COPY_CLASS(CThread)
  262. private:
  263. THR_ID m_ulThreadID;
  264. NTHR_ID m_lNativeID;
  265. BOOL m_bRunning;
  266. TWorker m_Worker;
  267. };
  268. template<class P = VOID, class R = UINT_PTR> using CStaticThread = CThread<__CFakeRunnerClass_, P, R>;
  269. template<class T> class CTlsObj
  270. {
  271. using TLocalMap = unordered_map<THR_ID, T*>;
  272. public:
  273. T* TryGet()
  274. {
  275. T* pValue = nullptr;
  276. {
  277. CReadLock locallock(m_lock);
  278. auto it = m_map.find(SELF_THREAD_ID);
  279. if(it != m_map.end())
  280. pValue = it->second;
  281. }
  282. return pValue;
  283. }
  284. template<typename ... _Con_Param> T* Get(_Con_Param&& ... construct_args)
  285. {
  286. T* pValue = TryGet();
  287. if(pValue == nullptr)
  288. {
  289. pValue = Construct(forward<_Con_Param>(construct_args) ...);
  290. CWriteLock locallock(m_lock);
  291. m_map[SELF_THREAD_ID] = pValue;
  292. }
  293. return pValue;
  294. }
  295. template<typename ... _Con_Param> T& GetRef(_Con_Param&& ... construct_args)
  296. {
  297. return *Get(forward<_Con_Param>(construct_args) ...);
  298. }
  299. T* SetNewAndGetOld(T* pValue)
  300. {
  301. T* pOldValue = TryGet();
  302. if(pValue != pOldValue)
  303. {
  304. if(pValue == nullptr)
  305. DoRemove();
  306. else
  307. {
  308. CWriteLock locallock(m_lock);
  309. m_map[SELF_THREAD_ID] = pValue;
  310. }
  311. }
  312. return pOldValue;
  313. }
  314. void Set(T* pValue)
  315. {
  316. T* pOldValue = SetNewAndGetOld(pValue);
  317. if(pValue != pOldValue)
  318. DoDelete(pOldValue);
  319. }
  320. void Remove()
  321. {
  322. T* pValue = TryGet();
  323. if(pValue != nullptr)
  324. {
  325. DoDelete(pValue);
  326. DoRemove();
  327. }
  328. }
  329. void Clear()
  330. {
  331. CWriteLock locallock(m_lock);
  332. if(!IsEmpty())
  333. {
  334. for(auto it = m_map.begin(), end = m_map.end(); it != end; ++it)
  335. DoDelete(it->second);
  336. m_map.clear();
  337. }
  338. }
  339. TLocalMap& GetLocalMap() {return m_map;}
  340. const TLocalMap& GetLocalMap() const {return m_map;}
  341. CTlsObj& operator = (T* p) {Set(p); return *this;}
  342. T* operator -> () {return Get();}
  343. const T* operator -> () const {return Get();}
  344. T& operator * () {return GetRef();}
  345. const T& operator * () const {return GetRef();}
  346. size_t Size () const {return m_map.size();}
  347. bool IsEmpty() const {return m_map.empty();}
  348. private:
  349. inline void DoRemove()
  350. {
  351. CWriteLock locallock(m_lock);
  352. m_map.erase(SELF_THREAD_ID);
  353. }
  354. static inline void DoDelete(T* pValue)
  355. {
  356. if(pValue != nullptr)
  357. delete pValue;
  358. }
  359. template<typename ... _Con_Param> static inline T* Construct(_Con_Param&& ... construct_args)
  360. {
  361. return new T(forward<_Con_Param>(construct_args) ...);
  362. }
  363. public:
  364. CTlsObj()
  365. {
  366. }
  367. CTlsObj(T* pValue)
  368. {
  369. Set(pValue);
  370. }
  371. ~CTlsObj()
  372. {
  373. Clear();
  374. }
  375. private:
  376. CSimpleRWLock m_lock;
  377. TLocalMap m_map;
  378. DECLARE_NO_COPY_CLASS(CTlsObj)
  379. };
  380. template<class T> class CTlsSimple
  381. {
  382. using TLocalMap = unordered_map<THR_ID, T>;
  383. static const T DEFAULT = (T)(0);
  384. public:
  385. BOOL TryGet(T& tValue)
  386. {
  387. BOOL isOK = FALSE;
  388. {
  389. CReadLock locallock(m_lock);
  390. auto it = m_map.find(SELF_THREAD_ID);
  391. if(it != m_map.end())
  392. {
  393. tValue = it->second;
  394. isOK = TRUE;
  395. }
  396. }
  397. return isOK;
  398. }
  399. T Get(T tDefault = DEFAULT)
  400. {
  401. T tValue;
  402. if(TryGet(tValue))
  403. return tValue;
  404. Set(tDefault);
  405. return tDefault;
  406. }
  407. T SetNewAndGetOld(T tValue)
  408. {
  409. T tOldValue;
  410. if(!TryGet(tOldValue))
  411. tOldValue = DEFAULT;
  412. else if(tValue != tOldValue)
  413. Set(tValue);
  414. return tOldValue;
  415. }
  416. void Set(T tValue)
  417. {
  418. CWriteLock locallock(m_lock);
  419. m_map[SELF_THREAD_ID] = tValue;
  420. }
  421. void Remove()
  422. {
  423. T tValue;
  424. if(TryGet(tValue))
  425. {
  426. CWriteLock locallock(m_lock);
  427. m_map.erase(SELF_THREAD_ID);
  428. }
  429. }
  430. void Clear()
  431. {
  432. CWriteLock locallock(m_lock);
  433. if(!IsEmpty())
  434. m_map.clear();
  435. }
  436. TLocalMap& GetLocalMap() {return m_map;}
  437. const TLocalMap& GetLocalMap() const {return m_map;}
  438. CTlsSimple& operator = (T t) {Set(t); return *this;}
  439. BOOL operator == (T t) {return Get() == t;}
  440. BOOL operator != (T t) {return Get() != t;}
  441. BOOL operator >= (T t) {return Get() >= t;}
  442. BOOL operator <= (T t) {return Get() <= t;}
  443. BOOL operator > (T t) {return Get() > t;}
  444. BOOL operator < (T t) {return Get() < t;}
  445. size_t Size () const {return m_map.size();}
  446. bool IsEmpty() const {return m_map.empty();}
  447. public:
  448. CTlsSimple()
  449. {
  450. }
  451. CTlsSimple(T tValue)
  452. {
  453. Set(tValue);
  454. }
  455. ~CTlsSimple()
  456. {
  457. Clear();
  458. }
  459. DECLARE_NO_COPY_CLASS(CTlsSimple)
  460. private:
  461. CSimpleRWLock m_lock;
  462. TLocalMap m_map;
  463. };