CcosThread.cpp 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587
  1. #include "CcosThread.h"
  2. #include "MsgQueue.h"
  3. #include <iostream>
  4. #include <sys/time.h>
  5. #include <errno.h>
  6. #include <ctime>
  7. #include <thread>
  8. #include <stdexcept>
  9. #include <sstream>
  10. // 线程安全日志锁
  11. static std::mutex g_logMutex;
  12. #define LOG(...) { std::lock_guard<std::mutex> lock(g_logMutex); std::cout << __VA_ARGS__ << std::endl; }
  13. #define ERR_LOG(...) { std::lock_guard<std::mutex> lock(g_logMutex); std::cerr << __VA_ARGS__ << std::endl; }
  14. // Thread_Base 实现
  15. Thread_Base::Thread_Base(void)
  16. : m_Base_Thread(0),
  17. m_ThreadID(0),
  18. m_ThreadRunning(false),
  19. m_ExitRequested(false),
  20. m_strName(""),
  21. m_pThreadLogger(nullptr)
  22. {
  23. LOG("Creating exit/work/run events");
  24. m_ExitFlag = LinuxEvent::CreateEvent(LinuxEvent::MANUAL_RESET, false);
  25. m_WorkFlag = LinuxEvent::CreateEvent(LinuxEvent::AUTO_RESET, false);
  26. m_RunFlag = LinuxEvent::CreateEvent(LinuxEvent::MANUAL_RESET, false);
  27. if (!m_ExitFlag || !m_WorkFlag || !m_RunFlag) {
  28. throw std::runtime_error("Failed to create event objects");
  29. }
  30. // 初始化条件变量
  31. if (pthread_cond_init(&m_StateCond, nullptr) != 0) {
  32. throw std::runtime_error("Failed to initialize condition variable");
  33. }
  34. }
  35. Thread_Base::~Thread_Base(void)
  36. {
  37. StopThread();
  38. pthread_cond_destroy(&m_StateCond);
  39. }
  40. void Thread_Base::SetLogger(PVOID pLoger) {
  41. m_pThreadLogger = pLoger;
  42. }
  43. PVOID Thread_Base::GetLogger() {
  44. return m_pThreadLogger;
  45. }
  46. void* Thread_Base::Thread_Base_Thread(void* pPara) {
  47. try {
  48. Thread_Base* handle = static_cast<Thread_Base*>(pPara);
  49. if (!handle) {
  50. ERR_LOG("Invalid thread parameter");
  51. return nullptr;
  52. }
  53. usleep(30000); // Sleep(30) 替代
  54. if (!handle->RegistThread()) {
  55. goto endwork_entry;
  56. }
  57. if (!handle->OnStartThread()) {
  58. goto endwork_entry;
  59. }
  60. LOG("Thread " << handle->m_strName << " starting");
  61. handle->SetThreadOnTheRun(true);
  62. while (true) {
  63. if (!handle->Exec()) {
  64. break;
  65. }
  66. int waitevent = handle->WaitTheIncommingEvent(0);
  67. if (waitevent == 0) {
  68. break;
  69. }
  70. }
  71. endwork_entry:
  72. LOG("Thread " << handle->m_strName << " exiting");
  73. handle->ThreadExitProcedure();
  74. return nullptr;
  75. }
  76. catch (const std::exception& e) {
  77. ERR_LOG("[Thread] Exception: " << e.what());
  78. }
  79. catch (...) {
  80. ERR_LOG("[Thread] Unknown exception");
  81. }
  82. return nullptr;
  83. }
  84. bool Thread_Base::StartThread(bool Sync, bool Inherit)
  85. {
  86. if (m_strName.empty()) {
  87. LOG("Warning: Thread name is empty");
  88. }
  89. // 检查线程是否已运行
  90. if (m_ThreadRunning) {
  91. #ifdef __GNU_LIBRARY__ // 仅GNU环境支持非标准函数
  92. int tryJoin = pthread_tryjoin_np(m_Base_Thread, nullptr);
  93. if (tryJoin == EBUSY) {
  94. LOG("Thread " << m_strName << " is already running");
  95. return true;
  96. }
  97. #endif
  98. // 线程已结束,清理资源
  99. LOG("Thread " << m_strName << " has ended, cleaning up");
  100. pthread_join(m_Base_Thread, nullptr);
  101. m_ThreadRunning = false;
  102. }
  103. m_ExitFlag->ResetEvent();
  104. // 初始化线程属性
  105. pthread_attr_t attr;
  106. if (pthread_attr_init(&attr) != 0) {
  107. ERR_LOG("Failed to initialize thread attributes");
  108. return false;
  109. }
  110. // 设置栈大小
  111. size_t stack_size = 2 * 1024 * 1024; // 2MB
  112. if (pthread_attr_setstacksize(&attr, stack_size) != 0) {
  113. ERR_LOG("Failed to set thread stack size");
  114. pthread_attr_destroy(&attr);
  115. return false;
  116. }
  117. // 配置调度继承
  118. if (Inherit) {
  119. pthread_attr_setinheritsched(&attr, PTHREAD_INHERIT_SCHED);
  120. }
  121. else {
  122. pthread_attr_setinheritsched(&attr, PTHREAD_EXPLICIT_SCHED);
  123. }
  124. // 创建线程
  125. LOG("Creating thread " << m_strName);
  126. int result = pthread_create(&m_Base_Thread, &attr, Thread_Base_Thread, this);
  127. pthread_attr_destroy(&attr);
  128. if (result != 0) {
  129. ERR_LOG("Failed to create thread " << m_strName << ", error: " << result);
  130. m_ExitFlag->SetEvent();
  131. return false;
  132. }
  133. m_ThreadRunning = true;
  134. m_ThreadID = static_cast<unsigned long>(m_Base_Thread);
  135. LOG("Thread " << m_strName << " created, ID: " << m_ThreadID);
  136. // 同步等待线程启动
  137. if (Sync) {
  138. LOG("Waiting for thread " << m_strName << " to start");
  139. return WaitTheThreadOnTheRun(300000); // 300秒超时
  140. }
  141. return true;
  142. }
  143. bool Thread_Base::StopThread(unsigned long timeperiod)
  144. {
  145. bool ret = true;
  146. // 设置退出标志 - 通过信号量通知
  147. m_ExitFlag->SetEvent();
  148. // 检查是否在同一个线程内调用
  149. if (pthread_equal(pthread_self(), m_Base_Thread)) {
  150. return true;
  151. }
  152. if (m_Base_Thread) {
  153. struct timespec ts;
  154. // 计算绝对超时时间
  155. if (clock_gettime(CLOCK_REALTIME, &ts) == -1) {
  156. perror("clock_gettime");
  157. return false;
  158. }
  159. // 将毫秒转换为秒和纳秒
  160. ts.tv_sec += timeperiod / 1000;
  161. ts.tv_nsec += (timeperiod % 1000) * 1000000;
  162. // 处理纳秒溢出
  163. if (ts.tv_nsec >= 1000000000) {
  164. ts.tv_sec += 1;
  165. ts.tv_nsec -= 1000000000;
  166. }
  167. #ifdef __GNU_LIBRARY__
  168. int joinResult = pthread_timedjoin_np(m_Base_Thread, nullptr, &ts);
  169. #else
  170. // 标准环境下使用循环等待
  171. int joinResult = ETIMEDOUT;
  172. const int interval = 100; // 100ms轮询
  173. for (unsigned long elapsed = 0; elapsed < timeperiod; elapsed += interval) {
  174. int tryJoin = pthread_tryjoin_np(m_Base_Thread, nullptr);
  175. if (tryJoin == 0) {
  176. joinResult = 0;
  177. break;
  178. }
  179. usleep(interval * 1000);
  180. }
  181. #endif
  182. if (joinResult == ETIMEDOUT) {
  183. ERR_LOG("Timeout stopping thread " << m_strName);
  184. // 尝试取消线程
  185. pthread_cancel(m_Base_Thread);
  186. // 尝试再次等待较短时间
  187. clock_gettime(CLOCK_REALTIME, &ts);
  188. ts.tv_sec += 5; // 额外等待5秒
  189. ts.tv_nsec = 0;
  190. if (pthread_timedjoin_np(m_Base_Thread, NULL, &ts) == 0) {
  191. ERR_LOG("Failed to stop thread " << m_strName << ", detaching");
  192. pthread_detach(m_Base_Thread);
  193. ret = false;
  194. }
  195. // 执行线程退出清理
  196. ThreadExitProcedure();
  197. }
  198. else if (joinResult != 0) {
  199. ERR_LOG("Failed to join thread " << m_strName << ", error: " << joinResult);
  200. ret = false;
  201. }
  202. m_Base_Thread = 0;
  203. }
  204. m_ThreadID = 0;
  205. m_ThreadRunning = false;
  206. return ret;
  207. }
  208. void Thread_Base::NotifyExit()
  209. {
  210. // 发送退出信号
  211. m_ExitFlag->SetEvent();
  212. }
  213. int Thread_Base::WaitTheIncommingEvent(unsigned long waittime) // 修改为 unsigned long 以匹配毫秒超时
  214. {
  215. std::vector<std::shared_ptr<LinuxEvent>> waitEvents = {
  216. m_ExitFlag,
  217. m_WorkFlag
  218. };
  219. if (m_Base_Thread)
  220. {
  221. // 等待两个事件中的一个发生
  222. int wait = LinuxEvent::WaitForMultipleEvents(waitEvents, waittime);
  223. if (wait == 0)
  224. {
  225. return 0; // ExitFlag触发
  226. }
  227. else if (wait == 1)
  228. {
  229. return 1; // WorkFlag触发
  230. }
  231. else
  232. {
  233. // 超时或错误
  234. return -1;
  235. }
  236. }
  237. // 无线程对象,默认认为退出
  238. return 0;
  239. }
  240. bool Thread_Base::WaitTheThreadEnd(unsigned long waittime)
  241. {
  242. if (m_Base_Thread == 0) {
  243. return true; // 没有线程存在,直接返回成功
  244. }
  245. // 计算绝对超时时间
  246. struct timespec ts;
  247. if (clock_gettime(CLOCK_REALTIME, &ts) == -1) {
  248. perror("clock_gettime");
  249. return false;
  250. }
  251. ts.tv_sec += waittime / 1000;
  252. ts.tv_nsec += (waittime % 1000) * 1000000;
  253. if (ts.tv_nsec >= 1000000000) {
  254. ts.tv_sec += 1;
  255. ts.tv_nsec -= 1000000000;
  256. }
  257. #ifdef __GNU_LIBRARY__
  258. int result = pthread_timedjoin_np(m_Base_Thread, nullptr, &ts);
  259. #else
  260. int result = ETIMEDOUT;
  261. const int interval = 100;
  262. for (unsigned long elapsed = 0; elapsed < waittime; elapsed += interval) {
  263. if (pthread_tryjoin_np(m_Base_Thread, nullptr) == 0) {
  264. result = 0;
  265. break;
  266. }
  267. usleep(interval * 1000);
  268. }
  269. #endif
  270. return result == 0;
  271. }
  272. bool Thread_Base::WaitTheThreadEndSign(unsigned long waittime)
  273. {
  274. return m_Base_Thread ? m_ExitFlag->Wait(waittime) : true;
  275. }
  276. bool Thread_Base::SetThreadOnTheRun(bool OnTheRun)
  277. {
  278. if (OnTheRun)
  279. {
  280. m_RunFlag->SetEvent();
  281. }
  282. else
  283. {
  284. m_RunFlag->ResetEvent();
  285. }
  286. return true;
  287. }
  288. void Thread_Base::NotifyThreadWork()
  289. {
  290. m_WorkFlag->SetEvent();
  291. // 发送退出信号
  292. }
  293. bool Thread_Base::WaitTheThreadOnTheRun(unsigned long waittime)
  294. {
  295. LOG("Waiting for thread " << m_strName << " to run");
  296. if (m_Base_Thread == 0) {
  297. ERR_LOG("Thread " << m_strName << " does not exist");
  298. return false;
  299. }
  300. return m_RunFlag->Wait(waittime);
  301. }
  302. // 虚函数默认实现
  303. bool Thread_Base::Exec() { return true; }
  304. bool Thread_Base::OnStartThread() { return true; }
  305. bool Thread_Base::OnEndThread() { return true; }
  306. bool Thread_Base::RegistThread() { return true; }
  307. bool Thread_Base::UnRegistThread() { return true; }
  308. pthread_t Thread_Base::GetTID() { return m_Base_Thread; }
  309. void Thread_Base::SetName(const char* pszThreadName) { if (pszThreadName) m_strName = pszThreadName; }
  310. std::shared_ptr<LinuxEvent> Thread_Base::GetWorkEvt()
  311. {
  312. return m_WorkFlag;
  313. }
  314. std::shared_ptr<LinuxEvent> Thread_Base::GetExitEvt()
  315. {
  316. return m_ExitFlag;
  317. }
  318. void Thread_Base::ThreadExitProcedure() {
  319. m_ThreadRunning = false;
  320. m_ExitFlag->SetEvent();
  321. OnEndThread();
  322. UnRegistThread();
  323. SetThreadOnTheRun(false);
  324. }
  325. // Work_Thread 实现
  326. Work_Thread::Work_Thread() {
  327. m_PauseEvt = LinuxEvent::CreateEvent(LinuxEvent::MANUAL_RESET, false);
  328. m_ResumeEvt = LinuxEvent::CreateEvent(LinuxEvent::MANUAL_RESET, false);
  329. m_SleepingEvt = LinuxEvent::CreateEvent(LinuxEvent::MANUAL_RESET, false);
  330. m_WorkingEvt = LinuxEvent::CreateEvent(LinuxEvent::MANUAL_RESET, false);
  331. if (!m_PauseEvt || !m_SleepingEvt || !m_ResumeEvt || !m_WorkingEvt) {
  332. throw std::runtime_error("Work_Thread event creation failed");
  333. }
  334. MsgQueue<ResDataObject>* p;
  335. p = new MsgQueue<ResDataObject>;
  336. m_pWorkQueReq = (void *)p;
  337. p = new MsgQueue<ResDataObject>;
  338. m_pWorkQueRes = (void*)p;
  339. }
  340. Work_Thread::~Work_Thread() {
  341. MsgQueue<ResDataObject>* p;
  342. p = (MsgQueue<ResDataObject> *)m_pWorkQueReq;
  343. delete p;
  344. m_pWorkQueReq = NULL;
  345. p = (MsgQueue<ResDataObject> *)m_pWorkQueRes;
  346. delete p;
  347. m_pWorkQueRes = NULL;
  348. }
  349. bool Work_Thread::PopReqDataObject(ResDataObject& obj) {
  350. return (bool)((MsgQueue<ResDataObject> *)m_pWorkQueReq)->DeQueue(obj);
  351. }
  352. bool Work_Thread::PushReqDataObject(ResDataObject& obj) {
  353. if (WaitTheThreadEndSign(0))
  354. {
  355. return false;//not possible to send it
  356. }
  357. //printf("Thread:%d,push REQ one\n", GetCurrentThreadId());
  358. ((MsgQueue<ResDataObject> *)m_pWorkQueReq)->InQueue(obj);
  359. NotifyThreadWork();
  360. return true;
  361. }
  362. bool Work_Thread::PopResDataObject(ResDataObject& obj) {
  363. //printf("Thread:%d,Pop Res one\n", GetCurrentThreadId());
  364. return (bool)((MsgQueue<ResDataObject> *)m_pWorkQueRes)->DeQueue(obj);
  365. }
  366. bool Work_Thread::PushResDataObject(ResDataObject& obj) {
  367. if (WaitTheThreadEndSign(0))
  368. {
  369. return false;//not possible to send it
  370. }
  371. //printf("Thread:%d,push Res one\n",GetCurrentThreadId());
  372. ((MsgQueue<ResDataObject> *)m_pWorkQueRes)->InQueue(obj);
  373. NotifyThreadWork();
  374. return true;
  375. }
  376. bool Work_Thread::WaitForInQue(unsigned long m_Timeout) {
  377. std::vector<std::shared_ptr<LinuxEvent>> waitEvents = {
  378. m_ExitFlag,
  379. m_PauseEvt,
  380. m_WorkFlag ,
  381. 0,
  382. 0
  383. };
  384. waitEvents[3] = ((MsgQueue<ResDataObject> *)m_pWorkQueReq)->GetNotifyHandle();
  385. waitEvents[4] = ((MsgQueue<ResDataObject> *)m_pWorkQueRes)->GetNotifyHandle();
  386. int ret = LinuxEvent::WaitForMultipleEvents(waitEvents,0);
  387. if ((ret >= WAIT_OBJECT_0) && (ret <= WAIT_OBJECT_0 + 1))
  388. {
  389. return false;
  390. }
  391. if (((MsgQueue<ResDataObject> *)m_pWorkQueReq)->WaitForInQue(0) == WAIT_OBJECT_0)
  392. {
  393. return true;
  394. }
  395. if (((MsgQueue<ResDataObject> *)m_pWorkQueRes)->WaitForInQue(0) == WAIT_OBJECT_0)
  396. {
  397. return true;
  398. }
  399. ret = LinuxEvent::WaitForMultipleEvents(waitEvents,m_Timeout);
  400. if ((ret >= WAIT_OBJECT_0 + 2) && (ret < (WAIT_OBJECT_0 + 5)))
  401. {
  402. return true;
  403. }
  404. return false;
  405. }
  406. bool Work_Thread::CheckForPause() {
  407. DWORD retEvt = 0;
  408. bool ret = false;
  409. std::vector<std::shared_ptr<LinuxEvent>> waitExp = {
  410. m_ExitFlag,
  411. m_PauseEvt
  412. };
  413. retEvt = LinuxEvent::WaitForMultipleEvents(waitExp, 0);
  414. if (retEvt == WAIT_OBJECT_0 + 1)
  415. {
  416. m_WorkingEvt->ResetEvent();
  417. m_SleepingEvt->SetEvent();
  418. ret = true;
  419. }
  420. return ret;
  421. }
  422. bool Work_Thread::WaitforResume() {
  423. bool Exret = false;
  424. std::vector<std::shared_ptr<LinuxEvent>> wait = {
  425. m_ExitFlag,
  426. m_ResumeEvt
  427. };
  428. DWORD ret = 0;
  429. ret = LinuxEvent::WaitForMultipleEvents(wait, -1);
  430. if (ret == WAIT_OBJECT_0 + 1)
  431. {
  432. //拿到Resume
  433. m_SleepingEvt->ResetEvent();
  434. m_WorkingEvt->SetEvent();
  435. Exret = true;
  436. }
  437. return Exret;
  438. }
  439. bool Work_Thread::PauseThread(unsigned long timeout) {
  440. bool ret = false;
  441. DWORD retEvt = 0;
  442. std::vector<std::shared_ptr<LinuxEvent>> wait = {
  443. m_ExitFlag,
  444. m_SleepingEvt
  445. };
  446. if (pthread_self() == GetTID()) {
  447. return false;
  448. }
  449. if (!WaitTheThreadOnTheRun(0)) {
  450. return false;
  451. }
  452. m_PauseEvt->SetEvent();
  453. //get response
  454. retEvt = LinuxEvent::WaitForMultipleEvents(wait, timeout);
  455. if (retEvt == WAIT_OBJECT_0 + 1)
  456. {
  457. ret = true;
  458. }
  459. m_PauseEvt->ResetEvent();
  460. return ret;
  461. }
  462. bool Work_Thread::ResumeThread(unsigned long timeout) {
  463. bool ret = false;
  464. DWORD retEvt = 0;
  465. std::vector<std::shared_ptr<LinuxEvent>> wait = {
  466. m_ExitFlag,
  467. m_WorkingEvt
  468. };
  469. if (pthread_self() == GetTID()) {
  470. return false;
  471. }
  472. if (!WaitTheThreadOnTheRun(0)) {
  473. return false;
  474. }
  475. m_ResumeEvt->SetEvent();
  476. //get response
  477. retEvt = LinuxEvent::WaitForMultipleEvents(wait, timeout);
  478. if (retEvt == WAIT_OBJECT_0 + 1)
  479. {
  480. ret = true;
  481. }
  482. m_ResumeEvt->ResetEvent();
  483. return ret;
  484. }
  485. // Ccos_Thread 实现
  486. Ccos_Thread::Ccos_Thread() {}
  487. Ccos_Thread::~Ccos_Thread() {}