dispatch_communicator.cpp 8.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235
  1. //
  2. // Created by zx on 2020/7/21.
  3. //
  4. #include "dispatch_communicator.h"
  5. Dispatch_communicator::~Dispatch_communicator(){}
  6. Error_manager Dispatch_communicator::dispatch_request(message::Dispatch_request_msg& request,
  7. message::Dispatch_response_msg& result)
  8. {
  9. /*
  10. * 检查request合法性,以及模块状态
  11. */
  12. if(request.base_info().sender()!=message::eMain||request.base_info().receiver()!=message::eDispatch)
  13. return Error_manager(ERROR,MINOR_ERROR,"dispatch request invalid");
  14. if(m_response_table.find(request)==true)
  15. return Error_manager(ERROR,MAJOR_ERROR," dispatch request repeated");
  16. //设置超时,若没有设置,默认300000 五分钟
  17. int timeout=request.base_info().has_timeout_ms()?request.base_info().timeout_ms():300000;
  18. //向测量节点发送测量请求,并记录请求
  19. Error_manager code;
  20. Communication_message message;
  21. message::Base_info base_msg;
  22. base_msg.set_msg_type(message::eDispatch_request_msg);
  23. base_msg.set_sender(message::eMain);
  24. base_msg.set_receiver(message::eDispatch);
  25. base_msg.set_timeout_ms(timeout);
  26. message.reset(base_msg,request.SerializeAsString());
  27. code=encapsulate_msg(&message);
  28. if(code!=SUCCESS)
  29. return code;
  30. //循环查询请求是否被处理
  31. auto start_time=std::chrono::system_clock::now();
  32. double time=0;
  33. do{
  34. //查询到记录
  35. message::Dispatch_response_msg response;
  36. ///查询是否存在,并且删除该记录,
  37. if(m_response_table.find(request,response))
  38. {
  39. //判断是否接收到回应,若回应信息被赋值则证明有回应
  40. if (response.has_base_info() && response.has_command_id())
  41. {
  42. message::Base_info response_base = response.base_info();
  43. //检查类型是否匹配
  44. if (response_base.msg_type() != message::eDispatch_response_msg) {
  45. return Error_manager(ERROR, CRITICAL_ERROR,
  46. "dispatch response basemsg type error");
  47. }
  48. //检查基本信息是否匹配
  49. if (response_base.sender() != message::eDispatch ||
  50. response_base.receiver() != message::eMain ||
  51. response.command_id() != request.command_id()) {
  52. return Error_manager(ERROR, MAJOR_ERROR,
  53. "dispatch response basemsg info error");
  54. }
  55. result = response;
  56. m_response_table.erase(request);
  57. return SUCCESS;
  58. }
  59. }
  60. else
  61. {
  62. //未查询到记录,任务已经被提前取消,记录被删除
  63. return Error_manager(FAILED,MINOR_ERROR,"dispatch request canceled");
  64. }
  65. auto end_time=std::chrono::system_clock::now();
  66. auto duration = std::chrono::duration_cast<std::chrono::microseconds>(end_time - start_time);
  67. time=1000.0*double(duration.count()) * std::chrono::microseconds::period::num / std::chrono::microseconds::period::den;
  68. std::this_thread::yield();
  69. usleep(1000);
  70. }while(time<double(timeout));
  71. m_response_table.erase(request);
  72. return Error_manager(RESPONSE_TIMEOUT,MINOR_ERROR,"dispatch request timeout");
  73. }
  74. /*
  75. * 提前取消请求
  76. */
  77. Error_manager Dispatch_communicator::cancel_request(message::Dispatch_request_msg& request)
  78. {
  79. message::Dispatch_response_msg t_response;
  80. if(m_response_table.find(request,t_response))
  81. {
  82. m_response_table.erase(request);
  83. }
  84. return SUCCESS;
  85. }
  86. Error_manager Dispatch_communicator::check_statu()
  87. {
  88. std::chrono::system_clock::time_point time_now=std::chrono::system_clock::now();
  89. auto durantion=time_now-m_dispatch_statu_time;
  90. if(m_dispatch_statu_msg.has_base_info()== false
  91. || durantion>std::chrono::seconds(5))
  92. {
  93. return Error_manager(ERROR,MINOR_ERROR,"调度节点通讯断开");
  94. }
  95. if(m_dispatch_statu_msg.dispatch_manager_status()==message::E_DISPATCH_MANAGER_FAULT)
  96. {
  97. return Error_manager(ERROR,MINOR_ERROR,"调度节点故障");
  98. }
  99. if(m_dispatch_statu_msg.dispatch_manager_status()==message::E_DISPATCH_MANAGER_UNKNOW)
  100. {
  101. return Error_manager(ERROR,MINOR_ERROR,"调度节点状态未知");
  102. }
  103. return SUCCESS;
  104. }
  105. Dispatch_communicator::Dispatch_communicator()
  106. {
  107. }
  108. Error_manager Dispatch_communicator::encapsulate_msg(Communication_message* message)
  109. {
  110. Error_manager code;
  111. //记录请求
  112. switch (message->get_message_type())
  113. {
  114. case Communication_message::eDispatch_request_msg:
  115. {
  116. message::Dispatch_request_msg request;
  117. if(false==request.ParseFromString(message->get_message_buf()))
  118. {
  119. code=Error_manager(ERROR,CRITICAL_ERROR,"request message parse failed");
  120. }
  121. m_response_table[request]=message::Dispatch_response_msg();
  122. //发送请求
  123. code= Communication_socket_base::encapsulate_msg(message);
  124. if(code!=SUCCESS)
  125. {
  126. m_response_table.erase(request);
  127. }
  128. break;
  129. }
  130. default:
  131. code= Error_manager(FAILED,CRITICAL_ERROR," measure发送任务类型不存在");
  132. break;
  133. }
  134. return code;
  135. }
  136. Error_manager Dispatch_communicator::execute_msg(Communication_message* p_msg)
  137. {
  138. if(p_msg== nullptr)
  139. return Error_manager(POINTER_IS_NULL,CRITICAL_ERROR,"dispatch response msg pointer is null");
  140. //测量response消息
  141. switch (p_msg->get_message_type())
  142. {
  143. ///测量结果反馈消息
  144. case Communication_message::eDispatch_response_msg:
  145. {
  146. message::Dispatch_response_msg response;
  147. response.ParseFromString(p_msg->get_message_buf());
  148. message::Dispatch_request_msg request=create_request_by_response(response);
  149. ///查询请求表是否存在,并且更新
  150. if(m_response_table.find_update(request,response)==false)
  151. {
  152. return Error_manager(ERROR,NEGLIGIBLE_ERROR,"dispatch response without request");
  153. }
  154. break;
  155. }
  156. ///测量系统状态
  157. case Communication_message::eDispatch_status_msg:
  158. {
  159. if(m_dispatch_statu_msg.ParseFromString(p_msg->get_message_buf())==false)
  160. {
  161. return Error_manager(ERROR,CRITICAL_ERROR,"调度模块状态消息解析失败");
  162. }
  163. m_dispatch_statu_time=std::chrono::system_clock::now();
  164. break;
  165. }
  166. }
  167. return SUCCESS;
  168. }
  169. /*
  170. * 检测消息是否可被处理
  171. */
  172. Error_manager Dispatch_communicator::check_msg(Communication_message* p_msg)
  173. {
  174. //通过 p_msg->get_message_type() 和 p_msg->get_receiver() 判断这条消息是不是给我的.
  175. //子类重载时, 增加自己模块的判断逻辑, 以后再写.
  176. if ( (p_msg->get_message_type() == Communication_message::Message_type::eDispatch_response_msg
  177. ||p_msg->get_message_type() == Communication_message::Message_type::eDispatch_status_msg)
  178. && p_msg->get_receiver() == Communication_message::Communicator::eMain )
  179. {
  180. return Error_code::SUCCESS;
  181. }
  182. else
  183. {
  184. //认为接受人
  185. return Error_code::INVALID_MESSAGE;
  186. }
  187. }
  188. /*
  189. * 心跳发送函数,重载
  190. */
  191. Error_manager Dispatch_communicator::encapsulate_send_data()
  192. {
  193. return SUCCESS;
  194. }
  195. //检查消息是否可以被解析, 需要重载
  196. Error_manager Dispatch_communicator::check_executer(Communication_message* p_msg)
  197. {
  198. return SUCCESS;
  199. }
  200. /*
  201. * 根据接收到的response,生成对应的 request
  202. */
  203. message::Dispatch_request_msg Dispatch_communicator::create_request_by_response(message::Dispatch_response_msg& response)
  204. {
  205. message::Dispatch_request_msg request;
  206. message::Base_info baseInfo;
  207. baseInfo.set_msg_type(message::eDispatch_request_msg);
  208. baseInfo.set_sender(response.base_info().receiver());
  209. baseInfo.set_receiver(response.base_info().sender());
  210. request.mutable_base_info()->CopyFrom(baseInfo);
  211. request.set_command_id(response.command_id());
  212. return request;
  213. }