/** * weight_dispatcher.cpp — v43 并行称重调度器实现 * * fix24 v43 核心:tb_num 确定即并行,重量不再等照片。 * * 调用链路(v43 enabled=1): * main.cpp: save_file() → tb_num 确定 * → weight_dispatch_async(plate, tb_num, st) [启动重量线程] * → upload_thread: plate_cap_info_upload() [照片上传] * → upload_thread: weight_dispatch_wait_result() [等待重量结果] * → upload_thread: db_insert_event_atomic() [原子写入两表] * → upload_thread: mqtt_client_publish() [MQTT 发布] */ #include "weight_dispatcher.h" #include "weight_scale.h" #include "database.h" #include "mqtt_client.h" #include #include #include #include #include #include #include #include // ==================== 外部依赖(来自 common.h / weight_scale.h) ==================== extern bool g_weight_enabled; extern bool g_mqtt_enabled; extern double g_weight_threshold_in; extern double g_weight_threshold_out; // ==================== 配置全局变量(来自 common.h) ==================== extern bool g_sync_parallel_enabled; // [sync] enabled extern int g_sync_reporter_timeout_sec; extern int g_sync_weight_retry_max; extern int g_sync_weight_retry_interval_sec; extern int g_sync_tb_num_dedup_sec; extern int g_sync_weight_snapshot_max_age_sec; // fix24 v43 Bug#18 // ==================== 内部状态 ==================== // tb_num → 重量结果(promise/future 模式) struct PendingWeightEntry { std::shared_ptr> prom; std::shared_future fut; std::atomic done{false}; }; static std::unordered_map> g_weight_results; static std::mutex g_weight_results_mtx; // 进行中任务计数 static std::atomic g_pending_count{0}; // tb_num 去重 map static std::unordered_map g_tb_num_dedup_map; static std::mutex g_tb_num_dedup_mtx; // 秤访问串行化锁(物理秤只有一个,多 worker 并发读会导致共享状态竞态) // 覆盖:weight_scale_read() + db_insert_weight_record(),保证每次读秤完整隔离 static std::mutex g_scale_access_mtx; // ==================== 内部辅助函数 ==================== /** * 重量任务线程入口函数 */ static void weight_worker_func(const std::string& plate, const std::string& tb_num, int station_type, std::shared_ptr entry) { // RAII 清理守卫:无论正常/异常退出,均保证 pending 计数递减 + 通知等待者 struct CleanupGuard { std::shared_ptr entry; WeightDispatchResult result; bool committed = false; ~CleanupGuard() { if (!committed) { // 异常路径:设置失败结果,防止等待者永久挂起 result.read_status = 2; result.upload_status = 2; // Bug#17: 终态失败,不留"待上传" result.response_msg = "worker异常退出"; try { entry->prom->set_value(result); } catch (...) {} entry->done.store(true); } g_pending_count.fetch_sub(1); } } guard{entry, {}, false}; WeightDispatchResult result; result.plate = plate; result.tb_num = tb_num; result.station_type = station_type; std::string in_out_str = (station_type == 1) ? "进站" : "出站"; std::cout << "[v43-重量] 开始任务: " << plate << " " << in_out_str << " 联单:" << tb_num << std::endl; try { // ===== 步骤1-3:读秤 + 上传 + 写DB(串行化,保护物理秤共享状态)===== { std::lock_guard scale_lock(g_scale_access_mtx); // ===== 步骤1:读取重量 ===== if (!g_weight_enabled) { std::cout << "[v43-重量] 称重未启用(flagWeight=0),跳过读取" << std::endl; result.read_status = 2; result.response_msg = "称重未启用"; result.upload_status = 2; // Bug#17: 读秤失败/跳过为终态,标记失败而非"待上传" } else { double threshold_kg = (station_type == 1) ? g_weight_threshold_in : g_weight_threshold_out; std::cout << "[v43-重量] 开始读取, 阈值:" << threshold_kg << "公斤" << std::endl; double weight_kg = 0.0; bool read_ok = false; // fix24 v44 P0: 三层兜底读取 // ① 预检测锁定(最快,0等待)——TCP线程已提前锁定稳定重量 if (!read_ok) { read_ok = weight_scale_read_prelocked(weight_kg, threshold_kg); } // ② 快照优先(Bug#18,中等)——复用采样池近期稳定读数 if (!read_ok) { read_ok = weight_scale_read_snapshot( weight_kg, threshold_kg, g_sync_weight_snapshot_max_age_sec); } // ③ 阻塞等待(兜底)——已在 weight_scale_read_snapshot 内部实现 if (read_ok && weight_kg > 0) { result.weight_kg = weight_kg; result.read_status = 1; std::cout << "[v43-重量] 读取成功: " << weight_kg << "公斤" << std::endl; // ===== 步骤2:上传重量到平台(带重试) ===== double weight_ton = weight_kg / 1000.0; bool upload_ok = false; // 调用 weight_scale.cpp 中的公开上传接口(带内部重试) int ret = upload_weight_with_retry_public(tb_num, static_cast(weight_ton)); upload_ok = (ret == 0); result.upload_status = upload_ok ? 1 : 2; result.response_msg = upload_ok ? "上传成功" : "上传失败(已重试)"; std::cout << "[v43-重量] " << (upload_ok ? "✅上传成功" : "❌上传失败") << " 联单:" << tb_num << " 重量:" << weight_ton << "吨" << std::endl; } else { result.read_status = 2; result.response_msg = "未获取有效重量"; // fix24 v43 Bug#17: 读秤失败是终态(无重试路径会自动捞取), // upload_status 必须写 2(失败),写 0 会在页面永久显示"待上传"误导 result.upload_status = 2; std::cout << "[v43-重量] ❌未读取到有效重量 联单:" << tb_num << std::endl; } } // ===== 步骤3:记录到 weight_records 表 ===== db_insert_weight_record(plate, tb_num, station_type, result.weight_kg, result.upload_status, result.response_msg); } // g_scale_access_mtx 在此释放,下一个 worker 可以读秤 // ===== 步骤4:通知等待者(Reporter / upload_thread) ===== entry->prom->set_value(result); entry->done.store(true); guard.result = result; guard.committed = true; // 正常路径已提交,析构只减计数 std::cout << "[v43-重量] 任务结束: " << plate << " 联单:" << tb_num << " 读取:" << (result.read_status == 1 ? "✅" : "❌") << " 上传:" << (result.upload_status == 1 ? "✅" : "❌") << std::endl; // ===== 步骤5:MQTT发布(fix27 Bug#30) ===== if (g_mqtt_enabled) { time_t now_t = time(nullptr); struct tm* tm_info = localtime(&now_t); char time_buf[64]; strftime(time_buf, sizeof(time_buf), "%Y-%m-%d %H:%M:%S", tm_info); std::string msg = "工程名: " + project_name + "\n" "车牌号: " + plate + "\n" "时间: " + std::string(time_buf) + "\n" "工地编号: " + point_number + "\n" "门号: " + throughway + "\n" "进出站: " + in_out_str + "\n" "车牌颜色: " + g_plate_color + "\n" "车型: " + g_vehicle_type + "\n" "三联单编号:" + tb_num + "\n" "一组照片(1张低位照片,1张高位照片)上传成功"; if (g_weight_enabled) { if (result.read_status == 1 && result.weight_kg > 0) { double weight_ton = result.weight_kg / 1000.0; msg += "\n称重数据: " + std::to_string(weight_ton) + "吨 (" + std::to_string(static_cast(result.weight_kg)) + "公斤) [" + (result.upload_status == 1 ? "上传成功" : "上传失败") + "]"; } else { msg += "\n称重数据: 未读取到有效重量"; } } // MQTT发布(带重试) bool mqtt_ok = false; for (int pub_att = 0; pub_att < 3 && !mqtt_ok; pub_att++) { if (mqtt_client_publish("", msg)) { mqtt_ok = true; std::cout << "[v43-MQTT] 发布成功: " << plate << " " << in_out_str << std::endl; } else if (pub_att < 2) { std::cerr << "[v43-MQTT] 发布失败,重试(" << (pub_att+1) << "/3)" << std::endl; std::this_thread::sleep_for(std::chrono::seconds(1)); } } if (!mqtt_ok) { std::cerr << "[v43-MQTT] 发布最终失败: " << plate << std::endl; } } } catch (const std::exception& e) { std::cerr << "[v43-重量] ❌worker异常: " << e.what() << " 联单:" << tb_num << std::endl; // CleanupGuard 析构自动处理:设置失败结果 + 减计数 } catch (...) { std::cerr << "[v43-重量] ❌worker未知异常 联单:" << tb_num << std::endl; // CleanupGuard 析构自动处理 } } // ==================== 公开接口实现 ==================== void weight_dispatcher_init() { std::cout << "[v43-调度器] 初始化 (enabled=" << (g_sync_parallel_enabled ? 1 : 0) << " timeout=" << g_sync_reporter_timeout_sec << "s" << " retry_max=" << g_sync_weight_retry_max << " dedup=" << g_sync_tb_num_dedup_sec << "s)" << std::endl; } void weight_dispatcher_shutdown() { // 等待所有进行中的任务完成(最多等 60s) int waited = 0; while (g_pending_count.load() > 0 && waited < 60) { std::this_thread::sleep_for(std::chrono::seconds(1)); waited++; } if (g_pending_count.load() > 0) { std::cerr << "[v43-调度器] 关闭警告: 仍有 " << g_pending_count.load() << " 个任务未完成" << std::endl; } else { std::cout << "[v43-调度器] 已关闭(所有任务已完成)" << std::endl; } } void weight_dispatch_async(const std::string& plate, const std::string& tb_num, int station_type) { if (!g_sync_parallel_enabled) { std::cerr << "[v43-调度器] 并行模式未启用,忽略 dispatch" << std::endl; return; } // tb_num 去重检查 if (weight_dispatch_is_tb_num_dup(tb_num)) { std::cout << "[v43-调度器] tb_num=" << tb_num << " 已处理过,跳过" << std::endl; return; } // 创建 promise/future auto entry = std::make_shared(); entry->prom = std::make_shared>(); entry->fut = entry->prom->get_future().share(); { std::lock_guard lk(g_weight_results_mtx); g_weight_results[tb_num] = entry; } g_pending_count.fetch_add(1); // 启动独立重量线程(detach,通过 future 传递结果) std::thread(weight_worker_func, plate, tb_num, station_type, entry).detach(); std::cout << "[v43-调度器] 已派发重量任务: " << plate << " 联单:" << tb_num << " 待处理:" << g_pending_count.load() << std::endl; } bool weight_dispatch_wait_result(const std::string& tb_num, int timeout_sec, WeightDispatchResult& result) { std::shared_future fut; { std::lock_guard lk(g_weight_results_mtx); auto it = g_weight_results.find(tb_num); if (it == g_weight_results.end()) { return false; // 未找到对应任务 } fut = it->second->fut; } // 等待结果 auto status = fut.wait_for(std::chrono::seconds(timeout_sec)); if (status == std::future_status::ready) { result = fut.get(); // fix27 Bug#29: 不清理 entry,支持多次调用(飞书消息+Reporter 都需要获取结果) // entry 由 weight_dispatch_cleanup() 统一清理,防止内存泄漏 return true; } // 超时:也不清理,让 worker 后续完成时仍可被获取 // 最终由 weight_dispatch_cleanup() 在照片上传失败/超时后清理 return false; } bool weight_dispatch_is_ready(const std::string& tb_num) { std::lock_guard lk(g_weight_results_mtx); auto it = g_weight_results.find(tb_num); if (it == g_weight_results.end()) return false; return it->second->done.load(); } int weight_dispatch_pending_count() { return g_pending_count.load(); } void weight_dispatch_cleanup(const std::string& tb_num) { std::lock_guard lk(g_weight_results_mtx); auto it = g_weight_results.find(tb_num); if (it != g_weight_results.end()) { std::cout << "[v43-调度器] 清理未消费 entry: " << tb_num << std::endl; g_weight_results.erase(it); } } bool weight_dispatch_is_tb_num_dup(const std::string& tb_num) { std::lock_guard lk(g_tb_num_dedup_mtx); auto it = g_tb_num_dedup_map.find(tb_num); if (it != g_tb_num_dedup_map.end()) { time_t elapsed = time(nullptr) - it->second; if (elapsed < g_sync_tb_num_dedup_sec) { return true; // 在去重窗口内 } } g_tb_num_dedup_map[tb_num] = time(nullptr); // 清理过期条目(防止内存泄漏) time_t now = time(nullptr); for (auto jt = g_tb_num_dedup_map.begin(); jt != g_tb_num_dedup_map.end(); ) { if (now - jt->second > g_sync_tb_num_dedup_sec * 2) { jt = g_tb_num_dedup_map.erase(jt); } else { ++jt; } } return false; }