CcosThread.cpp 14 KB

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