| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351 |
- /**
- * 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 <future>
- #include <unordered_map>
- #include <mutex>
- #include <condition_variable>
- #include <thread>
- #include <atomic>
- #include <iostream>
- #include <string>
- // ==================== 外部依赖(来自 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<std::promise<WeightDispatchResult>> prom;
- std::shared_future<WeightDispatchResult> fut;
- std::atomic<bool> done{false};
- };
- static std::unordered_map<std::string, std::shared_ptr<PendingWeightEntry>> g_weight_results;
- static std::mutex g_weight_results_mtx;
- // 进行中任务计数
- static std::atomic<int> g_pending_count{0};
- // tb_num 去重 map
- static std::unordered_map<std::string, time_t> 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<PendingWeightEntry> entry) {
- // RAII 清理守卫:无论正常/异常退出,均保证 pending 计数递减 + 通知等待者
- struct CleanupGuard {
- std::shared_ptr<PendingWeightEntry> 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<std::mutex> 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<float>(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<int>(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<PendingWeightEntry>();
- entry->prom = std::make_shared<std::promise<WeightDispatchResult>>();
- entry->fut = entry->prom->get_future().share();
- {
- std::lock_guard<std::mutex> 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<WeightDispatchResult> fut;
- {
- std::lock_guard<std::mutex> 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<std::mutex> 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<std::mutex> 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<std::mutex> 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;
- }
|