1
0

worker.cpp 5.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285
  1. #include <functional>
  2. #include <list>
  3. #include <vector>
  4. #include <mutex>
  5. #include <atomic>
  6. #include <memory>
  7. #include <thread>
  8. #include <ev++.h>
  9. #include "log.h"
  10. #include "worker.h"
  11. #include "message.h"
  12. #include "card_base.h"
  13. #include "card.h"
  14. #include "zloop.h"
  15. #include "area.h"
  16. loop_thread::loop_thread ()
  17. {
  18. m_thread.reset(new std::thread(std::bind(&loop_thread::run,this)));
  19. }
  20. void loop_thread::run()
  21. {
  22. loop_init();
  23. log_info("worker_thread start....");
  24. while(!check_stop_flag())
  25. {
  26. try
  27. {
  28. ev::dynamic_loop::run(0);
  29. }
  30. catch(const std::exception&e)
  31. {
  32. log_error("捕获到异常:%s",e.what());
  33. }
  34. catch(...)
  35. {
  36. log_error("捕获到未知异常");
  37. }
  38. }
  39. log_info("worker_thread exit....");
  40. }
  41. void loop_thread::on_async(const std::list<task*>&task_list)
  42. {
  43. for(task*t:task_list)
  44. {
  45. do_task(*t);
  46. }
  47. }
  48. void loop_thread::loop_init()
  49. {
  50. }
  51. void loop_thread::join()
  52. {
  53. m_thread->join();
  54. }
  55. void loop_thread::stop()
  56. {
  57. async_stop();
  58. }
  59. loop_thread::~loop_thread()
  60. {
  61. }
  62. struct hash_thread
  63. {
  64. int m_num_thread=4;
  65. void set_num_thread(int num_thread)
  66. {
  67. m_num_thread=num_thread;
  68. }
  69. int hash_code(uint32_t card_id)
  70. {
  71. return card_id*2003%m_num_thread;
  72. }
  73. };
  74. hash_thread g_hash;
  75. struct worker_thread: loop_thread ,visitor<std::shared_ptr<card_location_base>>
  76. {
  77. int m_thread_id=0;
  78. int m_card_list_version=-1;
  79. std::vector<std::shared_ptr<card_location_base>> m_local_card_list;
  80. worker_thread ()
  81. {
  82. m_local_card_list.reserve(128);
  83. }
  84. void set_thread_id(int id)
  85. {
  86. m_thread_id=id;
  87. }
  88. std::unique_ptr<ev::timer> card_timer_1s;
  89. void loop_init()
  90. {
  91. card_timer_1s.reset(new ev::timer(*this));
  92. card_timer_1s->set<worker_thread,&worker_thread::on_timeout>(this);
  93. card_timer_1s->start(1,1);
  94. }
  95. void update_local_cards()
  96. {
  97. int version=card_list::instance()->version();
  98. if(m_card_list_version!=version)
  99. {
  100. m_card_list_version=version;
  101. init_local_card_list();
  102. }
  103. }
  104. void on_timeout()
  105. {
  106. update_local_cards();
  107. for(auto&c:m_local_card_list)
  108. {
  109. c->on_timer();
  110. }
  111. }
  112. bool visit(std::shared_ptr<card_location_base> c)
  113. {
  114. if(g_hash.hash_code(c->m_id)==m_thread_id) //32bit id也可以
  115. {
  116. m_local_card_list.push_back(c);
  117. }
  118. return true;
  119. }
  120. void init_local_card_list()
  121. {
  122. m_local_card_list.clear();
  123. card_list::instance()->accept(*this);
  124. log_info("update local cards,count=%d",m_local_card_list.size());
  125. }
  126. void do_task(task&t)
  127. {
  128. switch(t.m_cmd_code)
  129. {
  130. case 0x843b://tof
  131. case 0x863b://tdoa
  132. log_info("card loc message%04X",t.m_cmd_code);
  133. card_list::instance()->on_message(this,t.body<message_locinfo>(),false);
  134. t.destroy();
  135. //card_message::on_loc_message(this,t.m_param1);
  136. break;
  137. case 0x853b://tof his
  138. case 0x873b://tdoa his
  139. log_info("site history message%04X",t.m_cmd_code);
  140. card_list::instance()->on_message(this,t.body<message_locinfo>(),true);
  141. //site_message::on_sync(this,t.m_param1);
  142. t.destroy();
  143. break;
  144. case 0x804c://ctrl site message
  145. log_info("ctrl site message%04X",t.m_cmd_code);
  146. t.destroy();
  147. break;
  148. case 0x10001://区域业务类型修改
  149. {
  150. update_local_cards();
  151. for(auto&c:m_local_card_list)
  152. {
  153. c->get_area_tool()->on_change_business(c,t);
  154. }
  155. auto&mcb=t.body<message_change_business>();
  156. std::shared_ptr<area> a=area_list::instance()->get(mcb.area_id);
  157. if(a && a->sub_frozen_count()==2)
  158. {
  159. a->set_business_list(std::move(mcb.new_list));
  160. a->sub_frozen_count();
  161. }
  162. if(mcb.ref_count.fetch_sub(1)==1)
  163. {
  164. t.destroy();
  165. }
  166. break;
  167. }
  168. case 0x10002://区域业务类型修改
  169. {
  170. auto&mcb=t.body<message_change_card_display>();
  171. auto c = card_list::instance()->get(mcb.m_card_id);
  172. c->m_display=mcb.m_display;
  173. t.destroy();
  174. break;
  175. }
  176. }
  177. }
  178. };
  179. struct worker_impl:worker
  180. {
  181. std::vector<std::shared_ptr<worker_thread>> m_threads;
  182. std::atomic<int> m_init_flag{-2};
  183. virtual void stop()
  184. {
  185. int exp=0;
  186. if(!m_init_flag.compare_exchange_strong (exp,1))
  187. return;
  188. for(auto&thr:m_threads)
  189. thr->stop();
  190. for(auto&thr:m_threads)
  191. thr->join();
  192. }
  193. worker_thread& hash(uint32_t i)
  194. {
  195. return *m_threads[g_hash.hash_code(i)];
  196. }
  197. virtual int num_thread()
  198. {
  199. return std::thread::hardware_concurrency();
  200. }
  201. bool running()
  202. {
  203. return m_init_flag.load()==0;
  204. }
  205. virtual void request(task*t)
  206. {
  207. hash(t->m_hash_id).async_request(t);
  208. }
  209. virtual void broadcast(task*tk)
  210. {
  211. for(auto&thr:m_threads)
  212. thr->async_request(tk);
  213. }
  214. void init(int num_thread)
  215. {
  216. int exp=-2;
  217. if(m_init_flag.compare_exchange_strong (exp,-1))
  218. {
  219. m_threads.resize(num_thread);
  220. for(int i=0;i<num_thread;i++)
  221. {
  222. m_threads[i].reset(new worker_thread());
  223. m_threads[i]->set_thread_id(i);
  224. }
  225. m_init_flag.store(0);
  226. }
  227. while(0!=m_init_flag.load())
  228. std::this_thread::yield();
  229. }
  230. };
  231. worker_impl _worker_impl;
  232. worker*worker::instance()
  233. {
  234. int num_thread=_worker_impl.num_thread();
  235. if(!_worker_impl.running())
  236. {
  237. log_info("worker thread count=%d",num_thread);
  238. g_hash.set_num_thread(num_thread);
  239. _worker_impl.init(num_thread);
  240. }
  241. return &_worker_impl;
  242. }