weight_dispatcher.cpp 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351
  1. /**
  2. * weight_dispatcher.cpp — v43 并行称重调度器实现
  3. *
  4. * fix24 v43 核心:tb_num 确定即并行,重量不再等照片。
  5. *
  6. * 调用链路(v43 enabled=1):
  7. * main.cpp: save_file() → tb_num 确定
  8. * → weight_dispatch_async(plate, tb_num, st) [启动重量线程]
  9. * → upload_thread: plate_cap_info_upload() [照片上传]
  10. * → upload_thread: weight_dispatch_wait_result() [等待重量结果]
  11. * → upload_thread: db_insert_event_atomic() [原子写入两表]
  12. * → upload_thread: mqtt_client_publish() [MQTT 发布]
  13. */
  14. #include "weight_dispatcher.h"
  15. #include "weight_scale.h"
  16. #include "database.h"
  17. #include "mqtt_client.h"
  18. #include <future>
  19. #include <unordered_map>
  20. #include <mutex>
  21. #include <condition_variable>
  22. #include <thread>
  23. #include <atomic>
  24. #include <iostream>
  25. #include <string>
  26. // ==================== 外部依赖(来自 common.h / weight_scale.h) ====================
  27. extern bool g_weight_enabled;
  28. extern bool g_mqtt_enabled;
  29. extern double g_weight_threshold_in;
  30. extern double g_weight_threshold_out;
  31. // ==================== 配置全局变量(来自 common.h) ====================
  32. extern bool g_sync_parallel_enabled; // [sync] enabled
  33. extern int g_sync_reporter_timeout_sec;
  34. extern int g_sync_weight_retry_max;
  35. extern int g_sync_weight_retry_interval_sec;
  36. extern int g_sync_tb_num_dedup_sec;
  37. extern int g_sync_weight_snapshot_max_age_sec; // fix24 v43 Bug#18
  38. // ==================== 内部状态 ====================
  39. // tb_num → 重量结果(promise/future 模式)
  40. struct PendingWeightEntry {
  41. std::shared_ptr<std::promise<WeightDispatchResult>> prom;
  42. std::shared_future<WeightDispatchResult> fut;
  43. std::atomic<bool> done{false};
  44. };
  45. static std::unordered_map<std::string, std::shared_ptr<PendingWeightEntry>> g_weight_results;
  46. static std::mutex g_weight_results_mtx;
  47. // 进行中任务计数
  48. static std::atomic<int> g_pending_count{0};
  49. // tb_num 去重 map
  50. static std::unordered_map<std::string, time_t> g_tb_num_dedup_map;
  51. static std::mutex g_tb_num_dedup_mtx;
  52. // 秤访问串行化锁(物理秤只有一个,多 worker 并发读会导致共享状态竞态)
  53. // 覆盖:weight_scale_read() + db_insert_weight_record(),保证每次读秤完整隔离
  54. static std::mutex g_scale_access_mtx;
  55. // ==================== 内部辅助函数 ====================
  56. /**
  57. * 重量任务线程入口函数
  58. */
  59. static void weight_worker_func(const std::string& plate,
  60. const std::string& tb_num,
  61. int station_type,
  62. std::shared_ptr<PendingWeightEntry> entry) {
  63. // RAII 清理守卫:无论正常/异常退出,均保证 pending 计数递减 + 通知等待者
  64. struct CleanupGuard {
  65. std::shared_ptr<PendingWeightEntry> entry;
  66. WeightDispatchResult result;
  67. bool committed = false;
  68. ~CleanupGuard() {
  69. if (!committed) {
  70. // 异常路径:设置失败结果,防止等待者永久挂起
  71. result.read_status = 2;
  72. result.upload_status = 2; // Bug#17: 终态失败,不留"待上传"
  73. result.response_msg = "worker异常退出";
  74. try { entry->prom->set_value(result); } catch (...) {}
  75. entry->done.store(true);
  76. }
  77. g_pending_count.fetch_sub(1);
  78. }
  79. } guard{entry, {}, false};
  80. WeightDispatchResult result;
  81. result.plate = plate;
  82. result.tb_num = tb_num;
  83. result.station_type = station_type;
  84. std::string in_out_str = (station_type == 1) ? "进站" : "出站";
  85. std::cout << "[v43-重量] 开始任务: " << plate << " " << in_out_str
  86. << " 联单:" << tb_num << std::endl;
  87. try {
  88. // ===== 步骤1-3:读秤 + 上传 + 写DB(串行化,保护物理秤共享状态)=====
  89. {
  90. std::lock_guard<std::mutex> scale_lock(g_scale_access_mtx);
  91. // ===== 步骤1:读取重量 =====
  92. if (!g_weight_enabled) {
  93. std::cout << "[v43-重量] 称重未启用(flagWeight=0),跳过读取" << std::endl;
  94. result.read_status = 2;
  95. result.response_msg = "称重未启用";
  96. result.upload_status = 2; // Bug#17: 读秤失败/跳过为终态,标记失败而非"待上传"
  97. } else {
  98. double threshold_kg = (station_type == 1) ? g_weight_threshold_in : g_weight_threshold_out;
  99. std::cout << "[v43-重量] 开始读取, 阈值:" << threshold_kg << "公斤" << std::endl;
  100. double weight_kg = 0.0;
  101. bool read_ok = false;
  102. // fix24 v44 P0: 三层兜底读取
  103. // ① 预检测锁定(最快,0等待)——TCP线程已提前锁定稳定重量
  104. if (!read_ok) {
  105. read_ok = weight_scale_read_prelocked(weight_kg, threshold_kg);
  106. }
  107. // ② 快照优先(Bug#18,中等)——复用采样池近期稳定读数
  108. if (!read_ok) {
  109. read_ok = weight_scale_read_snapshot(
  110. weight_kg, threshold_kg, g_sync_weight_snapshot_max_age_sec);
  111. }
  112. // ③ 阻塞等待(兜底)——已在 weight_scale_read_snapshot 内部实现
  113. if (read_ok && weight_kg > 0) {
  114. result.weight_kg = weight_kg;
  115. result.read_status = 1;
  116. std::cout << "[v43-重量] 读取成功: " << weight_kg << "公斤" << std::endl;
  117. // ===== 步骤2:上传重量到平台(带重试) =====
  118. double weight_ton = weight_kg / 1000.0;
  119. bool upload_ok = false;
  120. // 调用 weight_scale.cpp 中的公开上传接口(带内部重试)
  121. int ret = upload_weight_with_retry_public(tb_num, static_cast<float>(weight_ton));
  122. upload_ok = (ret == 0);
  123. result.upload_status = upload_ok ? 1 : 2;
  124. result.response_msg = upload_ok ? "上传成功" : "上传失败(已重试)";
  125. std::cout << "[v43-重量] " << (upload_ok ? "✅上传成功" : "❌上传失败")
  126. << " 联单:" << tb_num << " 重量:" << weight_ton << "吨" << std::endl;
  127. } else {
  128. result.read_status = 2;
  129. result.response_msg = "未获取有效重量";
  130. // fix24 v43 Bug#17: 读秤失败是终态(无重试路径会自动捞取),
  131. // upload_status 必须写 2(失败),写 0 会在页面永久显示"待上传"误导
  132. result.upload_status = 2;
  133. std::cout << "[v43-重量] ❌未读取到有效重量 联单:" << tb_num << std::endl;
  134. }
  135. }
  136. // ===== 步骤3:记录到 weight_records 表 =====
  137. db_insert_weight_record(plate, tb_num, station_type,
  138. result.weight_kg, result.upload_status, result.response_msg);
  139. } // g_scale_access_mtx 在此释放,下一个 worker 可以读秤
  140. // ===== 步骤4:通知等待者(Reporter / upload_thread) =====
  141. entry->prom->set_value(result);
  142. entry->done.store(true);
  143. guard.result = result;
  144. guard.committed = true; // 正常路径已提交,析构只减计数
  145. std::cout << "[v43-重量] 任务结束: " << plate << " 联单:" << tb_num
  146. << " 读取:" << (result.read_status == 1 ? "✅" : "❌")
  147. << " 上传:" << (result.upload_status == 1 ? "✅" : "❌") << std::endl;
  148. // ===== 步骤5:MQTT发布(fix27 Bug#30) =====
  149. if (g_mqtt_enabled) {
  150. time_t now_t = time(nullptr);
  151. struct tm* tm_info = localtime(&now_t);
  152. char time_buf[64];
  153. strftime(time_buf, sizeof(time_buf), "%Y-%m-%d %H:%M:%S", tm_info);
  154. std::string msg =
  155. "工程名: " + project_name + "\n"
  156. "车牌号: " + plate + "\n"
  157. "时间: " + std::string(time_buf) + "\n"
  158. "工地编号: " + point_number + "\n"
  159. "门号: " + throughway + "\n"
  160. "进出站: " + in_out_str + "\n"
  161. "车牌颜色: " + g_plate_color + "\n"
  162. "车型: " + g_vehicle_type + "\n"
  163. "三联单编号:" + tb_num + "\n"
  164. "一组照片(1张低位照片,1张高位照片)上传成功";
  165. if (g_weight_enabled) {
  166. if (result.read_status == 1 && result.weight_kg > 0) {
  167. double weight_ton = result.weight_kg / 1000.0;
  168. msg += "\n称重数据: " + std::to_string(weight_ton) + "吨 ("
  169. + std::to_string(static_cast<int>(result.weight_kg)) + "公斤) ["
  170. + (result.upload_status == 1 ? "上传成功" : "上传失败") + "]";
  171. } else {
  172. msg += "\n称重数据: 未读取到有效重量";
  173. }
  174. }
  175. // MQTT发布(带重试)
  176. bool mqtt_ok = false;
  177. for (int pub_att = 0; pub_att < 3 && !mqtt_ok; pub_att++) {
  178. if (mqtt_client_publish("", msg)) {
  179. mqtt_ok = true;
  180. std::cout << "[v43-MQTT] 发布成功: " << plate << " " << in_out_str << std::endl;
  181. } else if (pub_att < 2) {
  182. std::cerr << "[v43-MQTT] 发布失败,重试(" << (pub_att+1) << "/3)" << std::endl;
  183. std::this_thread::sleep_for(std::chrono::seconds(1));
  184. }
  185. }
  186. if (!mqtt_ok) {
  187. std::cerr << "[v43-MQTT] 发布最终失败: " << plate << std::endl;
  188. }
  189. }
  190. } catch (const std::exception& e) {
  191. std::cerr << "[v43-重量] ❌worker异常: " << e.what()
  192. << " 联单:" << tb_num << std::endl;
  193. // CleanupGuard 析构自动处理:设置失败结果 + 减计数
  194. } catch (...) {
  195. std::cerr << "[v43-重量] ❌worker未知异常 联单:" << tb_num << std::endl;
  196. // CleanupGuard 析构自动处理
  197. }
  198. }
  199. // ==================== 公开接口实现 ====================
  200. void weight_dispatcher_init() {
  201. std::cout << "[v43-调度器] 初始化 (enabled=" << (g_sync_parallel_enabled ? 1 : 0)
  202. << " timeout=" << g_sync_reporter_timeout_sec << "s"
  203. << " retry_max=" << g_sync_weight_retry_max
  204. << " dedup=" << g_sync_tb_num_dedup_sec << "s)" << std::endl;
  205. }
  206. void weight_dispatcher_shutdown() {
  207. // 等待所有进行中的任务完成(最多等 60s)
  208. int waited = 0;
  209. while (g_pending_count.load() > 0 && waited < 60) {
  210. std::this_thread::sleep_for(std::chrono::seconds(1));
  211. waited++;
  212. }
  213. if (g_pending_count.load() > 0) {
  214. std::cerr << "[v43-调度器] 关闭警告: 仍有 " << g_pending_count.load()
  215. << " 个任务未完成" << std::endl;
  216. } else {
  217. std::cout << "[v43-调度器] 已关闭(所有任务已完成)" << std::endl;
  218. }
  219. }
  220. void weight_dispatch_async(const std::string& plate,
  221. const std::string& tb_num,
  222. int station_type) {
  223. if (!g_sync_parallel_enabled) {
  224. std::cerr << "[v43-调度器] 并行模式未启用,忽略 dispatch" << std::endl;
  225. return;
  226. }
  227. // tb_num 去重检查
  228. if (weight_dispatch_is_tb_num_dup(tb_num)) {
  229. std::cout << "[v43-调度器] tb_num=" << tb_num << " 已处理过,跳过" << std::endl;
  230. return;
  231. }
  232. // 创建 promise/future
  233. auto entry = std::make_shared<PendingWeightEntry>();
  234. entry->prom = std::make_shared<std::promise<WeightDispatchResult>>();
  235. entry->fut = entry->prom->get_future().share();
  236. {
  237. std::lock_guard<std::mutex> lk(g_weight_results_mtx);
  238. g_weight_results[tb_num] = entry;
  239. }
  240. g_pending_count.fetch_add(1);
  241. // 启动独立重量线程(detach,通过 future 传递结果)
  242. std::thread(weight_worker_func, plate, tb_num, station_type, entry).detach();
  243. std::cout << "[v43-调度器] 已派发重量任务: " << plate << " 联单:" << tb_num
  244. << " 待处理:" << g_pending_count.load() << std::endl;
  245. }
  246. bool weight_dispatch_wait_result(const std::string& tb_num,
  247. int timeout_sec,
  248. WeightDispatchResult& result) {
  249. std::shared_future<WeightDispatchResult> fut;
  250. {
  251. std::lock_guard<std::mutex> lk(g_weight_results_mtx);
  252. auto it = g_weight_results.find(tb_num);
  253. if (it == g_weight_results.end()) {
  254. return false; // 未找到对应任务
  255. }
  256. fut = it->second->fut;
  257. }
  258. // 等待结果
  259. auto status = fut.wait_for(std::chrono::seconds(timeout_sec));
  260. if (status == std::future_status::ready) {
  261. result = fut.get();
  262. // fix27 Bug#29: 不清理 entry,支持多次调用(飞书消息+Reporter 都需要获取结果)
  263. // entry 由 weight_dispatch_cleanup() 统一清理,防止内存泄漏
  264. return true;
  265. }
  266. // 超时:也不清理,让 worker 后续完成时仍可被获取
  267. // 最终由 weight_dispatch_cleanup() 在照片上传失败/超时后清理
  268. return false;
  269. }
  270. bool weight_dispatch_is_ready(const std::string& tb_num) {
  271. std::lock_guard<std::mutex> lk(g_weight_results_mtx);
  272. auto it = g_weight_results.find(tb_num);
  273. if (it == g_weight_results.end()) return false;
  274. return it->second->done.load();
  275. }
  276. int weight_dispatch_pending_count() {
  277. return g_pending_count.load();
  278. }
  279. void weight_dispatch_cleanup(const std::string& tb_num) {
  280. std::lock_guard<std::mutex> lk(g_weight_results_mtx);
  281. auto it = g_weight_results.find(tb_num);
  282. if (it != g_weight_results.end()) {
  283. std::cout << "[v43-调度器] 清理未消费 entry: " << tb_num << std::endl;
  284. g_weight_results.erase(it);
  285. }
  286. }
  287. bool weight_dispatch_is_tb_num_dup(const std::string& tb_num) {
  288. std::lock_guard<std::mutex> lk(g_tb_num_dedup_mtx);
  289. auto it = g_tb_num_dedup_map.find(tb_num);
  290. if (it != g_tb_num_dedup_map.end()) {
  291. time_t elapsed = time(nullptr) - it->second;
  292. if (elapsed < g_sync_tb_num_dedup_sec) {
  293. return true; // 在去重窗口内
  294. }
  295. }
  296. g_tb_num_dedup_map[tb_num] = time(nullptr);
  297. // 清理过期条目(防止内存泄漏)
  298. time_t now = time(nullptr);
  299. for (auto jt = g_tb_num_dedup_map.begin(); jt != g_tb_num_dedup_map.end(); ) {
  300. if (now - jt->second > g_sync_tb_num_dedup_sec * 2) {
  301. jt = g_tb_num_dedup_map.erase(jt);
  302. } else {
  303. ++jt;
  304. }
  305. }
  306. return false;
  307. }