worker.cpp 6.1 KB

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