#include "mqtt_emms2.h"
#include "../ems/jsonct.h"
#include "../common.h"
#include "../proto_forward.h"
#include "define.h"
#include "lz4.h"
#include "../log_module/log_module.h"
#include "../cmdclient.h"
#include "sqlite3.h"
#include "../snapshoot/snapshoot.h"
#include <stdatomic.h>
#include <stdbool.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <sys/types.h>
#include <sys/ipc.h>
#include <sys/msg.h>
#include "collector-api.h"
#include "../alarm_control/alarm_control.h"
#ifdef EN_ACCIDENT
#include "../accident/accident.h"
#endif
#include "../frpc_proxy/frpc_proxy.h"
#include <dirent.h>
#include <fnmatch.h>
#include <unistd.h>
#include "../opcua/ua_client.h"
#include "../cc_lc_com/pub_fun.h"
#include "../ota_upgrade/ota_upgrade.h"
#include "../rate/rate.h"
#include "../multilanguage.h"
#include "../opt_log.h"
#include "../business_log/business_log.h"

mqtt_emms2_var_t* mqtt_emms2_var_list[MQTT_EMMS2_MAX_THREAD_NUM] = {0};
int               mqtt_emms2_var_list_num = 0;
extern int north_control;

mqtt_emms2_var_t* get_mqtt_emms2_var_by_id(char* id)
{
    for(int i = 0; i < mqtt_emms2_var_list_num; i++)
    {
        if (strcmp(mqtt_emms2_var_list[i]->station_id, id) == 0)
        {
            return mqtt_emms2_var_list[i];
        }
    }
    return NULL;
}

//
#define DATA_INFO_TO_STRING(tag, l_buf)                                    \
    do {                                                                   \
        if ((tag).data_type == 0) /* int */                                \
            lbuff_sprintf(&(l_buf), "\"%s\":%ld,", (tag).name, (tag).read_cache.to_int);  \
        else                                                               \
            lbuff_sprintf(&(l_buf), "\"%s\":%.2f,", (tag).name, (tag).read_cache.to_float); \
    } while (0)

#define TAG_INFO_TO_STRING(tag, l_buf)      \
    do                                                  \
    {                                                   \
        if (tag->data_type == TYPE_TAG_INT)             \
            lbuff_sprintf(&l_buf, "\"%s\":%ld,",tag->name,tag->read_cache.to_int);       \
        else                                            \
            lbuff_sprintf(&l_buf,"\"%s\":%.2f,",tag->name,tag->read_cache.to_float);         \
    } while (0);

#define TAG_VALUE_TO_STRING_ARRAY(tag, l_buf)      \
    do                                                  \
    {                                                   \
        if (tag->data_type == TYPE_TAG_INT)             \
            lbuff_sprintf(&l_buf, "%ld,",tag->read_cache.to_int);       \
        else                                            \
            lbuff_sprintf(&l_buf, "%.2f,",tag->read_cache.to_float);         \
    } while (0);

    #define EMS_SLAVE_TAG_VALUE_TO_STRING_ARRAY(tag, l_buf)      \
    do {                                                     \
        if ((tag).data_type == 0) /* int */                  \
            lbuff_sprintf(&(l_buf), "%ld,", (tag).read_cache.to_int);       \
        else                                                 \
            lbuff_sprintf(&(l_buf), "%.2f,", (tag).read_cache.to_float);    \
    } while (0)

static const int backoff_intervals[] = {1, 10, 30, 60, 600, 1800, 3600};

extern int cabinet_info_tab_count_g;
extern cabinet_info_tab_t cabinet_info_tab[];
extern cabinet_info_t cabinet_info;
pthread_mutex_t cab_info_mutex = PTHREAD_MUTEX_INITIALIZER;


static int report_all_ems_data_by_type(mqtt_emms2_var_t *var,const char *type,EMMS2_TAG_TYPE tag_type,int seq);
static int report_dev_data_by_type_and_name(mqtt_emms2_var_t *var,const char *type,char** dev_name,int dev_num,EMMS2_TAG_TYPE tag_type,int seq);
static int report_all_alarm(mqtt_emms2_var_t *var,int seq);
// static char *get_bms_no_cell_no(char *cellno);
static void report_all_data(mqtt_emms2_var_t *var);
static void mqtt_emms2_devInfo(mqtt_emms2_var_t *var);
static void mqtt_emms2_PolicDataGet(mqtt_emms2_var_t *var);

static alarm_report_t* pop_alarm_from_report_queue(mqtt_emms2_var_t* var);
static int get_alarm_report_queue_size(mqtt_emms2_var_t* var);
int check_and_upload_templates_and_alarm(mqtt_emms2_var_t* var);
/*******************************************Base64************************************************/
// Base64 解码表
// static const char base64_table[] = "ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789+/";

// 将 Base64 字符转换为对应的索引值
static int base64_char_to_index(char c)
{
    if (c >= 'A' && c <= 'Z') return c - 'A';
    if (c >= 'a' && c <= 'z') return c - 'a' + 26;
    if (c >= '0' && c <= '9') return c - '0' + 52;
    if (c == '+') return 62;
    if (c == '/') return 63;
    return -1; // 无效的 Base64 字符
}

// Base64 解码函数
int base64_decode(const char *base64_str, char *decoded_str) {
    int len = strlen(base64_str);
    int padding = 0; // Base64 字符串的填充字符数

    // 检查 Base64 字符串是否以 '=' 结尾，以确定填充字符数
    if (len > 0 && base64_str[len - 1] == '=') padding++;
    if (len > 1 && base64_str[len - 2] == '=') padding++;

    // 计算解码后的字符串长度
    int decoded_len = (len / 4) * 3 - padding;

    // 解码 Base64 字符串
    for (int i = 0, j = 0; i < len; i += 4, j += 3) {
        int a = base64_char_to_index(base64_str[i]);
        int b = base64_char_to_index(base64_str[i + 1]);
        int c = base64_char_to_index(base64_str[i + 2]);
        int d = base64_char_to_index(base64_str[i + 3]);

        if (a == -1 || b == -1 || c == -1 || d == -1) {
            // 包含无效的 Base64 字符
            return -1;
        }

        decoded_str[j] = (a << 2) | (b >> 4);
        if (j + 1 < decoded_len) decoded_str[j + 1] = ((b & 0x0f) << 4) | (c >> 2);
        if (j + 2 < decoded_len) decoded_str[j + 2] = ((c & 0x03) << 6) | d;
    }

    // 添加字符串结束符
    decoded_str[decoded_len] = '\0';

    return decoded_len;
}

/************************************************************************************************/


static int load_mqtt_emms2_info(mqtt_emms2_var_t *pcfg, char* station_id)
{
    if (pcfg == NULL || station_id == NULL)
        return -1;

    int re = 0;
    cJSON *config = NULL;
    char* f_data = NULL;
    char* file_path = NULL;
    
    if (strncmp(station_id, MQTT_EMMS2_STATION_ID, strlen(MQTT_EMMS2_STATION_ID)) == 0)
    {
        file_path = MQTT_EMMS2_CONFIG_PATH;
    }
    else if (strncmp(station_id, MQTT_STACTRL_STATION_ID, strlen(MQTT_STACTRL_STATION_ID)) == 0)
    {
        file_path = MQTT_STACTRL_CONFIG_PATH;
    }
    else
    {
        ems_syslog(LOG_ERR, "station_id error, station_id:%s", station_id);
        re = -1;
        goto emms2_exit;
    }
     f_data = read_file_data(file_path);
    if (f_data == NULL)
    {
        ems_syslog(LOG_CRIT, "[云/站控]:加载mqtt配置文件失败,文件名:%s", file_path);
        BUSINESS_LOG(BLOG_ERR, NORTH_ID_0054, station_id, "[北向]:[云/站控]:[%s] 读取mqtt服务器配置文件失败, 文件名：[%s]", station_id, file_path);
        re = -1;
        goto emms2_exit;
    }
    else
    {
        BUSINESS_LOG(LOG_NOTICE, NORTH_ID_0054, station_id, "[北向]:[云/站控]:[%s] 读取mqtt服务器配置文件成功, 文件名：[%s]", station_id, file_path);
    }

    config = json_parse_string_with_comments(f_data);
    if (!config)
    {
        BUSINESS_LOG(BLOG_ERR, NORTH_ID_0055, station_id, "[北向]:[云/站控]:[%s] 解析mqtt服务器配置文件失败, 文件名：[%s]", station_id, file_path);
        ems_syslog(LOG_CRIT, "[云/站控]:解析mqtt配置文件失败,文件内容异常,文件名:%s", file_path);
        re = -2;
        goto emms2_exit;
    }
    else
    {
        BUSINESS_LOG(LOG_NOTICE, NORTH_ID_0055, station_id, "[北向]:[云/站控]:[%s] 解析mqtt服务器配置文件成功, 文件名：[%s]", station_id, file_path);
    }

    cJSON* obj = cJSON_GetObjectItemCaseSensitive(config, "host");
    if (obj!=NULL)
        pcfg->base_param.addr = strdup(obj->valuestring);
    obj = cJSON_GetObjectItemCaseSensitive(config, "port");
    if (obj!=NULL)
        pcfg->base_param.port = obj->valueint;
    obj = cJSON_GetObjectItemCaseSensitive(config, "user");
    if (obj!=NULL)
        pcfg->base_param.user = strdup(obj->valuestring);
    obj = cJSON_GetObjectItemCaseSensitive(config, "pass");
    if (obj!=NULL)
        pcfg->base_param.pass = strdup(obj->valuestring);
    obj = cJSON_GetObjectItemCaseSensitive(config, "reportInterval");
    if (obj!=NULL)
        pcfg->base_param.reportInterval = obj->valueint;

    obj = cJSON_GetObjectItemCaseSensitive(config, "changeReportInterval");
    if (obj!=NULL)
        pcfg->base_param.changeReportInterval = obj->valueint;
    obj = cJSON_GetObjectItemCaseSensitive(config, "saveDisconnectData");
    if (obj!=NULL)
        pcfg->base_param.saveDisconnectData = obj->valueint;

    obj = cJSON_GetObjectItemCaseSensitive(config, "enable_tls");
    if (obj != NULL)
        pcfg->base_param.tls = obj->valueint;

    obj = cJSON_GetObjectItemCaseSensitive(config, "disable_lz4");
    if (obj != NULL)
        pcfg->base_param.disable_lz4 = obj->valueint;

    obj = cJSON_GetObjectItemCaseSensitive(config, "change_deadband");
    if (obj != NULL)
        pcfg->base_param.change_deadband = obj->valueint;
    else
        pcfg->base_param.change_deadband = 5;

    obj = cJSON_GetObjectItemCaseSensitive(config, "uploadConfigUrl");
    if (obj != NULL)
    {
        pcfg->base_param.config_url = strdup(obj->valuestring);
    }
    else
    {
        char tmp[256] = {0};
        snprintf(tmp, sizeof(tmp), "https://%s/api/emmsv2/manager/config/upload", pcfg->base_param.addr);
        pcfg->base_param.config_url = strdup(tmp);
    }
    pcfg->base_param.operation = get_opeartion_status(station_id);

emms2_exit:
    if (config != NULL)
        cJSON_Delete(config);
    if (f_data != NULL)
        free(f_data);
    return re;
}

/**************************************************************************************/

void init_offline_db_table(mqtt_emms2_var_t* var, char* station_id)
{
    char* zErrMsg   = NULL;
    char  sql[256]  = {0};
    char dbname[256] = {0};
    var->offline_db = (offline_db_s*)malloc(sizeof(offline_db_s));
    if (var->offline_db == NULL)
    {
        ems_syslog(LOG_ERR, "malloc offline_db error");
        goto INIT_ERR;
    }
    snprintf(dbname, sizeof(dbname), "%s_%s.db", EMMS2_HISTORY_DB_PRE, station_id);

    int ret = sqlite3_open(dbname, &var->offline_db->db);
    if (ret != SQLITE_OK)
    {
        ems_syslog(LOG_ERR, "sqlite3_open error:%s", sqlite3_errmsg(var->offline_db->db));
        goto INIT_ERR;
    }

    snprintf(sql, sizeof(sql), " CREATE TABLE IF NOT EXISTS " EMMS2_TAGBLE_NAME " (ID INTEGER PRIMARY KEY, topic TEXT NOT NULL, payload TEXT NOT NULL);");
    ems_syslog(LOG_DEBUG, "Creat Table :" EMMS2_TAGBLE_NAME);
     
    sqlite3_exec(var->offline_db->db, sql, NULL, NULL, &zErrMsg);
    if(zErrMsg)
    {
        ems_syslog(LOG_ERR, "Creat Table: %s, errMsg:%s", station_id, zErrMsg);
        sqlite3_free(zErrMsg);
        goto INIT_ERR;
    }

    snprintf(sql,
             sizeof(sql),
             "CREATE TRIGGER IF NOT EXISTS  prevent_insert "
             "BEFORE INSERT ON " EMMS2_TAGBLE_NAME
             " FOR EACH ROW "
             "BEGIN "
             "    SELECT RAISE(ABORT, 'Maximum row limit reached') "
             "    WHERE (SELECT COUNT(*) FROM " EMMS2_TAGBLE_NAME
             " ) >= %d; "
             "END;",
             MQTT_OFFLINE_DB_MAX_COUNT);

    // 执行SQL语句
    sqlite3_exec(var->offline_db->db, sql, NULL, NULL, &zErrMsg);
    if (zErrMsg)
    {
        ems_syslog(LOG_ERR, "Creat TRIGGER: %s, errMsg:%s", sql, zErrMsg);
        sqlite3_free(zErrMsg);
        goto INIT_ERR;
    }

    pthread_mutex_init(&var->offline_db->history_mutex, NULL);
    pthread_cond_init(&var->offline_db->history_cond, NULL);
    return;
INIT_ERR:
    if (var->offline_db != NULL)
    {
        sqlite3_close(var->offline_db->db);
        free(var->offline_db);
        var->offline_db = NULL;
    }
    return;
}

static void init_mqtt_alarm_report_list(mqtt_emms2_var_t* var)
{
    var->alarm_report_list = (alarm_report_list_s*)malloc(sizeof(alarm_report_list_s));
    if (var->alarm_report_list == NULL)
    {
        ems_syslog(LOG_ERR, "malloc mqtt_alarm_report_list error");
        return;
    }
    INIT_LIST_HEAD(&var->alarm_report_list->alarm_report_queue);
    pthread_mutex_init(&var->alarm_report_list->alarm_report_queue_mutex, NULL);
    var->alarm_report_list->alarm_report_queue_max_size = 1000;
    var->alarm_report_list->alarm_report_queue_size     = 0;
}

static void init_pub_dev_list(mqtt_emms2_var_t* var)
{
    var->pub_dev = malloc(sizeof(pub_dev_t));
    memset(var->pub_dev, 0, sizeof(pub_dev_t));

    int dev_count = 0;
    proto_forward_t* proto_f = get_proto_forward_var();
    for (int i = 0; i < proto_f->channels_size; i++)
    {
        dev_count += proto_f->channels[i]->devs_size;
    }

    var->pub_dev->pub_dev_list_num = 0;
    var->pub_dev->pub_dev_list     = malloc(sizeof(pub_dev_list_t) * dev_count);

    cJSON* root = NULL;
    char*  file = read_file_data(MQTT_PUB_DEV_LIST_CONFIG);
    if (file == NULL)
    {
        ems_syslog(LOG_ERR, "read_file_data %s error", MQTT_PUB_DEV_LIST_CONFIG);
        goto INIT_ERR;
    }
    root = cJSON_Parse(file);
    if (root == NULL)
    {
        ems_syslog(LOG_ERR, "cJSON_Parse %s error", MQTT_PUB_DEV_LIST_CONFIG);
        goto INIT_ERR;
    }
    cJSON* pdl = cJSON_GetObjectItem(root, "pub_dev_list");
    if (pdl == NULL || cJSON_IsArray(pdl) == false)
    {
        ems_syslog(LOG_ERR, "cJSON_GetObjectItem %s error", MQTT_PUB_DEV_LIST_CONFIG);
        goto INIT_ERR;
    }
    int size = cJSON_GetArraySize(pdl);
    if (size <= 0)
    {
        ems_syslog(LOG_ERR, "pub_dev_list size <= 0");
        goto INIT_ERR;
    }
    
    for (int i = 0; i < proto_f->channels_size; i++)
    {
        for (int j = 0; j < proto_f->channels[i]->devs_size; j++)
        {
            device_t* dev = proto_f->channels[i]->devs[j];
            for (int k = 0; k < size; k++)
            {
                cJSON* item = cJSON_GetArrayItem(pdl, k);
                if (item == NULL)
                    continue;
                cJSON* type = cJSON_GetObjectItem(item, "type");
                if (type == NULL || cJSON_IsString(type) == false || strcmp(dev->dev_type, type->valuestring) != 0)
                    continue;
                cJSON* tags = cJSON_GetObjectItem(item, "tags");
                if (tags == NULL || cJSON_IsArray(tags) == false)
                    continue;
                int tags_size = cJSON_GetArraySize(tags);
                if (tags_size <= 0)
                    continue;

                pub_dev_list_t* p = &var->pub_dev->pub_dev_list[var->pub_dev->pub_dev_list_num++];
                p->device         = dev;
                p->pub_tag_num    = 0;
                p->tag_pub_list   = malloc(sizeof(tag_t*) * tags_size);

                for (int m = 0; m < tags_size; m++)
                {
                    cJSON* tag = cJSON_GetArrayItem(tags, m);
                    if (tag == NULL || cJSON_IsString(tag) == false)
                        continue;
                    tag_t* t = find_tag_by_dev_and_node(proto_f, dev->no, tag->valuestring);
                    if (t != NULL)
                    {
                        p->tag_pub_list[p->pub_tag_num++] = t;
                    }
                }
            }
        }
    }
    if (root) cJSON_Delete(root);
    if (file) free(file);
    return;
INIT_ERR:
    if (root) cJSON_Delete(root);
    if (file) free(file);
    var->pub_dev->pub_dev_list_num = 0;
    if (var->pub_dev->pub_dev_list) free(var->pub_dev->pub_dev_list);
}

static int sqlite3_exec_command(sqlite3* db, char* sql)
{
    char *zErrMsg = NULL;
    char** azResult = NULL;

    int rc = sqlite3_get_table(db, sql, &azResult, NULL, NULL, &zErrMsg);
    if (rc != SQLITE_OK)
    {
        if (rc == SQLITE_CONSTRAINT)
            ems_syslog(LOG_NOTICE, "sqlite3_exec_command error: %s", zErrMsg);
        else
            ems_syslog(LOG_ERR, "sqlite3_exec_command error: %d, %s", rc, zErrMsg);
    }
    sqlite3_free_table(azResult);
    sqlite3_free(zErrMsg);

    return rc != SQLITE_OK ? -1 : 0;
}

int mqtt_publish_message(const mqtt_emms2_var_t* var, char* pub_topic, char* data, int data_len)
{
    int ret = -99;
    if (var->base_param.disable_lz4)
    {
        ret = mqtt_session_publish(var->session, pub_topic, data, data_len);
    }
    else
    {
        int   src_size        = data_len + 1;
        int   max_dst_size    = LZ4_compressBound(src_size);
        char* compressed_data = calloc(1, max_dst_size);

        const int compressed_data_size = LZ4_compress_default(data, compressed_data, src_size, max_dst_size);
        if (compressed_data_size <= 0)
        {
            BUSINESS_LOG(BLOG_ERR, NORTH_ID_0052, var->station_id, "[北向]:[云/站控]:[%s] LZ4_compress err", var->station_id);
            ems_syslog(LOG_ERR, "LZ4_compress err");
            free(compressed_data);
            return -1;
        }
        else
        {
            BUSINESS_LOG(LOG_NOTICE, NORTH_ID_0052, var->station_id, "[北向]:[云/站控]:[%s] LZ4_compress success", var->station_id);
        }
        // printf("We successfully compressed some data! Ratio: %.2f\n",
        //     (float) compressed_data_size/src_size);
        char topic_lz4[256] = {0};
        sprintf(topic_lz4, "%s/lz4/%d", pub_topic, data_len + 1);
        ret = mqtt_session_publish(var->session, topic_lz4, compressed_data, compressed_data_size);
        free(compressed_data);
    }
    ems_syslog(LOG_NOTICE, "%s[ret=%d]---->%s\n", pub_topic, ret, data);

    if(ret == 0)
    {
        BUSINESS_LOG(LOG_NOTICE, NORTH_ID_0053, var->station_id, "[北向]:[云/站控]:[%s] mqtt_publish_message success", var->station_id);
    }
    else
    {
        BUSINESS_LOG(BLOG_ERR, NORTH_ID_0053, var->station_id, "[北向]:[云/站控]:[%s] mqtt_publish_message err", var->station_id);
    }

    return ret;
}


void package_based_field(mqtt_emms2_var_t *var,char *dist,char *functionId,int seq,char *ExData)//seq=0会自动产生累加值 ExData会直接添加到大括号内部
{
    int seq_t = 0;
    if (seq == 0) 
        seq_t = var->seq++;
    else
        seq_t = seq;

    sprintf(dist,"{\"funcId\":\"%s\",\"lcSN\":\"%s\",\"seq\": %d,\"time\":%ld,%s}",functionId,var->topic_sn,seq_t,time(NULL),ExData);
}

void write_response(mqtt_emms2_var_t *var,int seq,RESPONSE_RESULT result,char* function) 
{
    char data_all[1024] = {0};
    char data_info[64] = {0};
    sprintf(data_info,"\"result\":%d",result);

    package_based_field(var,data_all,function,seq,data_info);

    char pub_topic[512] = {0};
    sprintf(pub_topic,"emms2/"TYPE_SET_RESP"/%s/%s",var->topic_sn,function);
    ems_syslog(LOG_INFO,"pub topic [%s] data:[%s]",pub_topic,data_all);
    if(mqtt_publish_message(var,pub_topic,data_all,strlen(data_all))!=0)
        ems_syslog(LOG_WARNING,"send data error!!");
}

static void snapshoot_response(mqtt_emms2_var_t *var,int seq,RESPONSE_RESULT result,char* function,long long int size,char *md5) 
{
    char data_all[1024] = {0};
    char data_info[64] = {0};
    sprintf(data_info,"\"result\":%d,\"info\":{\"size\":%lld,\"md5\":\"%s\"}",result,size,md5);

    package_based_field(var,data_all,function,seq,data_info);

    char pub_topic[512] = {0};
    sprintf(pub_topic,"emms2/"TYPE_SET_RESP"/%s/%s",var->topic_sn,function);
    ems_syslog(LOG_INFO,"pub topic [%s] data:[%s]",pub_topic,data_all);
    if(mqtt_publish_message(var,pub_topic,data_all,strlen(data_all))!=0)
        ems_syslog(LOG_WARNING,"send data error!!");
}

static int set_telecontrol(char *name,double value)
{
    int ret = -1;
    proto_forward_t *var_proto = get_proto_forward_var();     
    tag_t *tag_find = find_tag_by_name(var_proto,name);
    char addtmp[256] = {0};
    if ((tag_find)&&(tag_find->point_type == PROPERTY_REALDATA))//没办法判断是不是读写数据 
    {
        tag_value_t tag_val={0};
        tag_val.type = TYPE_TAG_FLOAT;
        tag_val.value.to_float = value;
        ret = write_data_to_tag_by_p(tag_find,&tag_val,1.0);
        if (ret < 0)
        {
            goto ERR;
        }
        
        snprintf(addtmp, sizeof(addtmp), ":%s/%f", tag_find->name, value);
        server_add_new_log(LOG_CODE_EMMS2_WRITE_POINT,((device_t*)tag_find->dev_ptr)->no,LOG_MODE_TRIG,1,addtmp);
        return ret;
    }
ERR:
    snprintf(addtmp, sizeof(addtmp), ":%s/%f", name, value);
    server_add_new_log(LOG_CODE_EMMS2_WRITE_POINT,"EMS",LOG_MODE_TRIG,2,addtmp);
    return -1;
}

static void save_forbidden_alarm_codes(cJSON *alarmCodes)
{
    int size = cJSON_GetArraySize(alarmCodes);
    int* alarm_codes = (int*)malloc(sizeof(int) * size);
    for (size_t i = 0; i < size; i++)
    {
        cJSON* code = cJSON_GetArrayItem(alarmCodes, i);
        if (code != NULL)
        {
            alarm_codes[i] = code->valueint;
        }
    }
    add_forbidden_alarm_code(alarm_codes, size);
    free(alarm_codes);
}

off_t check_file_size(char *pwd)
{
    struct stat fileStat;  
    if (stat(pwd, &fileStat) != 0) {  
        perror("stat");  
        return 1;  
    }  
    // 文件大小以字节为单位  
    return fileStat.st_size; 
}

void queue_smg_free(QUEUE_MSG_S *queue_msg)
{
    if (queue_msg == NULL) return;

    if (queue_msg->host != NULL)
    {
        free(queue_msg->host);
    }
    if (queue_msg->user != NULL)
    {
        free(queue_msg->user);
    }
    if (queue_msg->passwd != NULL)
    {
        free(queue_msg->passwd);
    }
    if (queue_msg->file != NULL)
    {
        free(queue_msg->file);
    }
    if (queue_msg->no != NULL)
    {
        free(queue_msg->no);
    }
    if (queue_msg->date != NULL)
    {
        free(queue_msg->date);
    }
    if (queue_msg->function_id != NULL)
    {
        free(queue_msg->function_id);
    }
    if (queue_msg->httpurl != NULL)
    {
        free(queue_msg->httpurl);
    }
    free(queue_msg);
}

int msg_mqueue_deal(mqtt_emms2_var_t *var,QUEUE_MSG_S *queue_msg)
{
    int ret = 0;
    if (queue_msg==NULL)
    {
        return -1;
    }
    char file_pwd_tmp[200] = {0}, file_pwd_tmp2[256] = {0}, md5_tmp[256] = {0};
    long long  file_size = 0;
    sprintf(file_pwd_tmp,"%s/%s/%s.csv", get_history_dir(), queue_msg->date,queue_msg->no);
    sprintf(file_pwd_tmp2,"%s_tmp",file_pwd_tmp);
    ems_syslog(LOG_NOTICE,"check file:%s", file_pwd_tmp);
    if (access(file_pwd_tmp, F_OK) == -1) {  
        ems_syslog(LOG_NOTICE,"file '%s' already no exists.\n", file_pwd_tmp); 
        ret = -1;
    }
    else
    {
        if(sync_snapshoot_data_and_block()!=0) //操作快照中直接返回
            return -1;
        copy_file(file_pwd_tmp,file_pwd_tmp2); //当前先拷贝一份在计算md5 拷贝完需要恢复快照
        complete_and_start_snapshoot(); //操作完成
        get_file_md5(file_pwd_tmp2,md5_tmp);
        ems_syslog(LOG_NOTICE, "md5sum:%s", md5_tmp);
        file_size = check_file_size(file_pwd_tmp2);
        ret = 0;
    }
    if (ret == 0) //快照数据接口和通用不一样 要单独处理
    {
        snapshoot_response(var,queue_msg->seq,EMMS2_SUCCESS,queue_msg->function_id,file_size,md5_tmp);
    }
    else
    {
        snapshoot_response(var,queue_msg->seq,EMMS2_ERROR,queue_msg->function_id,file_size,NULL);
        return -1;
    }
    //给响应之后后开始上传文件
    if(queue_msg->httpurl == NULL) // ftp
    {
        char passwd_tmp[128] = {0};
        base64_decode(queue_msg->passwd,passwd_tmp);
        ems_syslog(LOG_NOTICE, "passwd_tmp:%s", passwd_tmp);
        char ftp_url[256] = {0},userpwd[256] = {0};
        snprintf(ftp_url,256,"ftp://%s:%d/%s",queue_msg->host,queue_msg->port,queue_msg->file);
        snprintf(userpwd,256,"%s:%s",queue_msg->user,passwd_tmp);
        if (ftp_upload_file(ftp_url,userpwd,file_pwd_tmp2)==0)
        {
            ems_syslog(LOG_NOTICE, "file upload success!");
            if (access(file_pwd_tmp2, F_OK) == 0) { //上传完成删除文件
                remove(file_pwd_tmp2);
            } else {
                perror("file is inexistence!");
            }
        }
        else
        {
            ems_syslog(LOG_NOTICE, "file upload error!");
        }
    }
    else // http post
    { 
        if (access(file_pwd_tmp2, F_OK) != 0) // 文件不存在
        { 
            return -1;
        }
        int ret = http_post_upload_file(queue_msg->httpurl, queue_msg->file, file_pwd_tmp2, NULL, NULL);
        remove(file_pwd_tmp2);//上传完成删除文件
        if (ret != 0)
        {
            ems_syslog(LOG_NOTICE, "http post upload data error!");
            return -1;
        }
    }
    return 0;
}

void *snapshoot_data_task(void *arg)
{
    pthread_detach(pthread_self());
    QUEUE_MSG_S *queue_msg = arg;
    mqtt_emms2_var_t *var = queue_msg->var;
    if (var->snapshoot_thread_state == 1) return NULL;//改之前再查一遍

    var->snapshoot_thread_state = 1;
    if(msg_mqueue_deal(var,queue_msg)!=0)
        ems_syslog(LOG_NOTICE, "msg_mqueue_deal error!");
    queue_smg_free(queue_msg); 
    var->snapshoot_thread_state = 0;
    return NULL;
}


int deal_snapshoot_data(mqtt_emms2_var_t *var,cJSON *root,int seq,char *function_id)
{
    cJSON *info = cJSON_GetObjectItem(root, "info");
    if (!info)
    {
        return -1;
    }
    cJSON* no = cJSON_GetObjectItem(info, "no");
    if (!no)
    {
        ems_syslog(LOG_ERR, "no no!");
        return -1;
    }
    cJSON* date = cJSON_GetObjectItem(info, "date");
    if (!date)
    {
        ems_syslog(LOG_ERR, "no date!");
        return -1;
    }
    cJSON* file = cJSON_GetObjectItem(info, "file");
    if (!file)
    {
        ems_syslog(LOG_ERR, "no file!");
        return -1;
    }

    QUEUE_MSG_S* queue_msg = calloc(1, sizeof(QUEUE_MSG_S)); //线程中要释放
    queue_msg->file  = strdup(file->valuestring);

    // http 优先
    cJSON* uploadurl = cJSON_GetObjectItem(info, "uploadurl");
    if (uploadurl && strlen(uploadurl->valuestring) > 0)
    {
        queue_msg->httpurl = strdup(uploadurl->valuestring);
        queue_msg->host    = NULL;
        queue_msg->port    = 0;
        queue_msg->user    = NULL;
        queue_msg->passwd  = NULL;
    }
    else // ftp
    {
        cJSON* host   = cJSON_GetObjectItem(info, "host");
        cJSON* port   = cJSON_GetObjectItem(info, "port");
        cJSON* user   = cJSON_GetObjectItem(info, "user");
        cJSON* passwd = cJSON_GetObjectItem(info, "passwd");

        if (!host || !port || !user || !passwd )
        {
            free (queue_msg);
            return -1;
        }
        queue_msg->httpurl = NULL;
        queue_msg->host   = strdup(host->valuestring);
        queue_msg->port   = port->valueint;
        queue_msg->user   = strdup(user->valuestring);
        queue_msg->passwd = strdup(passwd->valuestring);
        
    }

    queue_msg->no = strdup(no->valuestring);
    queue_msg->date = strdup(date->valuestring);
    queue_msg->seq = seq;
    queue_msg->function_id = strdup(function_id);
    queue_msg->var = var;
    ems_syslog(LOG_INFO, "emms2get host:%s port:%d,user:%s passwd:%s file:%s no:%s date:%s",queue_msg->host,queue_msg->port,queue_msg->user,queue_msg->passwd,\
        queue_msg->file,queue_msg->no,queue_msg->date);
    if (var->snapshoot_thread_state == 0)
    {
        pthread_t tid;
        if (0 == pthread_create(&tid,NULL,snapshoot_data_task,queue_msg))
        {
            pthread_setname_np(tid, "snapshoot_task");
        }
    }
    else 
    {
        ems_syslog(LOG_ERR, "snapshoot_data_task is working!");
        return -1;
    }
    return 0;
}
#ifdef EN_ACCIDENT
/**
 * @description: 事故反演数据上传完成请求
 * @param {struct mqtt_emms2_var_t} *var 
 * @param {char} *function 
 * @param {char} *file_name
 * @param {char} *file_md5
 * @return {int} 0成功，其他失败
 */
int accident_finish_upload(mqtt_emms2_var_t* var, char* function, char* file_name, long long int size, char* file_md5)
{
    char data_all[1024] = {0};
    char data_info[512] = {0};

    snprintf(data_info, sizeof(data_info), "\"messages\":{\"fileName\":\"%s\",\"md5\":\"%s\",\"fileSize\":%lld}", file_name, file_md5, size);

    package_based_field(var,data_all,function,0,data_info);

    char pub_topic[512] = {0};
    snprintf(pub_topic, sizeof(pub_topic), "emms2/"TYPE_POST"/%s/%s", var->topic_sn, function);
    ems_syslog(LOG_WARNING,"[zxf]pub topic [%s] data:[%s]",pub_topic,data_all);
    if(mqtt_publish_message(var,pub_topic,data_all,strlen(data_all))!=0)
    {
        ems_syslog(LOG_WARNING,"[zxf]send accident finish upload error!!");
        return -1;
    }
        
    return 0;
}

/**
 * @description: ftp上传事故文件
 * @param {struct mqtt_emms2_var_t} *var 
 * @param {QUEUE_MSG_S} *queue_msg
 * @param {char} *md5
 * @return {int} 0成功，其他失败
 */
int msg_accident_file_ftp_send(mqtt_emms2_var_t *var, QUEUE_MSG_S *queue_msg, struct Accident_dispose_t *clould_info)
{
    if (queue_msg == NULL  || clould_info == NULL || clould_info->taos_timeid == 0 || strlen(clould_info->md5) == 0)
    {
        return -1;
    }
    // char file_pwd_tmp[200] = {0};
    char new_file_path[256] = {0};                  // 新文件路径
    long long  file_size = 0;
    //调用从md5获取到具体的文件路径的API 
    sprintf(new_file_path, "%s/%ld_%d.tar.gz",ACCIDENT_ZIP_PATH, clould_info->taos_timeid, clould_info->code_id);

    ems_syslog(LOG_WARNING,"[zxf]check file:%s", new_file_path);
    if (access(new_file_path, F_OK) == -1) 
    {  
        ems_syslog(LOG_WARNING,"[zxf]file '%s' already no exists.\n", new_file_path); 
        return -1;
    }

    ems_syslog(LOG_WARNING, "[zxf]File copy successfully to '%s'", new_file_path);

    file_size = check_file_size(new_file_path);

    // 上传文件
    if(queue_msg->httpurl == NULL) // ftp
    {
        char passwd_tmp[128] = {0};
        base64_decode(queue_msg->passwd, passwd_tmp);
        ems_syslog(LOG_WARNING, "[zxf]passwd_tmp:%s", passwd_tmp);
        char ftp_url[256] = {0}, userpwd[256] = {0};
        snprintf(ftp_url,256,"ftp://%s:%d/%s",queue_msg->host,queue_msg->port,queue_msg->file);
        snprintf(userpwd,256,"%s:%s",queue_msg->user,passwd_tmp);
        if (ftp_upload_file(ftp_url,userpwd,new_file_path)==0)
        {
            ems_syslog(LOG_WARNING, "[zxf]file upload success!");
            if(accident_finish_upload(var, FUNCTION_ACCIDENT_FINISH_UPLOAD, queue_msg->file, file_size, clould_info->md5) < 0)
            {
                remove(new_file_path);
                ems_syslog(LOG_WARNING, "[zxf]send accident finish upload message error!");
                return -1;
            }
            remove(new_file_path);
            update_accident_send_file_time();
            return 0;
        }
        else
        {
            remove(new_file_path);
            ems_syslog(LOG_WARNING, "[zxf]file upload error!");
            return -1;
        }
    }
    else // http post
    { 
        int ret = http_post_upload_file(queue_msg->httpurl, queue_msg->file, new_file_path, NULL, NULL);
        // remove(file_pwd_tmp);//上传完成删除文件
        if (ret == 0)
        {
            ems_syslog(LOG_WARNING, "[zxf]http file upload success!");
            if(accident_finish_upload(var, FUNCTION_ACCIDENT_FINISH_UPLOAD, queue_msg->file, file_size, clould_info->md5) < 0)
            {
                ems_syslog(LOG_WARNING, "[zxf]http send accident finish upload message error!");
                return -1;
            }
            update_accident_send_file_time();
            return 0;
        }
        else
        {
            ems_syslog(LOG_WARNING, "[zxf]http post upload data error!");
            return -1;
        }
    }

    return 0;
}

/**
 * @description: 处理事故反演数据请求消息
 * @param {struct mqtt_emms2_var_t} *var 
 * @param {cJSON} *root
 * @return {int} 0成功，其他失败
 */
int deal_accident_upload_info(mqtt_emms2_var_t *var, cJSON *root)
{
    // if(accident_upload_control.accident_upload_state != ACCIDENT_STATE_WAITING_FTP)
    // {
    //     return -1;
    // }
    cJSON *info = cJSON_GetObjectItem(root, "info");
    if (!info)
    {
        return -1;
    }

    cJSON* file = cJSON_GetObjectItem(info, "file");
    if (!file)
    {
        ems_syslog(LOG_ERR, "[zxf]no file!");
        return -1;
    }

    QUEUE_MSG_S* ftp_msg = calloc(1, sizeof(QUEUE_MSG_S));
    ftp_msg->file = strdup(file->valuestring);
    // http 优先
    cJSON* uploadurl = cJSON_GetObjectItem(info, "uploadurl");
    if (uploadurl && strlen(uploadurl->valuestring) > 0)
    {
        ftp_msg->httpurl = strdup(uploadurl->valuestring);
        ftp_msg->host    = NULL;
        ftp_msg->port    = 0;
        ftp_msg->user    = NULL;
        ftp_msg->passwd  = NULL;
        ems_syslog(LOG_ERR, "[zxf]emms2get accident ftp httpurl%s file:%s",ftp_msg->httpurl,ftp_msg->file);
        ftp_msg->file = strdup(file->valuestring);
    }
    else // ftp
    {
        cJSON* host   = cJSON_GetObjectItem(info, "host");
        cJSON* port   = cJSON_GetObjectItem(info, "port");
        cJSON* user   = cJSON_GetObjectItem(info, "user");
        cJSON* passwd = cJSON_GetObjectItem(info, "passwd");

        if (!host || !port || !user || !passwd )
        {
            free (ftp_msg);
            return -1;
        }
        ftp_msg->httpurl = NULL;
        ftp_msg->host   = strdup(host->valuestring);
        ftp_msg->port   = port->valueint;
        ftp_msg->user   = strdup(user->valuestring);
        ftp_msg->passwd = strdup(passwd->valuestring);
        ftp_msg->var    = var;
        ems_syslog(LOG_ERR, "[zxf]emms2get accident ftp host:%s port:%d,user:%s passwd:%s file:%s",ftp_msg->host,ftp_msg->port,ftp_msg->user,ftp_msg->passwd,ftp_msg->file);
    }
    
    if(msg_accident_file_ftp_send(var, ftp_msg, clould_info) < 0)
    {
        ems_syslog(LOG_ERR, "[zxf]msg_accident_file_ftp_send false!");
        queue_smg_free(ftp_msg);
        return -1;
    }
    queue_smg_free(ftp_msg);
    accident_upload_control.accident_upload_state = ACCIDENT_STATE_WAITING_CLEAR;
    accident_upload_control.rcv_accident_clear_flag = false;
    ems_syslog(LOG_ERR, "[zxf]accident state ACCIDENT_STATE_WAITING_CLEAR!");
    return 0;
}

/**
 * @description: 处理事故反演本地事故记录清除消息
 * @param {struct cJSON} *root 
 * @return {int} 0成功，其他失败
 */
int deal_accident_clear(cJSON *root)
{
    // if(accident_upload_control.accident_upload_state != ACCIDENT_STATE_WAITING_CLEAR)
    // {
    //     return -1;
    // }
    cJSON *info = cJSON_GetObjectItem(root, "info");
    if (!info)
    {
        ems_syslog(LOG_ERR, "[zxf]no info!");
        return -1;
    }

    cJSON* report_md5 = cJSON_GetObjectItem(info, "md5");
    if (!report_md5)
    {
        ems_syslog(LOG_ERR, "[zxf]no report md5!");
        return -1;
    }

    char* report_md5_str = report_md5->valuestring;  

    if (strlen(report_md5_str) >= sizeof(accident_upload_control.rcv_md5))  
    {  
        ems_syslog(LOG_ERR, "[zxf]MD5 string too long!");  
        return -1;  
    }  

    strncpy(accident_upload_control.rcv_md5, report_md5_str, sizeof(accident_upload_control.rcv_md5));  
    ems_syslog(LOG_ERR, "[zxf]accident rcv celar md5 %s!", accident_upload_control.rcv_md5);
    // accident_upload_control.md5[sizeof(accident_upload_control.md5)-1] = '\0';  

    int del_res = accident_del_api(report_md5_str);
    ems_syslog(LOG_INFO, "[zxf]accident del api res = [%d]!", del_res);
    accident_upload_control.rcv_accident_clear_flag = true;

    return 0;
}

/**
 * @description: 事故反演历史事故上报
 * @param {struct mqtt_emms2_var_t} *var 
 * @param {char} *function
 * @param {char} *file_md5
 * @return {int} 0成功，其他失败
 */
int accident_info_report_upload(mqtt_emms2_var_t* var, char* function, char* file_md5)
{
    char data_all[1024] = {0};
    char data_info[64] = {0};
    
    char new_file_path[256] = {0};                  // 新文件路径
    long long  file_size = 0;
    //调用从md5获取到具体的文件路径的API 
    sprintf(new_file_path, "%s/%ld_%d.tar.gz",ACCIDENT_ZIP_PATH, clould_info->taos_timeid, clould_info->code_id);

    ems_syslog(LOG_WARNING,"[zxf]check file:%s", new_file_path);
    if (access(new_file_path, F_OK) == -1) 
    {  
        ems_syslog(LOG_WARNING,"[zxf]file '%s' already no exists.\n", new_file_path); 
        return -1;
    }

    ems_syslog(LOG_WARNING, "[zxf]File copy successfully to '%s'", new_file_path);

    file_size = check_file_size(new_file_path);

    sprintf(data_info,"\"messages\":{\"md5\":\"%s\", \"fileSize\":%lld}",file_md5, file_size);

    package_based_field(var,data_all,function,0,data_info);

    char pub_topic[512] = {0};
    sprintf(pub_topic,"emms2/"TYPE_POST"/%s/%s",var->topic_sn,function);
    ems_syslog(LOG_WARNING,"[zxf]pub topic [%s] data:[%s]",pub_topic,data_all);
    if(mqtt_publish_message(var,pub_topic,data_all,strlen(data_all))!=0)
    {
        ems_syslog(LOG_WARNING,"[zxf]send accident info report data error!!");
        return -1;
    }
        
    update_accident_send_md5_time();
    accident_upload_control.accident_upload_state = ACCIDENT_STATE_WAITING_FTP;
    ems_syslog(LOG_WARNING,"[zxf]accident state ACCIDENT_STATE_WAITING_FTP!");
    return 0;
}
#endif
extern int upload_template_request(const mqtt_emms2_var_t* var, cJSON* info);
extern int upload_alarm_requset(const mqtt_emms2_var_t* var, cJSON* info);

char* get_payload_data(char* topic, char* payload, int payloadLen, char* message_type, char* function_id, char* tmp_c)
{
    char* get_data = NULL;
    if (strstr(topic, "lz4"))
    {
        int ret_dec = -1;
        size_t multiplier = LZ4_DEFAULT_BUFFER_MULTIPLIER;
        while (ret_dec < 0 && multiplier <= LZ4_MAX_BUFFER_MULTIPLIER) 
        {
            size_t get_data_len = payloadLen * multiplier;
            if (get_data != NULL) {
                free(get_data);
                get_data = NULL;
            }

            get_data = malloc(get_data_len);
            if (get_data == NULL) {
                ems_syslog(LOG_ERR, "Memory allocation failed!");
                return NULL;
            }
            memset(get_data, 0, get_data_len);
            ret_dec = LZ4_decompress_safe(payload, get_data, payloadLen, get_data_len);
            if (ret_dec < 0) {
                multiplier *= 10;
                ems_syslog(LOG_DEBUG, "Buffer too small, retrying with multiplier: %zu", multiplier);
            }
        }
        if (ret_dec < 0)
        {
            ems_syslog(LOG_ERR, "Not enough space!!!");
            if (get_data != NULL) free(get_data);
            return NULL;
        }
        ems_syslog(LOG_DEBUG, "decompress  payload %s", get_data);
        sscanf(topic, "emms2/%[^/]/%[^/]/%[^/]/lz4", message_type, tmp_c, function_id);
    }
    else
    {
        get_data = malloc(payloadLen + 1);
        memset(get_data, 0, payloadLen + 1);
        strncpy(get_data, payload, payloadLen);
        sscanf(topic, "emms2/%[^/]/%[^/]/%s", message_type, tmp_c, function_id);
    }
    return get_data;
}

int bulid_and_send_multi_rate_info(mqtt_emms2_var_t *var, int seq, struct tags_matrix *matrix)
{
    if (!var || !matrix) 
    {
        ems_syslog(LOG_ERR, "rate:[zxf]bulid_get_multi_rate_info invalid input parameters");
        return -1;
    }

    // 创建info对象
    cJSON *info = cJSON_CreateObject();
    if (!info) 
    {
        ems_syslog(LOG_ERR, "rate:[zxf]Failed to create info object");
        return -1;
    }

    // 添加dev_no字段
    cJSON_AddStringToObject(info, "dev_no", matrix->dev_no);

    // 创建tags数组
    cJSON *tags = cJSON_CreateArray();
    if (!tags) 
    {
        ems_syslog(LOG_ERR, "rate:[zxf]Failed to create tags array");
        cJSON_Delete(info);
        return -1;
    }
    cJSON_AddItemToObject(info, "tags", tags);

    // 遍历所有行
    for (int i = 0; i < matrix->row_num; i++) 
    {
        struct tags_matrix_row *row = &matrix->rows[i];
        
        // 创建行数组
        cJSON *row_array = cJSON_CreateArray();
        if (!row_array) 
        {
            ems_syslog(LOG_ERR, "rate:[zxf]Failed to create row array");
            continue;
        }
        cJSON_AddItemToArray(tags, row_array);

        // 添加月份和日期 [mon, day]
        cJSON *date_array = cJSON_CreateArray();
        cJSON_AddItemToArray(date_array, cJSON_CreateNumber(row->mon));
        cJSON_AddItemToArray(date_array, cJSON_CreateNumber(row->day));
        cJSON_AddItemToArray(row_array, date_array);

        // 遍历所有节点
        for (int j = 0; j < row->Node_num; j++) 
        {
            struct tags_matrix_Node *node = &row->Node[j];
            
            // 创建时间类型数组 ["HH:MM", type]
            cJSON *node_array = cJSON_CreateArray();
            char time_str[6];
            snprintf(time_str, sizeof(time_str), "%02d:%02d", node->hour, node->min);
            
            cJSON_AddItemToArray(node_array, cJSON_CreateString(time_str));
            cJSON_AddItemToArray(node_array, cJSON_CreateNumber(node->type));
            cJSON_AddItemToArray(row_array, node_array);
        }
    }

    // 将info对象转换为字符串
    char *info_str = cJSON_PrintUnformatted(info);
    cJSON_Delete(info);
    if (!info_str) 
    {
        ems_syslog(LOG_ERR, "rate:[zxf]Failed to print info JSON");
        return -1;
    }

    // 构建完整报文
    struct lnxall_buff ex_data;
    lbuff_init(&ex_data, 512);

    lbuff_sprintf(&ex_data, "\"info\":%s", info_str);
    free(info_str);

    struct lnxall_buff  data_all;
    lbuff_init(&data_all,(int)strlen(ex_data.bufptr));

    lbuff_sprintf(&data_all,"{\"funcId\":\"%s\",\"lcSN\":\"%s\",\"seq\": %d,\"time\":%ld,%s}",FUNCTION_RATE_GET,var->topic_sn, seq,time(NULL), ex_data.bufptr);

    char pub_topic[512] = {0};
    sprintf(pub_topic,"emms2/"TYPE_GET_RESP"/%s/%s", var->topic_sn,FUNCTION_RATE_GET);
    ems_syslog(LOG_ERR,"rate:多费率报文内容: pub topic [%s] data:[%s]",pub_topic,data_all.bufptr);
    if(mqtt_publish_message(var,pub_topic,data_all.bufptr,strlen(data_all.bufptr))!=0)
        ems_syslog(LOG_WARNING,"rate:send data error!!");
    lbuff_free(&data_all);
    lbuff_free(&ex_data);

    return 0;
}

void *get_rate_data_task(void *arg)
{
    pthread_detach(pthread_self());
    RATE_MSG_S *rate_msg = arg;
    mqtt_emms2_var_t *var = rate_msg->var;

    struct tags_matrix tmp_tags_matrix = {0};

    // TEST 写固定的召测报文

    // 设置设备编号
    strncpy(tmp_tags_matrix.dev_no, rate_msg->dev_no, DEV_NO_LEN-1);

    // // 读取设备费率信息
    int read_res = rate_cfg_read(rate_msg->dev_no, &tmp_tags_matrix);

    if(read_res == 0)
    {
        bulid_and_send_multi_rate_info(var, rate_msg->seq, &tmp_tags_matrix);
    }
    else
    {
        //发送读取失败报文
        char data_info[64] = {0};
        char data_all[512] = {0};
        sprintf(data_info,"\"result\":%d,\"info\":{}",-1);
        package_based_field(var,data_all,FUNCTION_RATE_GET,rate_msg->seq,data_info);

        char pub_topic[512] = {0};
        sprintf(pub_topic,"emms2/"TYPE_GET_RESP"/%s/%s",var->topic_sn,FUNCTION_RATE_GET);
        ems_syslog(LOG_ERR,"pub topic [%s] data:[%s]",pub_topic,data_all);
        if(mqtt_publish_message(var,pub_topic,data_all,strlen(data_all))!=0)
            ems_syslog(LOG_WARNING,"send data error!!");
    }

    if (rate_msg->dev_no != NULL)
    {
        free(rate_msg->dev_no);
    }
    if(rate_msg)
    {
        free(rate_msg);
    }

    for(int i = 0; i < tmp_tags_matrix.row_num; i++)
    {
        if(tmp_tags_matrix.rows[i].Node)
            free(tmp_tags_matrix.rows[i].Node);
    }

    if(tmp_tags_matrix.rows)
    {
        free(tmp_tags_matrix.rows);
    }

    return NULL;
}

int get_multi_rate(mqtt_emms2_var_t *var, cJSON *root, int seq)
{
    if (!root) 
    {
        ems_syslog(LOG_ERR, "[zxf]Input get_multi_rate JSON is NULL");
        return -1;
    }

    // 获取info对象
    cJSON *info = cJSON_GetObjectItemCaseSensitive(root, "info");
    if (!info || !cJSON_IsObject(info)) 
    {
        ems_syslog(LOG_ERR, "[zxf]Invalid 'info' field");
        return -1;
    }

    // 获取dev_no字段
    cJSON *dev_no = cJSON_GetObjectItemCaseSensitive(info, "dev_no");
    if (!dev_no || !cJSON_IsString(dev_no)) 
    {
        ems_syslog(LOG_ERR, "[zxf]Invalid 'dev_no' string field");
        return -1;
    }

    RATE_MSG_S* rate_msg = calloc(1, sizeof(RATE_MSG_S)); //线程中要释放
    rate_msg->seq = seq;
    rate_msg->var = var;
    rate_msg->dev_no = strdup(dev_no->valuestring);

    ems_syslog(LOG_INFO, "[zxf]dev_no = [%s]", rate_msg->dev_no);

    pthread_t tid;
    if (0 == pthread_create(&tid,NULL,get_rate_data_task,rate_msg))
    {
        pthread_setname_np(tid, "get_rate_task");
    }
    
    return 0;
}

static int mqtt_handle_recv_msg(void *obj, mqtt_message_t *mqtt_msg)
{
    mqtt_emms2_var_t *var = obj;
    static int first_heart_flag = 0;
    ems_syslog(LOG_INFO, "%s, received MQTT topic:%s payload length:%d",var->station_id, mqtt_msg->topic, mqtt_msg->payloadLen);
    
    char message_type[64] = {0},tmp_c[64] = {0},function_id[64] = {0};
    char *get_data = get_payload_data(mqtt_msg->topic, mqtt_msg->payload, mqtt_msg->payloadLen, message_type, function_id, tmp_c);
    
    ems_syslog(LOG_INFO, "received message_type:%s function_id:%s", message_type, function_id);
    
    if( get_data == NULL)
    {
        return -1;
    }

    cJSON *root = cJSON_Parse(get_data);
    if (!root)
    {
        ems_syslog(LOG_ERR, "json parse error");
        if (get_data != NULL)
        {
            free(get_data);
        }
        return -1;
    }
    cJSON *item = cJSON_GetObjectItem(root, "seq");
    if (!item)
    {
        ems_syslog(LOG_INFO, "no seq!");
        goto END;
    }
    int seq = item->valueint;
    ems_syslog(LOG_INFO, "%s, seq %d!", var->station_id, seq);
    int ret = 0;
    if (strcmp(message_type, TYPE_SET)==0) //设置参数
    {
        if(((proto_forward_t*)get_pcs_ctrl_var()->device_layer_ptr)->opt_log != NULL)
        {
            char *root_str = NULL;
            root_str = cJSON_PrintUnformatted(root);
            opt_log_add(((proto_forward_t*)get_pcs_ctrl_var()->device_layer_ptr)->opt_log, var->station_id, "EMS", _("Set"), function_id, root_str);
        }

        if (strcmp(function_id, FUNCTION_HISTORY)==0)
        {
            if(deal_snapshoot_data(var,root,seq,function_id)<0)
            {
                ret = -1;
                goto SET_END;
            }
            else
            {
                goto END;
            }
        }
        ems_syslog(LOG_INFO, "emms2 "TYPE_SET);
        if(north_control != -1)
        {
            ems_syslog(LOG_ERR, "rcv schedule msg but north_control = [%d]!", north_control);
            ret = -2;
            goto SET_END;
        }
        if (strcmp(function_id, FUNCTION_TELECONTROL)==0)
        {
            item = cJSON_GetObjectItem(root, "tags");
            if (!item)
            {
                ems_syslog(LOG_ERR, "no seq!");
                ret = -1;
                goto SET_END;
            }
            int tag_num = cJSON_GetArraySize(item);
            for (size_t i = 0; i < tag_num; i++)
            {
                cJSON *tag = cJSON_GetArrayItem(item,i);
                char name_ems[64] = {0};
                sprintf(name_ems,DEV_NO_EMS".%s",tag->string);
                ems_syslog(LOG_NOTICE, "set:%s -> %.2f",name_ems,tag->valuedouble);
                set_telecontrol(name_ems,tag->valuedouble);    
                ret = 0;//写了就成功            
            }
        }
        else if(strcmp(function_id, FUNCTION_SUBTELECONTROL)==0)
        {
            cJSON *queryList = cJSON_GetObjectItem(root, "messages");
            if (!queryList)
            {
                ems_syslog(LOG_ERR, "no queryList!");
                ret = -1;
                goto SET_END;
            }
            int dev_num = cJSON_GetArraySize(queryList);
            for (size_t i = 0; i < dev_num; i++)
            {
                cJSON *devn = cJSON_GetArrayItem(queryList,i);
                cJSON *no = cJSON_GetObjectItem(devn, "no");
                if (no!=NULL)
                {
                    cJSON *tags = cJSON_GetObjectItem(devn, "tags");
                    int tag_num = cJSON_GetArraySize(tags);
                    for (size_t j = 0; j < tag_num; j++)
                    {
                        cJSON* tag = cJSON_GetArrayItem(tags, j);
                        if (strcmp(tag->string, TAG_BMS_SET_CONNECT) == 0) // 特殊处理 todo
                        {
                            int v = tag->valueint;
                            ret = pe_set_value(no->valuestring, FUNCTION_SET_CONNECT, v);
                        }
                        else
                        {
                            char name_tag[64] = {0};
                            sprintf(name_tag, "%s.%s", no->valuestring, tag->string);
                            ems_syslog(LOG_NOTICE, "set:%s -> %.2f", name_tag, tag->valuedouble);
                            set_telecontrol(name_tag, tag->valuedouble);
                            ret = 0; //写了就成功
                        }

                    }
                }
            }
        }
        else if(strcmp(function_id, FUNCTION_SCHEDULE)==0)
        {
            cJSON *messages = cJSON_GetObjectItem(root, "messages");
            if (!messages)
            {
                ems_syslog(LOG_ERR, "no messages!");
                ret = -1;
                goto SET_END;
            }
            if (parse_all_plan(messages)<0)
            {
                ret = -1;
                goto SET_END;
            }
            g_reflush_all_plan = 1;
            save_all_plan_to_file(messages);
            server_add_new_log(LOG_CODE_EMMS2_WRITE_PLAN,"EMS",LOG_MODE_TRIG,1,NULL);
        }else if(strcmp(function_id, FUNCTION_ALARM_DISABLE)==0)
        { 
            cJSON *messages = cJSON_GetObjectItem(root, "messages");
            if (!messages)
            {
                ems_syslog(LOG_ERR, "no messages!");
                ret = -1;
                goto SET_END;
            }
            cJSON *alarmCodes = cJSON_GetObjectItem(messages, "alarmCodes");
            if (!alarmCodes)
            {
                ems_syslog(LOG_ERR, "no alarmCodes!");
                ret = -1;
                goto SET_END;
            }
            save_forbidden_alarm_codes(alarmCodes);
            server_add_new_log(LOG_CODE_EMMS2_DISABLE_ALARM,"EMS",LOG_MODE_TRIG,1,NULL);
        }
        else if (strcmp(function_id, FUNCTION_SUB_CTRL)==0)
        {
            cJSON *info = cJSON_GetObjectItem(root, "info");
            if (info)
            {
                cJSON *timeout = cJSON_GetObjectItem(info, "timeout");
                if (!timeout)
                {
                    ems_syslog(LOG_ERR, "no timeout use default!");
                    var->pub_dev->end  = DEFAULT_PUB_TIME_OUT + time(NULL);
                }
                else
                {
                    var->pub_dev->end  = timeout->valueint + time(NULL);
                }

                cJSON *intervalTime = cJSON_GetObjectItem(info, "period");
                if (!intervalTime)
                {
                    ems_syslog(LOG_ERR, "no intervalTime use default!");
                    var->pub_dev->period = DEFAULT_PUB_INTERVAL_TIME;
                }
                else
                {
                    var->pub_dev->period  = intervalTime->valueint;
                }
            }
            else
            {
                ems_syslog(LOG_ERR, "no info everything use default!");
                var->pub_dev->end  = DEFAULT_PUB_TIME_OUT;
                var->pub_dev->period  = DEFAULT_PUB_INTERVAL_TIME;
            }
            ems_syslog(LOG_INFO, "var->pub_dev->period  :%ld ,var->pub_dev->end  :%ld",var->pub_dev->period ,var->pub_dev->end );
            //action不用管超时自动停止
        }
        else if (strcmp(function_id, FUNCTION_PORXYFRPC) == 0)
        {
            ret = frpc_proxy_cmd(var,root,seq);
            if (ret < 0)
            {
                goto SET_END;
            }
            goto END;
        }
        else if (strcmp(function_id, FUNCTION_CABINETINFO) == 0)
        {
            cJSON* message = cJSON_GetObjectItem(root, "message");
            mqtt_set_cabinet_info(var, message, seq);
            goto END;
        }
        else if (strcmp(function_id, FUNCTION_CONFIG) == 0 || strcmp(function_id, FUNCTION_ID_UPGRADE) == 0)
        {
            goto END;
        }
        else if (strcmp(function_id, FUNCTION_RATE_SET)==0 && 0 == strcmp(message_type, TYPE_SET))
        {
            ems_syslog(LOG_ERR, "[zxf]rcv Rate Set message!");
            report_rate_set_return(var,TYPE_GET_RESP,EMMS2_TAG_TYPE_R,seq, rate_mqtt_analysis_cfg(root));
        }
        else if (strcmp(function_id, FUNCTION_ID_OTAUPDATE) == 0){  
            ota_upgrade_cmd(var, root, seq);
            goto END;
        }
#ifdef EN_ACCIDENT
        else if (strcmp(function_id, FUNCTION_ACCIDENT_CFG)==0)
        {
            ems_syslog(LOG_ERR, "[zxf]rcv AccidentConfig message!");
            ret = deal_accident_rule_config(root);
            if(ret != 0)
            {
                goto SET_END;
            }
            else
            {
                int save_result = save_accident_rule_config_file(root);
                if(save_result != 0)
                {
                    ems_syslog(LOG_ERR, "[zxf]save accident rule json file error!");
                    ret = 2;
                    goto SET_END;
                }
            }
        }
#endif
        else if (strcmp(function_id, FUNCTION_RATE_SET)==0)
        {
            ems_syslog(LOG_ERR, "[zxf]rcv Rate Set message!");
            report_rate_set_return(var,TYPE_GET_RESP,EMMS2_TAG_TYPE_R,seq, rate_mqtt_analysis_cfg(root));
        
        }
        else
        {
            ret = -1;
        }
        
SET_END:
        if (ret == 0)
        {
            write_response(var,seq,EMMS2_SUCCESS,function_id);
        }
        else if (ret == -1)
        {
            write_response(var,seq,EMMS2_ERROR,function_id);
        }
        else
        {
            write_response(var,seq,EMMS2_ERROR,function_id);
            server_add_new_log(LOG_CODE_EMMS2_NO_CTRL_PERMISSON,"EMS",LOG_MODE_TRIG,1,NULL);
        }
        
    }
    else if (strcmp(message_type, TYPE_GET)==0) //获取数据
    {       
        ems_syslog(LOG_INFO, "emms2 "TYPE_GET);
        if (strcmp(function_id, FUNCTION_TELEMETRY)==0)
        {
            ems_syslog(LOG_INFO, "emms2 get ems read data ");
           //上报EMS所有 只读数据
           report_all_ems_data_by_type(var,TYPE_GET_RESP,EMMS2_TAG_TYPE_R,seq);

        //    server_add_new_log(LOG_CODE_EMMS2_READ,"EMS",LOG_MODE_TRIG,1,NULL);
        }
        else if (strcmp(function_id, FUNCTION_TELECONTROL)==0)
        {
            ems_syslog(LOG_INFO, "emms2 get ems write data ");
           //上报EMS所有 读写数据
           report_all_ems_data_by_type(var,TYPE_GET_RESP,EMMS2_TAG_TYPE_RW,seq);
        //    server_add_new_log(LOG_CODE_EMMS2_READ,"EMS",LOG_MODE_TRIG,2,NULL);
        }
        else if ((strcmp(function_id, FUNCTION_SUBTELEMETRY)==0)||(strcmp(function_id, FUNCTION_SUBTELECONTROL)==0))
        {
            ems_syslog(LOG_INFO, "emms2 get dev data ");
            EMMS2_TAG_TYPE tag_type;
            if (strcmp(function_id, FUNCTION_SUBTELEMETRY)==0)  tag_type = EMMS2_TAG_TYPE_R;
            else tag_type = EMMS2_TAG_TYPE_RW;
            cJSON *queryList = cJSON_GetObjectItem(root, "messages");
            if (!queryList)
            {
                ems_syslog(LOG_ERR, "no queryList!");
                goto END;
            }
            int dev_num = cJSON_GetArraySize(queryList);
            char dev_name[MAX_MQTT_EMMS2_DEV_NUM][EMMS2_MAX_BASE_SHORT_LEN] = {{0}};
            char *dev_pr[MAX_MQTT_EMMS2_DEV_NUM] = {0};
            ems_syslog(LOG_INFO, "dev_num %d",dev_num);
            for (int i = 0; i < dev_num; i++)
            {
                cJSON *devn = cJSON_GetArrayItem(queryList,i);
                cJSON *no = cJSON_GetObjectItem(devn, "no");
                if (no!=NULL)
                {
                   
                    if(no->valuestring){
                        strcpy(dev_name[i],no->valuestring);
                        dev_pr[i] = dev_name[i];
                        ems_syslog(LOG_INFO, "no [%d][%s][%s]",i,no->valuestring,dev_name[i]);
                    }
                }
                // server_add_new_log(LOG_CODE_EMMS2_READ,no->valuestring,LOG_MODE_TRIG,3,no->valuestring);
            }
            if (dev_num>0)        
                report_dev_data_by_type_and_name(var,TYPE_GET_RESP,dev_pr,dev_num,tag_type,seq);
        }
        else if (strcmp(function_id, FUNCTION_ALL_ALARM)==0)
        {
            report_all_alarm(var,seq);
            // server_add_new_log(LOG_CODE_EMMS2_READ,"EMS",LOG_MODE_TRIG,4,NULL);
        }
        else if (strcmp(function_id, FUNCTION_SCHEDULE) == 0) // 返回7天策略
        {
            ems_syslog(LOG_INFO, "emms2 get ems read data: FUNCTION_SCHEDULE ");
            if(get_cabinet_info()->is_slave == 0)
            {
                mqtt_emms2_PolicDataGet(var); // 返回7天策略
            }
            else
            {
                ems_syslog(LOG_INFO, "emms2 get ems read data: FUNCTION_SCHEDULE in slave mode ,no response ");
            }
        }
        else if( strcmp(function_id, FUNCTION_TEMPLATE) == 0) 
        {
            cJSON* info = cJSON_GetObjectItem(root, "info");
            ems_syslog(LOG_INFO, "MQTT recv template get request, station %s", var->station_id);
            if (upload_template_request(var, info) != 0)
            {
                ems_syslog(LOG_ERR, "MQTT handle template get request fail, station %s", var->station_id);
            }
        }
        else if (strcmp(function_id, FUNCTION_ALARMFILE)==0)
        {
            cJSON* info = cJSON_GetObjectItem(root, "info");
            upload_alarm_requset(var, info);
        }
        else if (strcmp(function_id, FUNCTION_RATE_GET)==0)
        {
            get_multi_rate(var, root, seq);
        }
    }
    else if (strcmp(message_type, TYPE_POST_RESP)==0) //上报响应
    {
        
        if (strcmp(function_id, FUNCTION_ID_LOGIN)==0)
        {
            var->emms2_status = EMMS2_DEVREPORTING; //登陆成功
            ems_syslog(LOG_CRIT, "[云/站控]:%s已回复登陆请求,数据开始上报", var->station_id);
            set_alarm_need_report(true);
        }
        else if ((strcmp(function_id, FUNCTION_HEARTBEAT)==0))//首次收到心跳需要重新上报一次设备表
        {
            var->offlinecount--;
            if (!first_heart_flag)
            {
                mqtt_emms2_devInfo(var);
                report_all_data(var);
                first_heart_flag = 1;
                ems_syslog(LOG_INFO,"first heart report over!!");
            }            
        }
        else if (strcmp(function_id, FUNCTION_SUBTELEMETRY)==0)
        {
            cond_signal(&var->cond,&var->mutex);
        }
        else if (strcmp(function_id, FUNCTION_CONTINUE_DATA) == 0)
        {
            cond_signal(&var->offline_db->history_cond, &var->offline_db->history_mutex);
        }
        else if( strcmp(function_id, FUNCTION_CABINETINFO) == 0)
        {
            sync_cabinet_info_response(var, root);
        }
#ifdef EN_ACCIDENT
        else if (0 == strcmp(var->station_id, MQTT_EMMS2_STATION_ID) && strcmp(function_id, FUNCTION_ACCIDENT_UPLOAD_INFO)==0)
        {
            if(deal_accident_upload_info(var, root)<0)
            {
                accident_upload_control.accident_upload_state = ACCIDENT_STATE_INIT;
            }
        }
        else if (0 == strcmp(var->station_id, MQTT_EMMS2_STATION_ID) && strcmp(function_id, FUNCTION_ACCIDENT_CLEAR)==0)
        {
            if(deal_accident_clear(root)<0)
            {
                accident_upload_control.accident_upload_state = ACCIDENT_STATE_INIT;
            }
        }
#endif
    }

END:
    if (root != NULL)
    {
        cJSON_Delete(root);
    }
    if (get_data != NULL)
    {
        free(get_data);
    }
    return 0;
    
}



int get_network_type();
int get_4g_iccid(char* dst, int max_len) ;
int get_vaild_4g_iccid(char* dst, int max_len);

static int get_4g_info(int* net, char* iccids, char* valid, int max_len)
{
    *net = get_network_type();
    if (get_4g_iccid(iccids, max_len) <= 0)
    {
        ems_syslog(LOG_ERR, "get_4g_iccid failed");
    }
    
    if (get_vaild_4g_iccid(valid, max_len) <= 0)
    {
        ems_syslog(LOG_ERR, "get_vaild_4g_iccid failed ");
    }
    return 0;
}

static void mqtt_emms2_login(mqtt_emms2_var_t *var)//登录
{
    int fd;
    // cabinet_info_t *cat_var = get_cabinet_info();
    char* data_all = malloc(4096 * 2);
    char* data_info = malloc(2048 * 2);

    // Read 4G location file
    fd = open(LOCATION_4G_FILEPATH, O_RDONLY);
    if (fd >= 0) {
        ssize_t len, lsize;
        char buf[128];
        len = read(fd, buf, sizeof(buf));
        if (len > 0) {
            char * dst;
            lsize = sizeof(cabinet_info.cab_location);
            if (len >= lsize)
                len = lsize - 1;
            pthread_mutex_lock(&cab_info_mutex);
            dst = cabinet_info.cab_location;
            memcpy(dst, buf, len);
            pthread_mutex_unlock(&cab_info_mutex);
            dst[len] = '\0';
        }
        close(fd);
    }
    char iccids[64] = {0};
    char valid[64]  = {0};
    int  net         = 0;
    get_4g_info(&net, iccids, valid, sizeof(iccids));

    pthread_mutex_lock(&cab_info_mutex);
    sprintf(data_info, "\"info\":{\"id\":\"%s\",\"name\":\"%s\",\"station\":\"%s\",\"capacity\":%d,\"power\":%d,\"time\":\"%s\",\"local\":\"%s\",\"station_name\":\"%s\",\
    \"vendor\":\"LNXALL\",\"model\":\"%s\",\"softVersion\":\"%s\",\"tenant_id\":\"%s\",\"mode\":%d, \"network\":%d, \"iccids\":\"%s\", \"validCard\":\"%s\"}",
            cabinet_info.cab_id,
            cabinet_info.name,
            cabinet_info.cab_location,
            cabinet_info.cab_capacity,
            cabinet_info.cab_power,
            cabinet_info.cab_time,
            cabinet_info.cab_local,
            cabinet_info.cab_station_name,
            get_board_name(),
            SW_VERSION,
            cabinet_info.tenant_id,
            cabinet_info.ctrl_mode,
            net,
            iccids,
            net == 1 ? valid : "");
    pthread_mutex_unlock(&cab_info_mutex);
    // sprintf(data_info,"\"info\":{",);

    // for(int i =0;i<cabinet_info_tab_count_g;i++)
    // {
    //     char *key=NULL;
    //     key=cabinet_info_tab[i].cab_key;
    //     if(key==NULL)continue;
    //     if(strcmp(cabinet_info_tab[i].cab_type,"string")==0)
    //     {
    //         sprintf(data_tmp,"\"%s\":\"%s\",",key,((char *)&cabinet_info + cabinet_info_tab[i].cab_member_ofs));
    //         strcat(data_info,data_tmp);
    //     }
    //     else if(strcmp(cabinet_info_tab[i].cab_type,"int")==0)
    //     {
    //         sprintf(data_tmp,"\"%s\":%d,",key,*((int *)((char *)&cabinet_info + cabinet_info_tab[i].cab_member_ofs)));
    //         strcat(data_info,data_tmp);
    //     }
    // }
    // if (strlen(data_info)>strlen("\"info\":{"))
    // {
    //     data_info[strlen(data_info)-1] = '\0';
    // }
    
    // strcat(data_info,"}");

    package_based_field(var,data_all,FUNCTION_ID_LOGIN,0,data_info);

    char pub_topic[512] = {0};
    sprintf(pub_topic,"emms2/"TYPE_POST"/%s/"FUNCTION_ID_LOGIN,var->topic_sn);
    ems_syslog(LOG_INFO,"pub topic [%s] data:[%s]",pub_topic,data_all);
    if(mqtt_publish_message(var,pub_topic,data_all,strlen(data_all))!=0)
        ems_syslog(LOG_WARNING,"send data error!!");

    free(data_info);
    free(data_all);
}

static void mqtt_emms2_get_dido_info_from_config(cJSON* data_info,bool enable_lc)
{
    if (data_info == NULL)
    {
        return;
    }

    cJSON* dev_class = cJSON_GetObjectItem(data_info, "dev_class");
    if (dev_class == NULL)
    {
        return;
    }
    int dev_num = cJSON_GetArraySize(dev_class);

    for (int i = 0; i < dev_num; i++)
    {
        cJSON* class = cJSON_GetArrayItem(dev_class, i);
        cJSON* devs  = cJSON_GetObjectItem(class, "devs");

        if(devs == NULL)
        {
            continue;
        }

        int dev_count = cJSON_GetArraySize(devs);
        for (int j = 0; j < dev_count; j++)
        {
            cJSON* dev = cJSON_GetArrayItem(devs, j);
            cJSON* no  = cJSON_GetObjectItem(dev, "no");

            if (!(no && cJSON_IsString(no) && strstr(no->valuestring, "DIDO") != NULL))
                continue;

            cJSON* template = cJSON_GetObjectItem(dev, "template");

            if (template == NULL || cJSON_IsString(template) == false || strlen(template->valuestring) == 0)
                continue;
            char template_path[128] = {0};
            if (template->valuestring[0] == '/') //绝对路径
            {
                snprintf(template_path, 128, "%s", template->valuestring);
            }
            else
            {
                snprintf(template_path, 128, "%s/%s", TEMPLATE_DIR, template->valuestring);
            }

            char* pdata = read_file_data(template_path);
            if (pdata == NULL)
            {
                proto_syslog(LOG_ERR, "read cfg file %s error", template_path);
                continue;
            }

            cJSON* root = json_parse_string_with_comments(pdata);
            free(pdata);
            if (!root)
            {
                proto_syslog(LOG_ERR, "parse pdata %s error", pdata);
                continue;
            }

            cJSON* di = cJSON_GetObjectItemCaseSensitive(root, "read_reg");
            if (di)
            {
                int    di_cnt       = cJSON_GetArraySize(di);
                cJSON* dev_read_reg = cJSON_CreateArray();
                for (int k = 0; k < di_cnt; k++)
                {
                    cJSON* tag = cJSON_GetArrayItem(di, k);
                    if (tag == NULL) continue;
                    cJSON* tag_name = cJSON_GetObjectItemCaseSensitive(tag, "tag");
                    cJSON* name     = cJSON_GetObjectItemCaseSensitive(tag, "name");
                    if (tag_name && name)
                    {
                        cJSON* p = cJSON_CreateObject();
                        cJSON_AddStringToObject(p, "tag", tag_name->valuestring);
                        cJSON_AddStringToObject(p, "name", name->valuestring);
                        cJSON_AddItemToArray(dev_read_reg, p);
                    }
                }
                cJSON_AddItemToObject(dev, "read_reg", dev_read_reg);
            }

            cJSON* do_ = cJSON_GetObjectItemCaseSensitive(root, "write_reg");
            if (do_)
            {
                cJSON* dev_write_reg = cJSON_CreateArray();
                int    do_snt        = cJSON_GetArraySize(do_);
                for (int k = 0; k < do_snt; k++)
                {
                    cJSON* tag      = cJSON_GetArrayItem(do_, k);
                    cJSON* tag_name = cJSON_GetObjectItemCaseSensitive(tag, "tag");
                    cJSON* name     = cJSON_GetObjectItemCaseSensitive(tag, "name");
                    if (tag_name && name)
                    {
                        cJSON* p = cJSON_CreateObject();
                        cJSON_AddStringToObject(p, "tag", tag_name->valuestring);
                        cJSON_AddStringToObject(p, "name", name->valuestring);
                        cJSON_AddItemToArray(dev_write_reg, p);
                    }
                }
                cJSON_AddItemToObject(dev, "write_reg", dev_write_reg);
            }
            cJSON_Delete(root);
        }
    }
}

int ems_slave_dev_config_merge(cJSON *cfg1)
{
    int cur_sum = get_pcs_ctrl_var()->lc_list_len;
    lc_it *lc_cfg = get_pcs_ctrl_var()->lc_list;
    int ret = 0;

    for (int i = 0; i < cur_sum; i++)
    {
        if (lc_cfg[i].slave == NULL)
        {
            ems_syslog(LOG_INFO, "ems_slave is DISABLED");
            continue;
        }
        if (lc_cfg[i].type == eLC_TYPE_EMS)
        {
            const char *p = slave_ems_get_dev_info(lc_cfg[i].slave);
            if (p == NULL)
            {
                ems_syslog(LOG_ERR, "Failed to get dev_info for slave %s", lc_cfg[i].slave->no);
                ret++;
                continue;
            }
            cJSON *current = string_to_json_and_add_prefix(p, lc_cfg[i].no, lc_cfg[i].name);
            if (current == NULL)
            {
                ems_syslog(LOG_ERR, "Failed to convert string to JSON for slave %s", lc_cfg[i].slave->no);
                ret++;
                continue;
            }

            if (dev_cfg_merge(cfg1, current, 0) != 0)
            {
                ems_syslog(LOG_ERR, "merge ems slave dev_info failed for slave %s", lc_cfg[i].slave->no);
                ret = 0;
            }
            cJSON_Delete(current);
        }
    }

    return ret;
}

int update_group_and_parent_for_lc_dev(cJSON *dev, cJSON *config_dev)
{
    cJSON *dev_no = cJSON_GetObjectItemCaseSensitive(dev, "no");
    cJSON *dev_tree = cJSON_GetObjectItemCaseSensitive(config_dev, "dev_tree");
    cJSON *lc_no = cJSON_GetObjectItemCaseSensitive(config_dev, "no");
    if ((dev_no == NULL) || (dev_tree == NULL) || (lc_no == NULL))
    {
        return -1;
    }

    const char *dot = strstr(dev_no->valuestring, ".");
    if (dot == NULL)
    {
        return -2;
    }//
    if ((strlen(lc_no->valuestring) != (dot - dev_no->valuestring)) || (memcmp(dev_no->valuestring, lc_no->valuestring, dot - dev_no->valuestring) != 0))
    {
        return -3;
    }
    dot++;
    //
    int cnt = cJSON_GetArraySize(dev_tree);
    for (int i = 0; i < cnt; ++i)
    {
        cJSON *cfg_dev = cJSON_GetArrayItem(dev_tree, i);
        cJSON *cfg_no = cJSON_GetObjectItemCaseSensitive(cfg_dev, "no");
        cJSON *cfg_group = cJSON_GetObjectItemCaseSensitive(cfg_dev, "group");
        if ((cfg_no == NULL) || (cfg_group == NULL))
        {
            continue;
        }
        
        if (strcmp(cfg_no->valuestring, dot) == 0)
        {
            cJSON *group = cJSON_GetObjectItemCaseSensitive(dev, "group");
            if (group)
            {
                cJSON_SetValuestring(group, cfg_group->valuestring);
            }
            else
            {
                cJSON_AddStringToObject(dev, "group", cfg_group->valuestring);
            }
            cJSON *parent = cJSON_GetObjectItemCaseSensitive(cfg_dev, "parent");
            if (parent != NULL)
            {
                char parent_buf[64] = {0};
                snprintf(parent_buf, sizeof(parent_buf), "%s.%s", lc_no->valuestring, parent->valuestring);
                cJSON_AddStringToObject(dev, "parent", parent_buf);
            }
            return 0;
        }
    }
    return -99;
}

int is_slave_dev(const char *dev_no)
{
    return ((strstr(dev_no, "LC") == dev_no) && (strstr(dev_no, ".") != NULL)) ? 1 : 0;
}

cJSON *update_dev_group_info_with_config(cJSON *dev_group_info, const char *config_path)
{
    char *config_data = read_file_data(config_path);
    if (!config_data)
    {
        ems_syslog(LOG_ERR, "Failed to read config file: %s", config_path);
        return NULL;
    }
    cJSON *config_json = cJSON_Parse(config_data);
    free(config_data);
    if (!config_json)
    {

        ems_syslog(LOG_ERR, "Failed to parse config file: %s", config_path);
        return NULL;
    }

    cJSON *dev_tree = cJSON_GetObjectItemCaseSensitive(config_json, "dev_tree");
    if (!dev_tree || !cJSON_IsArray(dev_tree))
    {
        ems_syslog(LOG_ERR, "No valid 'dev_tree' array found in config.");
        cJSON_Delete(config_json);
        return NULL;
    }
    cJSON *dev_class = cJSON_GetObjectItemCaseSensitive(dev_group_info, "dev_class");
    if (!dev_class || !cJSON_IsArray(dev_class))
    {
        ems_syslog(LOG_ERR, "No valid 'dev_class' array found in dev_group_info");
        cJSON_Delete(config_json);
        return NULL;
    }

    int dev_class_count = cJSON_GetArraySize(dev_class);

    for (int i = 0; i < dev_class_count; ++i)
    {
        cJSON *class_item = cJSON_GetArrayItem(dev_class, i);
        if (!class_item || !cJSON_IsObject(class_item))
            continue;

        cJSON *devs = cJSON_GetObjectItemCaseSensitive(class_item, "devs");
        if (!devs || !cJSON_IsArray(devs))
            continue;

        int dev_count = cJSON_GetArraySize(devs);
        for (int j = 0; j < dev_count; ++j)
        {
            cJSON *dev = cJSON_GetArrayItem(devs, j);
            if (!dev || !cJSON_IsObject(dev))
                continue;

            cJSON *no = cJSON_GetObjectItemCaseSensitive(dev, "no");

            const char *dev_no = no->valuestring;

            int config_dev_count = cJSON_GetArraySize(dev_tree);
            if (0 == is_slave_dev(dev_no)) // 不是从机上面的设备
            {
                for (int l = 0; l < config_dev_count; ++l)
                {
                    cJSON *config_dev = cJSON_GetArrayItem(dev_tree, l);
                    cJSON *config_no = cJSON_GetObjectItemCaseSensitive(config_dev, "no");
                    //
                    if (strcmp(dev_no, config_no->valuestring) == 0)
                    {
                        cJSON *config_group = cJSON_GetObjectItemCaseSensitive(config_dev, "group");
                        if (config_group && cJSON_IsString(config_group))
                        {
                            cJSON *group = cJSON_GetObjectItemCaseSensitive(dev, "group");
                            if (group)
                            {
                                cJSON_SetValuestring(group, config_group->valuestring);
                            }
                            else
                            {
                                cJSON_AddStringToObject(dev, "group", config_group->valuestring);
                            }
                            cJSON *parent = cJSON_GetObjectItemCaseSensitive(config_dev, "parent");
                            if (parent != NULL)
                            {
                                cJSON_AddStringToObject(dev, "parent", parent->valuestring);
                            }
                        }
                        else
                        {
                            ems_syslog(LOG_WARNING, "Configuration for device '%s' is missing 'group' field.", dev_no);
                        }
                        break;
                    }
                }
            }
            else
            {
                for (int l = 0; l < config_dev_count; ++l)
                {
                    
                    cJSON *config_dev = cJSON_GetArrayItem(dev_tree, l);
                    cJSON *config_type = cJSON_GetObjectItemCaseSensitive(config_dev, "type");
                    if (strcmp(config_type->valuestring, "LC") == 0)
                    {
                        if (0 == update_group_and_parent_for_lc_dev(dev, config_dev))
                        {
                            break;
                        }
                    }

                }
            }
        }
    }
    cJSON_Delete(config_json);
    return dev_group_info;
}

static void mqtt_emms2_devInfo(mqtt_emms2_var_t *var) //设备表
{
    // cabinet_info_t *cat_var = get_cabinet_info();
    char *all_dev_info = NULL;
    proto_forward_t *proto_var = get_proto_forward_var();
    if (proto_var->dev_info_string == NULL)
    {
        ems_syslog(LOG_ERR,"dev_info_string is NULL!!");
        return;
    }
    //

    cJSON* dev_info = cJSON_Parse(proto_var->dev_info_string);
    char* p_data = NULL;
    if (cabinet_info.ctrl_mode == EMS_MODE_ZC && (p_data = read_file_data(EMS_LC_ALL_DEVICE_CFG)) != NULL)
    {
        // 主从模式
        cJSON* cfg1 = json_parse_string_with_comments(p_data);
        dev_cfg_merge(cfg1, dev_info, 0);                      // 合并主从配置文件


        if (g_usercfg_variant.EnDataCollect) {
            ems_slave_dev_config_merge(cfg1);
        }
        
        mqtt_emms2_get_dido_info_from_config(cfg1, true);   // 添加DIDO名称
        //
        cJSON_Delete(dev_info);
        dev_info = cfg1;
        free(p_data);
    }
    else
    {
        // 对等模式
        mqtt_emms2_get_dido_info_from_config(dev_info, false); // 添加DIDO名称
    }
    update_dev_group_info_with_config(dev_info ,DEV_GROUP_INFO_PATH);
    update_dev_rate_info(dev_info);
    all_dev_info = cJSON_PrintUnformatted(dev_info);
    cJSON_Delete(dev_info);

    if (all_dev_info == NULL)
    {
        ems_syslog(LOG_ERR, "mqtt send DeviceInfo err, all_dev_info is NULL!!");
        return;
    }
    char* info = (char*)malloc(strlen(all_dev_info) + 32);
    sprintf(info, "\"info\":%s", all_dev_info);

    char *data_all = calloc(1,strlen(info)+256);
    package_based_field(var,data_all,FUNCTION_DEVINFO,0,info);

    char pub_topic[512] = {0};
    sprintf(pub_topic,"emms2/"TYPE_POST"/%s/"FUNCTION_DEVINFO,var->topic_sn);
    ems_syslog(LOG_INFO,"pub topic [%s] data:[%s]",pub_topic,data_all);
    if(mqtt_publish_message(var,pub_topic,data_all,strlen(data_all))!=0)
        ems_syslog(LOG_WARNING,"send data error!!");
    free(all_dev_info);
    free(info);
    free(data_all);
    //
}
int get_4g_sign();
static void mqtt_emms2_heartBeat(mqtt_emms2_var_t *var) //心跳
{
    char data_all[1024] = {0};
    char data_info[1024] = {0};

    int net = get_network_type();
    sprintf(data_info, "\"info\":{\"alarmCnt\":%d, \"sign\":%d}", get_alarm_cnt(), net == 2 ? 31 : get_4g_sign());

    package_based_field(var,data_all,FUNCTION_HEARTBEAT,0,data_info);
    
    char pub_topic[512] = {0};
    sprintf(pub_topic,"emms2/"TYPE_POST"/%s/"FUNCTION_HEARTBEAT,var->topic_sn);
    ems_syslog(LOG_INFO,"pub topic [%s] data:[%s]",pub_topic,data_all);
    if(mqtt_publish_message(var,pub_topic,data_all,strlen(data_all))!=0)
        ems_syslog(LOG_WARNING,"send data error!!");
    else
        var->offlinecount++;
}


static void mqtt_emms2_subscribe_all(mqtt_emms2_var_t *var)
{
    mqtt_session_t *session = var->session;
    char sub_topic[EMMS2_MAX_BASE_LEN] = {0};

    sprintf(sub_topic,"emms2/"TYPE_POST_RESP"/%s/+",var->topic_sn); //登陆应答
    mqtt_session_subscribe(session, sub_topic);
    ems_syslog(LOG_INFO, "sub topic:%s", sub_topic);

    sprintf(sub_topic,"emms2/"TYPE_GET"/%s/+",var->topic_sn); //平台召测
    mqtt_session_subscribe(session, sub_topic);
    ems_syslog(LOG_INFO, "sub topic:%s", sub_topic);

    sprintf(sub_topic,"emms2/"TYPE_POST_RESP"/%s/+/lz4",var->topic_sn); //登陆应答
    mqtt_session_subscribe(session, sub_topic);
    ems_syslog(LOG_INFO, "sub topic:%s", sub_topic);

    sprintf(sub_topic,"emms2/"TYPE_GET"/%s/+/lz4",var->topic_sn); //平台召测
    mqtt_session_subscribe(session, sub_topic);
    ems_syslog(LOG_INFO, "sub topic:%s", sub_topic);

    if ((strncmp(var->station_id, MQTT_EMMS2_STATION_ID, strlen(MQTT_EMMS2_STATION_ID)) == 0 && cabinet_info.en_mqtt_emms2_write) ||
        (strncmp(var->station_id, MQTT_STACTRL_STATION_ID, strlen(MQTT_STACTRL_STATION_ID)) == 0 && cabinet_info.en_mqtt_stactrl_write))
    {
        sprintf(sub_topic,"emms2/"TYPE_SET"/%s/+",var->topic_sn); //平台设置
        mqtt_session_subscribe(session, sub_topic);
        ems_syslog(LOG_INFO, "sub topic:%s", sub_topic);

        sprintf(sub_topic,"emms2/"TYPE_SET"/%s/+/lz4",var->topic_sn); //平台设置
        mqtt_session_subscribe(session, sub_topic);
        ems_syslog(LOG_INFO, "sub topic:%s", sub_topic);
    }
}
static void mqtt_emms2_PolicDataGet(mqtt_emms2_var_t *var)
{

#define POLIC_DATA_GET_NUM 7 // 获取7天的策略数据 可修改

    char pub_topic[512] = {0};    // topic
    extern all_plan g_all_plan[]; // 策略缓冲区

    time_t t = time(NULL);

    struct tm target_time = {};  
    localtime_r(&t, &target_time);
    struct tm *local_time = &target_time;

    int year = local_time->tm_year + 1900;
    int month = local_time->tm_mon + 1;
    int day = local_time->tm_mday;

    int ret = 0;   // 检查是否内存获取失败 1:是 0:否
    int statu = 0; // 是否遍历到目标日期策略 1:是 0 否
    int j = 0;     // 记录获取数量 受限于:POLIC_DATA_GET_NUM

    struct lnxall_buff output;
    lbuff_init(&output, 1024);

    // 添加消息头部
    if (-1 == lbuff_sprintf(&output, "{\"funcId\":\"%s\",\"lcSN\":\"%s\",\"seq\": %d,\"time\":%ld,\"messages\":[", FUNCTION_SCHEDULE, var->topic_sn, 0, time(NULL)))
        ret = 1;

    for (int i = 0; i < MAX_DAY_YEAR && j < POLIC_DATA_GET_NUM; i++)
    {
        // 遍历到当前日期
        if (j == 0 && (g_all_plan[i].year == year) && (g_all_plan[i].mounth == month) && (g_all_plan[i].date == day))
        {
            statu = 1;
        }

        if (statu)
        {

            j += 1;
            if(g_all_plan[i].year == 0 && g_all_plan[i].mounth == 0 && g_all_plan[i].date == 0)
                continue;
            // 当天时间写入
            if (-1 == lbuff_sprintf(&output, "{\"Date\":\"%04d-%02d-%02d\",\"Schedule\":[",
                                    g_all_plan[i].year, g_all_plan[i].mounth, g_all_plan[i].date))
                ret = 1;

            // 获取当天策略写入策略
            for (int k = 0; k < NUM_TOU_SEG_MAX; k++)
            {
                if (g_all_plan[i].schedule[k].StHour == 0 && g_all_plan[i].schedule[k].StMin == 0 &&
                    g_all_plan[i].schedule[k].EnHour == 0 && g_all_plan[i].schedule[k].EnMin == 0)
                {
                    break; // 所有时间为0，终止循环
                }

                const char *separator = "";    // json 尾部 ,
                if (k < NUM_TOU_SEG_MAX - 1 && // 确保 k + 1 不会越界
                    !(g_all_plan[i].schedule[k + 1].StHour == 0 && g_all_plan[i].schedule[k + 1].StMin == 0 &&
                      g_all_plan[i].schedule[k + 1].EnHour == 0 && g_all_plan[i].schedule[k + 1].EnMin == 0))
                {
                    separator = ",";
                }

                if (-1 == lbuff_sprintf(&output, "{\"StartTime\":\"%02d:%02d\",\"EndTime\":\"%02d:%02d\",\"Power\":%d,\"RePower\":%d,\"SOCUpperLimit\":%d,\"SOCLowerLimit\":%d}%s",
                                        g_all_plan[i].schedule[k].StHour, g_all_plan[i].schedule[k].StMin,
                                        g_all_plan[i].schedule[k].EnHour, g_all_plan[i].schedule[k].EnMin,
                                        g_all_plan[i].schedule[k].Power,
                                        g_all_plan[i].schedule[k].RePower,
                                        g_all_plan[i].schedule[k].SocUp,
                                        g_all_plan[i].schedule[k].SocLower,
                                        separator))
                {
                    ret = 1;
                }
            }

            // json 边界判断
            if (i < MAX_DAY_YEAR && j < POLIC_DATA_GET_NUM)
                lbuff_sprintf(&output, "]},");
            else
                lbuff_sprintf(&output, "]}");
        }
    }

    // 结束补充
    lbuff_sprintf(&output, "]}");

    if (ret) // 内存申请失败 跳过消息发送
    {
        ems_syslog(LOG_WARNING, "mqtt_emms2_PolicDataGet:Error: Out of space Failed to request memory space!!");
        goto PolicDataGet_END;
    }
    //  topic组装
    sprintf(pub_topic, "emms2/" TYPE_GET_RESP "/%s/" FUNCTION_SCHEDULE, var->topic_sn);


    // 接口发送
    if (0 != mqtt_publish_message(var, pub_topic, output.bufptr, strlen(output.bufptr)))
        ems_syslog(LOG_WARNING, "mqtt_emms2_PolicDataGet:Error: send data error!!");
PolicDataGet_END:
    ems_syslog(LOG_WARNING, "send data");

    // 内存释放
    lbuff_free(&output);
}
#if 0
static bool check_data_change_deadzone(tag_t *tag)
{
    if (((tag_white_list*)tag->some_pr[EMMS2_TAG_EX_PR_INDEX]) == NULL)
    {
        ems_syslog(LOG_NOTICE, "[%s] emms2 tag pr 1 is null", tag->name);
        return false;
    }
    //ems_syslog(LOG_NOTICE, "[%s] emms2 tag pr 1 is not null", tag->name);
    float deadZone = (((tag_white_list*)tag->some_pr[EMMS2_TAG_EX_PR_INDEX])->deadZone);
    if (deadZone<0)
        return false;

    if(tag->data_type == TYPE_TAG_INT)
    {
        if(abs(tag->read_cache.to_int - tag->last_read_cache.to_int)>=deadZone)
            return true;
    }
    else 
    {
        if(fabsf(tag->read_cache.to_float - tag->last_read_cache.to_float)>=deadZone)
            return true;
    }

    return false;
}

#endif
#if 0
static void mqtt_cloud_data_cb(void *context,tag_t *tag) 
{
    if(!check_data_change_deadzone(tag))return;

    mqtt_emms2_var_t *var = context;
    if(var->emms2_status == EMMS2_HEARTBEAT)
    {
        char *dev_no = NULL;
        if (strstr(((device_t *)(tag->dev_ptr))->no,"CELL"))
        {
            dev_no = get_bms_no_cell_no(((device_t *)(tag->dev_ptr))->no);
        }
        else
        {
            dev_no = ((device_t*)tag->dev_ptr)->no;
        }

        char *needreport=NULL;
        if (hash_intptr_findptr(var->dev_hash,dev_no,strlen(dev_no), (void **)&needreport) < 0)
        { 
            needreport  = calloc(1,1);
            *needreport = 1;
            int ret = hash_intptr_addptr(&var->dev_hash,dev_no,strlen(dev_no),needreport);
            if (ret < 0)
            {
                ems_syslog(LOG_NOTICE, "hash_intptr_addptr err!");
            }
        }
        else
        {
            *needreport = 1;
        }
        return ;
    }
}
#endif
// 获取网络端口状态
char * get_net_string(char *data_info)
{
    char *data_tmp = NULL;

    sprintf(data_info,"\"messages\":[");
    cmd_call(PATH_NET_STATUS, 3, &data_tmp, 1);
    if (data_tmp!=NULL)
    {
        strcat(data_info,data_tmp);
        free(data_tmp);
    }
    
    strcat(data_info,"]");  
    return data_info;
}
extern int get_alarm_string(struct lnxall_buff * lbuf, int type);
int report_all_alarm(mqtt_emms2_var_t *var,int seq)//回复所有告警
{
    struct lnxall_buff data_info;
    lbuff_init(&data_info, 8192);
    get_alarm_string(&data_info,0);
    char* data_all = calloc(1, data_info.curlen + 256);
    package_based_field(var,data_all,FUNCTION_ALL_ALARM,seq,data_info.bufptr);
    char pub_topic[512] = {0};
    sprintf(pub_topic,"emms2/"TYPE_GET_RESP"/%s/"FUNCTION_ALL_ALARM,var->topic_sn);
    ems_syslog(LOG_INFO,"pub topic [%s] data:[%s]",pub_topic,data_all);
    int ret = mqtt_publish_message(var,pub_topic,data_all,strlen(data_all));
    lbuff_free(&data_info);
    free(data_all);
    return ret;
}
extern int get_all_ems_sn(struct lnxall_buff* buf);
extern int get_is_rx_broadcast_data(void);
static void check_ems_num_change(mqtt_emms2_var_t *var) {

  double grid_online = 0;
  if (dev_get_dev_tag_float_with_result(DEV_NO_GRID_MT, STATE_ONLINE,
                                        &grid_online) != 0 ||
      get_is_rx_broadcast_data() == true) {
    return;
  }

  struct lnxall_buff buf;
  lbuff_init(&buf, 1024);
  lbuff_sprintf(&buf, "\"info\":[");
  get_all_ems_sn(&buf);
  if (buf.bufptr[buf.curlen - 1] == ',') {
    buf.bufptr[buf.curlen - 1] = ']';
  } else {
    lbuff_sprintf(&buf, "]");
  }
  char pub_topic[512] = {0};
  char *data_all = calloc(1, buf.curlen + 256);
  package_based_field(var, data_all, FUNCTION_GRID_GROUP, var->seq, buf.bufptr);
  sprintf(pub_topic, "emms2/" TYPE_POST "/%s/" FUNCTION_GRID_GROUP,
          var->topic_sn);
  ems_syslog(LOG_INFO, "pub topic [%s] data:[%s]", pub_topic, data_all);
  if (mqtt_publish_message(var, pub_topic, data_all, strlen(data_all)) != 0) {
    ems_syslog(LOG_ERR, "mqtt publish ems sn failed");
  }
  lbuff_free(&buf);
  free(data_all);
}

static void check_alarm_info(mqtt_emms2_var_t *var) 
{
    if (get_alarm_report_queue_size(var) <= 0)
    {
        return;
    }

    char data_tmp[1024]      = {0};
    struct lnxall_buff data_info;
    lbuff_init(&data_info, 8192);
    lbuff_sprintf(&data_info, "\"messages\":[");
    alarm_report_t* alarm = NULL;
    while ((alarm = pop_alarm_from_report_queue(var)) != NULL)
    {
        alarm_info_t* alarm_info = alarm->alarm;
        if (alarm_info != NULL)
        {
            snprintf(data_tmp,
                     sizeof(data_tmp),
                     "{\"alarmCode\":\"%d\",\"reasonCode\":%d,\"severity\":%d,\"entityInstance\":\"%s\",\"additionalParam\":\"%s\", \"raiseTime\":%" PRIu64 "},",
                     alarm_info->alarm_code,
                     alarm_info->reason_code,
                     alarm->status == 1 ? alarm_info->severity : EMMS2_SEVERITY_CLEAR,
                     alarm->entity,
                     alarm->additional == NULL ? "" : alarm->additional,
                     alarm->status == 1 ? (alarm->raise / 1000) : (alarm->recovery / 1000));
            lbuff_append(&data_info, data_tmp, strlen(data_tmp));
        }

        destory_alarm_report(alarm);
    }

    if (data_info.bufptr[data_info.curlen -1] == ',')
        data_info.bufptr[data_info.curlen -1] = ']';
    char* data_all = malloc(data_info.curlen + 256);
    package_based_field(var, data_all, FUNCTION_ALARM, 0, data_info.bufptr);
    char pub_topic[512] = {0};
    sprintf(pub_topic, "emms2/" TYPE_POST "/%s/" FUNCTION_ALARM, var->topic_sn);
    ems_syslog(LOG_INFO, "pub topic [%s] data:[%s]", pub_topic, data_all);
    int ret = mqtt_publish_message(var, pub_topic, data_all, strlen(data_all));
    if (ret != 0)
    {
        ems_syslog(LOG_ERR, "mqtt publish failed, ret:%d", ret);
    }
    lbuff_free(&data_info);
    free(data_all);
}

static void get_cell_data_by_bms_name(char *bms_no, struct lnxall_buff * bufp)
{
    struct lnxall_buff lbuf_volt;
    struct lnxall_buff lbuf_temp;
    struct lnxall_buff lbuf_pole;
    proto_forward_t *var_proto = get_proto_forward_var();
    unsigned int len_volt, len_temp, len_pole;
    char no[128] = {0};
    lbuff_init(&lbuf_volt, 4096);
    lbuff_init(&lbuf_temp, 4096);
    lbuff_init(&lbuf_pole, 4096);


    lbuff_sprintf(&lbuf_volt, "\"CellVol\":[");
    len_volt = lbuf_volt.curlen;

    lbuff_sprintf(&lbuf_temp, "\"CellTemp\":[");
    len_temp = lbuf_temp.curlen;

    lbuff_sprintf(&lbuf_pole, "\"PoleTemp\":[");
    len_pole = lbuf_pole.curlen;
    sprintf(no,"%s_CELL",bms_no);
    for (size_t i = 0; i < var_proto->channels_size; i++)
    {
        for (size_t j = 0; j < var_proto->channels[i]->devs_size; j++)
        {
            if (strstr(var_proto->channels[i]->devs[j]->no,no))
            {
                for (size_t k = 0; k < var_proto->channels[i]->devs[j]->tags.ro_tags_size; k++)
                {
                    
                    if(strstr(var_proto->channels[i]->devs[j]->tags.ro_tags[k]->name,"RkCeVolt")) //电压
                    {
                        TAG_VALUE_TO_STRING_ARRAY(var_proto->channels[i]->devs[j]->tags.ro_tags[k], lbuf_volt);
                    }
                    if(strstr(var_proto->channels[i]->devs[j]->tags.ro_tags[k]->name,"RkCeTemp")) //温度
                    {
                        TAG_VALUE_TO_STRING_ARRAY(var_proto->channels[i]->devs[j]->tags.ro_tags[k], lbuf_temp);
                    }
                    if(strstr(var_proto->channels[i]->devs[j]->tags.ro_tags[k]->name,"PoleTemp")) //温度
                    {
                        TAG_VALUE_TO_STRING_ARRAY(var_proto->channels[i]->devs[j]->tags.ro_tags[k], lbuf_pole);
                    }
                }
            }
        }
    }

    if (lbuf_volt.curlen > len_volt)
        lbuff_setlen(&lbuf_volt, (int) (lbuf_volt.curlen - 1));
    lbuff_append(&lbuf_volt, "],", 2);

    if (lbuf_temp.curlen > len_temp)
        lbuff_setlen(&lbuf_temp, (int) (lbuf_temp.curlen - 1));
    lbuff_append(&lbuf_temp, "],", 2);

    if (lbuf_pole.curlen > len_pole)
        lbuff_setlen(&lbuf_pole, (int) (lbuf_pole.curlen - 1));
    lbuff_append(&lbuf_pole, "],", 2);

    lbuff_append(bufp, lbuf_volt.bufptr, (int) lbuf_volt.curlen);
    lbuff_append(bufp, lbuf_temp.bufptr, (int) lbuf_temp.curlen);
    lbuff_append(bufp, lbuf_pole.bufptr, (int) lbuf_pole.curlen);

    lbuff_free(&lbuf_volt);
    lbuff_free(&lbuf_temp);
    lbuff_free(&lbuf_pole);
}

static void get_ems_slave_cell_data_by_bms_name(char *bms_no, struct lnxall_buff *bufp)
{
    struct lnxall_buff lbuf_volt;
    struct lnxall_buff lbuf_temp;
    struct lnxall_buff lbuf_pole;
    int cur_sum = get_pcs_ctrl_var()->lc_list_len;
    lc_it *lc_cfg = get_pcs_ctrl_var()->lc_list;
    unsigned int len_volt, len_temp, len_pole;
    char no[128] = {0};
    lbuff_init(&lbuf_volt, 4096);
    lbuff_init(&lbuf_temp, 4096);
    lbuff_init(&lbuf_pole, 4096);

    lbuff_sprintf(&lbuf_volt, "\"CellVol\":[");
    len_volt = lbuf_volt.curlen;

    lbuff_sprintf(&lbuf_temp, "\"CellTemp\":[");
    len_temp = lbuf_temp.curlen;

    lbuff_sprintf(&lbuf_pole, "\"PoleTemp\":[");
    len_pole = lbuf_pole.curlen;
    sprintf(no, "%s_CELL", bms_no);
    for (int i = 0; i < cur_sum; i++)
    {
        if (lc_cfg[i].type == eLC_TYPE_EMS)
        {
            for (int j = 0; j < lc_cfg[i].slave->dev_sum; j++)
            {
                const device_t *dev_ptr = slave_ems_get_dev(lc_cfg[i].slave, lc_cfg[i].slave->devs[j].disp_name);
                if (strstr(dev_ptr->disp_name, no) || slave_ems_get_dev_oline(lc_cfg[i].slave, dev_ptr->disp_name) == 1)
                {
                    for (size_t k = 0; k < dev_ptr->tags.ro_tags_size; k++)
                    {
                        tag_t *tag = dev_ptr->tags.ro_tags[k];
                        if (strstr(tag->name, "RkCeVolt"))
                        {
                            EMS_SLAVE_TAG_VALUE_TO_STRING_ARRAY(*tag, lbuf_volt);
                        }
                        if (strstr(tag->name, "RkCeTemp"))
                        {
                            EMS_SLAVE_TAG_VALUE_TO_STRING_ARRAY(*tag, lbuf_temp);
                        }
                        if (strstr(tag->name, "PoleTemp"))
                        {
                            EMS_SLAVE_TAG_VALUE_TO_STRING_ARRAY(*tag, lbuf_pole);
                        }
                    }
                }
            }
        }
    }
    if (lbuf_volt.curlen > len_volt)
        lbuff_setlen(&lbuf_volt, (int) (lbuf_volt.curlen - 1));
    lbuff_append(&lbuf_volt, "],", 2);

    if (lbuf_temp.curlen > len_temp)
        lbuff_setlen(&lbuf_temp, (int) (lbuf_temp.curlen - 1));
    lbuff_append(&lbuf_temp, "],", 2);

    if (lbuf_pole.curlen > len_pole)
        lbuff_setlen(&lbuf_pole, (int) (lbuf_pole.curlen - 1));
    lbuff_append(&lbuf_pole, "],", 2);

    lbuff_append(bufp, lbuf_volt.bufptr, (int) lbuf_volt.curlen);
    lbuff_append(bufp, lbuf_temp.bufptr, (int) lbuf_temp.curlen);
    lbuff_append(bufp, lbuf_pole.bufptr, (int) lbuf_pole.curlen);

    lbuff_free(&lbuf_volt);
    lbuff_free(&lbuf_temp);
    lbuff_free(&lbuf_pole);
}

static void get_std_enum_list_value(mqtt_emms2_var_t* var, char* dev_no, struct lnxall_buff* l_buf)
{
    if (strcmp(var->station_id, MQTT_STACTRL_STATION_ID) != 0)
    {
        return;
    }
    int value = -1;
    if (strstr(dev_no, DEV_NO_BMS))
    {
        dev_get_function_int(dev_no, FUNCTION_GET_STATUS, &value);
        lbuff_sprintf(l_buf, "\"%s\":%d,", "StdRunStatus", value >= 0 ? value : -1);
        dev_get_function_int(dev_no, FUNCTION_GET_PRECHAR, &value);
        lbuff_sprintf(l_buf, "\"%s\":%d,", "StdPrechaStatus", value >= 0 ? value : -1);
    }
    else if (strstr(dev_no, DEV_NO_PCS) && strstr(dev_no, DEV_NO_PCS_MT) == NULL)
    {
        dev_get_function_int(dev_no, FUNCTION_GET_ONOFF, &value);
        lbuff_sprintf(l_buf, "\"%s\":%d,", "StdRunStatus", value >= 0 ? value : -1);
        dev_get_function_int(dev_no, FUNCTION_GET_CONNECT, &value);
        lbuff_sprintf(l_buf, "\"%s\":%d,", "StdGridStatus", value >= 0 ? value : -1);
        dev_get_function_int(dev_no, FUNCTION_GET_MODE, &value);
        lbuff_sprintf(l_buf, "\"%s\":%d,", "StdCtrlStatus", value >= 0 ? value : -1);
    }
}

void package_dev_data_by_type_and_name(mqtt_emms2_var_t *var,const char *type,char** dev_name,int dev_num,EMMS2_TAG_TYPE tag_type,int seq,char **data_all,char **pub_topic)
{
    struct lnxall_buff lbuf;
    unsigned int inlen;
    unsigned int len_tmp = 0;

    lbuff_init(&lbuf, 4096);
    lbuff_sprintf(&lbuf, "\"messages\":[");
    inlen = lbuf.curlen;
    
    proto_forward_t *var_proto = get_proto_forward_var();
    for (size_t i = 0; i < var_proto->channels_size; i++)
    {
        for (size_t j = 0; j < var_proto->channels[i]->devs_size; j++)
        {
            if ((strcmp(var_proto->channels[i]->devs[j]->no,DEV_NO_EMS)==0)||strstr(var_proto->channels[i]->devs[j]->no,"CELL"))
            {
                continue;
            }
            if (dev_num > 0)
            {
                for (int l = 0; l < dev_num; l++)
                {
                    ems_syslog(LOG_INFO, "dev_name[%s]", dev_name[l]);
                    if (dev_name[l]==NULL)
                    {
                        ems_syslog(LOG_INFO, "dev_name[%d] NULL", l);
                        continue;
                    }
                    
                    if (strcmp(var_proto->channels[i]->devs[j]->no,dev_name[l])==0)
                    {
                        ems_syslog(LOG_ERR, "emms2 get %s data!",dev_name[l]);
                        lbuff_sprintf(&lbuf, "{\"no\":\"%s\",\"tags\":{",dev_name[l]);
                        len_tmp = lbuf.curlen;
                        switch (tag_type)
                        {
                        case EMMS2_TAG_TYPE_R:
                            get_std_enum_list_value(var, var_proto->channels[i]->devs[j]->no, &lbuf);
                            if (strstr(var_proto->channels[i]->devs[j]->no,"BMS"))
                            {
                                get_cell_data_by_bms_name(var_proto->channels[i]->devs[j]->no, &lbuf);
                            }
                            for (size_t k = 0; k < var_proto->channels[i]->devs[j]->tags.sys_info_tags_size; k++)
                            {
                                // if ((var_proto->channels[i]->devs[j]->tags.sys_info_tags[k]->some_pr[EMMS2_TAG_EX_PR_INDEX]!=NULL)&&((((tag_white_list*)var_proto->channels[i]->devs[j]->tags.sys_info_tags[k]->some_pr[EMMS2_TAG_EX_PR_INDEX])->tag_type)==tag_type))
                                
                                if (var_proto->channels[i]->devs[j]->tags.sys_info_tags[k]->point_type<PROPERTY_ALARM)
                                {
                                    TAG_INFO_TO_STRING(var_proto->channels[i]->devs[j]->tags.sys_info_tags[k], lbuf);
                                }
                            }
                            for (size_t k = 0; k < var_proto->channels[i]->devs[j]->tags.ro_tags_size; k++)
                            {
                                if ((var_proto->channels[i]->devs[j]->tags.ro_tags[k]->point_type)<PROPERTY_ALARM)
                                {
                                    TAG_INFO_TO_STRING(var_proto->channels[i]->devs[j]->tags.ro_tags[k], lbuf);
                                }
                            }
                            break;
                        case EMMS2_TAG_TYPE_RW:
                            for (size_t k = 0; k < var_proto->channels[i]->devs[j]->tags.rw_tags_size; k++)
                            {
                                if (var_proto->channels[i]->devs[j]->tags.rw_tags[k]->point_type<PROPERTY_ALARM)
                                {
                                    TAG_INFO_TO_STRING(var_proto->channels[i]->devs[j]->tags.rw_tags[k], lbuf);
                                }
                            }
                            break;
                        case EMMS2_TAG_TYPE_ALARM:
                            for (size_t k = 0; k < var_proto->channels[i]->devs[j]->tags.ro_tags_size; k++)
                            {
                                if (var_proto->channels[i]->devs[j]->tags.ro_tags[k]->point_type == PROPERTY_ALARM)
                                {
                                    TAG_INFO_TO_STRING(var_proto->channels[i]->devs[j]->tags.ro_tags[k], lbuf);
                                }
                            }
                            break;
                        
                        default:
                            break;
                        }

                        if (lbuf.curlen > len_tmp)
                            lbuff_setlen(&lbuf, (int) (lbuf.curlen - 1));
                        lbuff_append(&lbuf, "}},", 3);
                    }  
                }
            }
            else
            {
                ems_syslog(LOG_INFO, "dev_name[%s]", var_proto->channels[i]->devs[j]->no);
                lbuff_sprintf(&lbuf, "{\"no\":\"%s\",\"tags\":{",var_proto->channels[i]->devs[j]->no);
                len_tmp = lbuf.curlen;

                switch (tag_type)
                {
                case EMMS2_TAG_TYPE_R:
                    get_std_enum_list_value(var, var_proto->channels[i]->devs[j]->no, &lbuf);
                    if (strstr(var_proto->channels[i]->devs[j]->no,"BMS"))
                    {
                        get_cell_data_by_bms_name(var_proto->channels[i]->devs[j]->no, &lbuf);
                    }
                    for (size_t k = 0; k < var_proto->channels[i]->devs[j]->tags.sys_info_tags_size; k++)
                    {
                        // if ((var_proto->channels[i]->devs[j]->tags.sys_info_tags[k]->some_pr[EMMS2_TAG_EX_PR_INDEX]!=NULL)&&((((tag_white_list*)var_proto->channels[i]->devs[j]->tags.sys_info_tags[k]->some_pr[EMMS2_TAG_EX_PR_INDEX])->tag_type)==tag_type))
                        
                        if (var_proto->channels[i]->devs[j]->tags.sys_info_tags[k]->point_type<PROPERTY_ALARM)
                        {
                            TAG_INFO_TO_STRING(var_proto->channels[i]->devs[j]->tags.sys_info_tags[k], lbuf);
                        }
                    }
                    for (size_t k = 0; k < var_proto->channels[i]->devs[j]->tags.ro_tags_size; k++)
                    {
                        if (var_proto->channels[i]->devs[j]->tags.ro_tags[k]->point_type<PROPERTY_ALARM)
                        {
                            TAG_INFO_TO_STRING(var_proto->channels[i]->devs[j]->tags.ro_tags[k], lbuf);
                        }
                    }
                    for (size_t k = 0; k < var_proto->channels[i]->devs[j]->tags.rw_tags_size; k++)
                    {
                        if (strcmp(var_proto->channels[i]->devs[j]->tags.rw_tags[k]->dev_tag_name,"EMS.ControlMode") == 0)
                        {
                            TAG_INFO_TO_STRING(var_proto->channels[i]->devs[j]->tags.rw_tags[k], lbuf);
                        }
                    }
                    break;
                case EMMS2_TAG_TYPE_RW:
                    for (size_t k = 0; k < var_proto->channels[i]->devs[j]->tags.rw_tags_size; k++)
                    {
                        if (var_proto->channels[i]->devs[j]->tags.rw_tags[k]->point_type<PROPERTY_ALARM)
                        {
                            TAG_INFO_TO_STRING(var_proto->channels[i]->devs[j]->tags.rw_tags[k], lbuf);
                        }
                    }
                    break;
                case EMMS2_TAG_TYPE_ALARM:
                    for (size_t k = 0; k < var_proto->channels[i]->devs[j]->tags.ro_tags_size; k++)
                    {
                        if (var_proto->channels[i]->devs[j]->tags.ro_tags[k]->point_type == PROPERTY_ALARM)
                        {
                            TAG_INFO_TO_STRING(var_proto->channels[i]->devs[j]->tags.ro_tags[k], lbuf);
                        }
                    }
                    break;
                
                default:
                    break;
                }

                if (lbuf.curlen > len_tmp)
                    lbuff_setlen(&lbuf, (int) (lbuf.curlen - 1));
                lbuff_append(&lbuf, "} },", 4);
            }
        }
    }

    if (g_usercfg_variant.EnDataCollect)
    {
        int cur_sum = get_pcs_ctrl_var()->lc_list_len;
        lc_it *lc_cfg = get_pcs_ctrl_var()->lc_list;
        for (int i = 0; i < cur_sum; i++)
        {
            if (lc_cfg[i].slave == NULL)
            {
                ems_syslog(LOG_INFO, "ems_slave is DISENBED");
                break;
            }
            if (lc_cfg[i].type == eLC_TYPE_EMS)
            {
                for (int j = 0; j < lc_cfg[i].slave->dev_sum; j++)
                {
                    const device_t *dev_ptr = slave_ems_get_dev(lc_cfg[i].slave, lc_cfg[i].slave->devs[j].no);
                    if (strcmp(dev_ptr->no, "BMS_CELLS") == 0)
                    {
                        continue;
                    }
                    if (slave_ems_get_dev_oline(lc_cfg[i].slave, dev_ptr->no) == 1)
                    {
                        size_t prefix_len = strlen(lc_cfg[i].slave->no);
                        size_t value_len = strlen(dev_ptr->no);

                        char *new_value = (char *)malloc(prefix_len + 1 + value_len + 1);

                        if (new_value)
                        {
                            strcpy(new_value, lc_cfg[i].slave->no);
                            strcat(new_value, ".");
                            strcat(new_value, dev_ptr->no);

                            lbuff_sprintf(&lbuf, "{\"no\":\"%s\",\"tags\":{\"Online\": 1,", new_value);
                            free(new_value);
                        }
                    }
                    switch (tag_type)
                    {
                    case EMMS2_TAG_TYPE_R:
                        get_std_enum_list_value(var, (char *)dev_ptr->no, &lbuf);
                        if (strstr(dev_ptr->no, "BMS"))
                        {
                            get_ems_slave_cell_data_by_bms_name((char *)dev_ptr->no, &lbuf);
                        }
                        for (int k = 0; k < dev_ptr->tags.ro_tags_size; k++)
                        {
                            DATA_INFO_TO_STRING(*(dev_ptr->tags.ro_tags[k]), lbuf);
                        }
                        break;
                    case EMMS2_TAG_TYPE_RW:
                        break;
                    default:
                        break;
                    }

                    if (lbuf.curlen > len_tmp)
                        lbuff_setlen(&lbuf, (int)(lbuf.curlen - 1));
                    lbuff_append(&lbuf, "} },", 4);
                }
            }
        }
    }

    if (lbuf.curlen > inlen)
        lbuff_setlen(&lbuf, (int) (lbuf.curlen - 1));
    lbuff_append(&lbuf, "]", 1);

    const char *function = NULL;
    switch (tag_type)
    {
    case EMMS2_TAG_TYPE_R:
        function = FUNCTION_SUBTELEMETRY;
        break;
    case EMMS2_TAG_TYPE_RW:
        function = FUNCTION_SUBTELECONTROL;

        break;
    default:
        break;
    }

    if (*data_all == NULL)
    {
        *data_all=calloc(1, lbuf.curlen + 512);
    }
    if (*pub_topic == NULL)
    {
        *pub_topic=calloc(1,512);
    }

    sprintf(*pub_topic,"emms2/%s/%s/%s",type,var->topic_sn,function);
    package_based_field(var,*data_all,(char *)function,seq, lbuf.bufptr);
    lbuff_free(&lbuf);
    ems_syslog(LOG_NOTICE, "data_all:::package topic [%s] data:[%s]",*pub_topic,*data_all);
}

static int report_dev_data_by_type_and_name(mqtt_emms2_var_t *var,const char *type,char** dev_name,int dev_num,EMMS2_TAG_TYPE tag_type,int seq) //dev_num 为0上报所有数采设备
{
    
    int ret = 0;
    char *data_all=NULL,*pub_topic = NULL;
    package_dev_data_by_type_and_name(var,type,dev_name,dev_num,tag_type,seq,&data_all,&pub_topic);
    if((ret = mqtt_publish_message(var,pub_topic,data_all,strlen(data_all)))!=0)
        ems_syslog(LOG_WARNING,"send data error!!");

    if (pub_topic!=NULL)
    {
        free(pub_topic);
    }
    if (data_all!=NULL)
    {
        free(data_all);
    }
    return ret;
}
static int saving_publish_error_data(mqtt_emms2_var_t* var, char* data_all);
static int report_all_dev_data_with_history(mqtt_emms2_var_t* var)
{
    int ret = 0;
    char *data_all=NULL,*pub_topic = NULL;
    package_dev_data_by_type_and_name(var,TYPE_POST,NULL,0,EMMS2_TAG_TYPE_R,0,&data_all,&pub_topic);
    if (var->emms2_status<EMMS2_HEARTBEAT)
    {
        ret = -1;
        goto CHECK_ERR;
    }
    ems_syslog(LOG_INFO,"report dev data!");
    if((ret = mqtt_publish_message(var,pub_topic,data_all,strlen(data_all)))!=0)
        ems_syslog(LOG_WARNING,"send data error!!");
CHECK_ERR:
    // var->seq_dev_data_post = var->seq;
    if( var->base_param.saveDisconnectData &&((ret<0) || (cond_wait(&var->cond,&var->mutex,var->cond_wait_sec,0)!=0) )) //超时存数据
    {
        ems_syslog(LOG_NOTICE, "publish data error saving !!");
        saving_publish_error_data(var, data_all);
    }

    if (pub_topic!=NULL)
    {
        free(pub_topic);
    }
    if (data_all!=NULL)
    {
        free(data_all);
    }
    return ret;
}

//这段数据data_all传入数据是压缩/处理后数据. 在调用此函数定义地方,注明内部会对数据做处理. 
static int saving_publish_error_data(mqtt_emms2_var_t* var, char* data_all)
{
    if (var->offline_db == NULL)
        init_offline_db_table(var, var->station_id);
    if (var->offline_db == NULL) {
        ems_syslog(LOG_ERR, "saving data error, init_offline_db_table error");
        return -1;
    }

    int ret = -1;
    char* pub_topic = NULL;
    char* data_sql  = NULL;
    char* func_i    = strstr(data_all, FUNCTION_SUBTELEMETRY);

    if (func_i == NULL) // !! 数据格式要求一致, 直接替换function name; FUNCTION_CONTINUE_DATA长度必须为12
        goto ERR0;

    memcpy(func_i, FUNCTION_CONTINUE_DATA, strlen(FUNCTION_CONTINUE_DATA));
    int pub_topic_len = 512;
    pub_topic = calloc(1, pub_topic_len);
    if (pub_topic == NULL)
    {
        ems_syslog(LOG_ERR, "saving data error, malloc pub_topic error");
        goto ERR0;
    }
    snprintf(pub_topic, pub_topic_len, "emms2/%s/%s/%s", TYPE_POST, var->topic_sn, FUNCTION_CONTINUE_DATA);

    data_sql = calloc(1, strlen(pub_topic) + strlen(data_all) + 256);
    if (data_sql == NULL)
    {
        ems_syslog(LOG_ERR, "saving data error, malloc data_sql error");
        goto ERR0;
    }
    sprintf(data_sql, "INSERT INTO '" EMMS2_TAGBLE_NAME "' VALUES(%ld,'%s','%s');", time(NULL), pub_topic, data_all);

    ret = sqlite3_exec_command(var->offline_db->db, data_sql);
    if ( ret != 0)
    {
        ems_syslog(LOG_ERR, "saving data error, sqlite3_exec_command error");
    }

ERR0:
    if(data_sql)
        free(data_sql);
    if(pub_topic)
        free(pub_topic);
    return ret;
}

static int report_all_dev_history_data(mqtt_emms2_var_t *var)
{
    if (var->offline_db == NULL)
        init_offline_db_table(var, var->station_id);
    if (var->offline_db == NULL) {
        ems_syslog(LOG_ERR, "report_all_dev_history_data error, init_offline_db_table error");
        return -1;
    }

    int ret = 0;
    char *data_all=NULL,*pub_topic = NULL;

    int nrow = 0;
    char *zErrMsg =NULL;
	char **azResult=NULL;
    sqlite3_get_table(var->offline_db->db, "SELECT topic,payload FROM '" EMMS2_TAGBLE_NAME "' ORDER BY ID ASC LIMIT 1;", &azResult, &nrow, NULL, &zErrMsg);
    if(zErrMsg)
	{
		printf("getHistoryData_backend_sqlite error:%s\n",zErrMsg);
		goto END;
	}
    if (nrow>0)
    {
        pub_topic = azResult[2];
        data_all = azResult[3];
        ems_syslog(LOG_DEBUG,"sqlite package topic [%s] data:[%s]",pub_topic,data_all);
    }
    else
    {
        ems_syslog(LOG_INFO, "sqlite history data is null");
        ret = 0;
        goto END;
    }
    ems_syslog(LOG_DEBUG,"report history data!");
    if((ret = mqtt_publish_message(var,pub_topic,data_all,strlen(data_all)))!=0)
        ems_syslog(LOG_WARNING,"send data error!!");
    
    // var->seq_dev_data_post = var->seq;
    if ((ret == 0) && (cond_wait(&var->offline_db->history_cond, &var->offline_db->history_mutex, var->cond_wait_sec, 0) == 0)) //完成则删除数据
    {
        ems_syslog(LOG_DEBUG,"send history success delete current data!!");
        sqlite3_exec_command(var->offline_db->db, "DELETE FROM '" EMMS2_TAGBLE_NAME "' WHERE ID IN (SELECT ID FROM '" EMMS2_TAGBLE_NAME "' ORDER BY ID ASC LIMIT 1);");
    }
    else
        ret = -1;

END:
    if (azResult)
        sqlite3_free_table(azResult);
    if (zErrMsg)
        sqlite3_free(zErrMsg);
    return ret;
}

int package_ems_data_by_type(mqtt_emms2_var_t *var,const char *type,EMMS2_TAG_TYPE tag_type,int seq,char **data_all,char **pub_topic)
{
    int ret = 0;
    struct lnxall_buff lbuf;
    unsigned int inlen;
    
    lbuff_init(&lbuf, 4096);
    lbuff_sprintf(&lbuf, "\"tags\":{");
    inlen = lbuf.curlen;
    
    proto_forward_t *var_proto = get_proto_forward_var();
    for (size_t i = 0; i < var_proto->channels_size; i++)
    {
        for (size_t j = 0; j < var_proto->channels[i]->devs_size; j++)
        {
            if (strcmp(var_proto->channels[i]->devs[j]->no,DEV_NO_EMS)==0)
            {
                switch (tag_type)
                {
                case EMMS2_TAG_TYPE_R:
                    for (size_t k = 0; k < var_proto->channels[i]->devs[j]->tags.sys_info_tags_size; k++)
                    {
                        // if ((var_proto->channels[i]->devs[j]->tags.sys_info_tags[k]->some_pr[EMMS2_TAG_EX_PR_INDEX]!=NULL)&&((((tag_white_list*)var_proto->channels[i]->devs[j]->tags.sys_info_tags[k]->some_pr[EMMS2_TAG_EX_PR_INDEX])->tag_type)==tag_type))
                        
                        if (var_proto->channels[i]->devs[j]->tags.sys_info_tags[k]->point_type<PROPERTY_ALARM)
                        {
                            TAG_INFO_TO_STRING(var_proto->channels[i]->devs[j]->tags.sys_info_tags[k], lbuf);
                        }
                    }
                    for (size_t k = 0; k < var_proto->channels[i]->devs[j]->tags.ro_tags_size; k++)
                    {
                        if (var_proto->channels[i]->devs[j]->tags.ro_tags[k]->point_type<PROPERTY_ALARM)
                        {
                            TAG_INFO_TO_STRING(var_proto->channels[i]->devs[j]->tags.ro_tags[k], lbuf);
                        }
                    }
                    for (size_t k = 0; k < var_proto->channels[i]->devs[j]->tags.rw_tags_size; k++)
                    {
                        if (strcmp(var_proto->channels[i]->devs[j]->tags.rw_tags[k]->dev_tag_name,"EMS.ControlMode") == 0)
                        {
                            TAG_INFO_TO_STRING(var_proto->channels[i]->devs[j]->tags.rw_tags[k], lbuf);
                        }
                    }
                    break;
                case EMMS2_TAG_TYPE_RW:
                    for (size_t k = 0; k < var_proto->channels[i]->devs[j]->tags.rw_tags_size; k++)
                    {
                        if (var_proto->channels[i]->devs[j]->tags.rw_tags[k]->point_type<PROPERTY_ALARM)
                        {
                            TAG_INFO_TO_STRING(var_proto->channels[i]->devs[j]->tags.rw_tags[k], lbuf);
                        }
                    }
                    break;
                case EMMS2_TAG_TYPE_ALARM:
                    for (size_t k = 0; k < var_proto->channels[i]->devs[j]->tags.ro_tags_size; k++)
                    {
                        if (var_proto->channels[i]->devs[j]->tags.ro_tags[k]->point_type == PROPERTY_ALARM)
                        {
                            TAG_INFO_TO_STRING(var_proto->channels[i]->devs[j]->tags.ro_tags[k], lbuf);
                        }
                    }
                    break;
                
                default:
                    break;
                }
            }
        }
    }

    if (lbuf.curlen > inlen)
        lbuff_setlen(&lbuf, (int) (lbuf.curlen - 1));
    lbuff_append(&lbuf, "}", 1);

    const char *function = NULL;
    switch (tag_type)
    {
    case EMMS2_TAG_TYPE_R:
        function = FUNCTION_TELEMETRY;
        break;
    case EMMS2_TAG_TYPE_RW:
        function = FUNCTION_TELECONTROL;
        break;
    default:
        break;
    }
    if (*data_all == NULL)
    {
        *data_all=calloc(1, lbuf.curlen + 512);
    }
    if (*pub_topic == NULL)
    {
        *pub_topic=calloc(1,512);
    }
    package_based_field(var,*data_all,(char *)function,seq, lbuf.bufptr);
    sprintf(*pub_topic,"emms2/%s/%s/%s",type,var->topic_sn,function);
    ems_syslog(LOG_INFO,"package topic [%s] data:[%s]",*pub_topic,*data_all);
    lbuff_free(&lbuf);
    return ret;
}
static int report_all_ems_data_by_type(mqtt_emms2_var_t *var,const char *type,EMMS2_TAG_TYPE tag_type,int seq)
{
    int ret = 0;
    char *data_all=NULL,*pub_topic = NULL;
    
    if (var->emms2_status<EMMS2_HEARTBEAT)
    {
        return -1;
    }
    ems_syslog(LOG_INFO,"report emms data!");
    package_ems_data_by_type(var,type,tag_type,seq,&data_all,&pub_topic);
    if((ret = mqtt_publish_message(var,pub_topic,data_all,strlen(data_all)))!=0)
        ems_syslog(LOG_WARNING,"send data error!!");

    if (pub_topic!=NULL)
    {
        free(pub_topic);
    }
    if (data_all!=NULL)
    {
        free(data_all);
    }
    return ret;
}



static void report_all_data(mqtt_emms2_var_t *var)
{
    ems_syslog(LOG_WARNING,"report all data!");
    report_all_ems_data_by_type(var,TYPE_POST,EMMS2_TAG_TYPE_R,0);
    report_dev_data_by_type_and_name(var,TYPE_POST,NULL,0,EMMS2_TAG_TYPE_R,0);
}

#if 0
void check_change_report(mqtt_emms2_var_t *var)
{
    tag_t *tag_return = NULL;
    change_list_t *change_info = NULL;
    pthread_rwlock_wrlock(&var->dev_list_rwlock);
    while(!list_empty(&var->change_report_list))
    {
        change_info = list_first_entry(&var->change_report_list, typeof(*change_info), list);
        tag_return = change_info->w_tag;
        list_del(&change_info->list);
        free(change_info);
    }
    pthread_rwlock_unlock(&var->dev_list_rwlock);
}
#else

typedef struct 
{
    // mqtt_emms2_var_t *var;
    char *(*buff)[MAX_MQTT_EMMS2_DEV_NUM];
    int num;
    char report_ems;
}hash_context_t;

static int iter_dev(struct hash_intptr *hash_p, int unknow, void *param)
{
    hash_context_t *hash_context = param;
    char *needreport = hash_p->hh_ptr;
    char *key=(char *)hash_p->hh_key;
    if (*needreport == 1)
    {
        ems_syslog(LOG_INFO,"take dev [%s] !",key);
        if (strcmp(key,DEV_NO_EMS)==0)
        {
            hash_context->report_ems = 1;
        }else
        {
            (*hash_context->buff)[hash_context->num++] = key;
        } 
        *needreport=0;       
    }
    
    return 0;
}

static int check_change_report(mqtt_emms2_var_t *var)
{

    char *dev_list[MAX_MQTT_EMMS2_DEV_NUM] = {NULL};
    hash_context_t hash_context = {.buff = &dev_list,.num=0,.report_ems=0};
    hash_intptr_iter(&var->dev_hash,iter_dev,&hash_context);
    int reported=-1;
    if (hash_context.report_ems==1)
    {
        report_all_ems_data_by_type(var,TYPE_POST,EMMS2_TAG_TYPE_R,0);
        reported = 0;
    }
    if (hash_context.num>0)
    {
        report_dev_data_by_type_and_name(var,TYPE_POST,dev_list,hash_context.num,EMMS2_TAG_TYPE_R,0);
        reported = 0;
    }
    return reported;    
}
#endif

static int load_changetag_config(mqtt_emms2_var_t *var)
{
    char*  file = read_file_data(MQTT_CHANGE_REPORT_LIST);
    if (file == NULL)
    {
        var->tag_config_list = malloc(2 * sizeof(tag_config_t));
        if (var->tag_config_list == NULL)
        {
            ems_syslog(LOG_ERR, "内存分配失败");
            return -1;
        }
        
        strcpy(var->tag_config_list[0].dev_name, "PCS");
        strcpy(var->tag_config_list[0].type_name, "PCS");
        strcpy(var->tag_config_list[0].tag_name, FUNCTION_GET_CONNECT);
        var->tag_config_list[0].func = 1;
        
        strcpy(var->tag_config_list[1].dev_name, "BMS");
        strcpy(var->tag_config_list[1].type_name, "BMS");
        strcpy(var->tag_config_list[1].tag_name, "RackCagDsgState");
        var->tag_config_list[1].func = 0;
        var->tag_config_count = 2;
        ems_syslog(LOG_INFO, "使用默认配置点位, %d个配置项", var->tag_config_count);
        return 0;
    }
    
    cJSON *root = cJSON_Parse(file);
    if (root == NULL)
    {
        ems_syslog(LOG_ERR, "JSON解析失败");
        free(file);
        return -1;
    }
    
    // 获取change_tag_list数组
    cJSON *change_tag_list = cJSON_GetObjectItem(root, "change_tag_list");
    if (change_tag_list == NULL || !cJSON_IsArray(change_tag_list))
    {
        ems_syslog(LOG_ERR, "JSON格式错误:缺少change_tag_list数组");
        free(file);
        cJSON_Delete(root);
        return -1;
    }
    
    // 计算总配置项数量
    int total_config_count = 0;
    cJSON *device_item = NULL;
    cJSON_ArrayForEach(device_item, change_tag_list)
    {
        cJSON *tags_array = cJSON_GetObjectItem(device_item, "tags");
        if (tags_array != NULL && cJSON_IsArray(tags_array))
        {
            total_config_count += cJSON_GetArraySize(tags_array);
        }
    }
    
    if (total_config_count == 0)
    {
        ems_syslog(LOG_ERR, "没有找到有效的配置项");
        free(file);
        cJSON_Delete(root);
        return -1;
    }
    
    // 分配内存
    var->tag_config_list = malloc(total_config_count * sizeof(tag_config_t));
    if (var->tag_config_list == NULL)
    {
        ems_syslog(LOG_ERR, "内存分配失败");
        free(file);
        cJSON_Delete(root);
        return -1;
    }
    
    var->tag_config_count = 0;
    cJSON_ArrayForEach(device_item, change_tag_list)
    {
        cJSON *type_obj = cJSON_GetObjectItem(device_item, "type");
        cJSON *no_obj = cJSON_GetObjectItem(device_item, "no");
        cJSON *tags_array = cJSON_GetObjectItem(device_item, "tags");
        
        if (type_obj == NULL || no_obj == NULL || tags_array == NULL || 
            !cJSON_IsString(type_obj) || !cJSON_IsString(no_obj) || !cJSON_IsArray(tags_array))
        {
            continue;
        }
        
        cJSON *tag_item = NULL;
        cJSON_ArrayForEach(tag_item, tags_array)
        {
            if (cJSON_IsString(tag_item))
            {
                snprintf(var->tag_config_list[var->tag_config_count].dev_name,sizeof(var->tag_config_list[var->tag_config_count].dev_name),"%s",no_obj->valuestring);
                snprintf(var->tag_config_list[var->tag_config_count].type_name,sizeof(var->tag_config_list[var->tag_config_count].type_name),"%s",type_obj->valuestring);
                snprintf(var->tag_config_list[var->tag_config_count].tag_name,sizeof(var->tag_config_list[var->tag_config_count].tag_name),"%s",tag_item->valuestring);
                var->tag_config_list[var->tag_config_count].func = 0;
                var->tag_config_count++;
            }
        }
    }
    
    free(file);
    cJSON_Delete(root);
    
    ems_syslog(LOG_INFO, "成功加载点位配置文件，共%d个配置项", var->tag_config_count);
    return 0;
}

static void init_change_report_list(mqtt_emms2_var_t* var)
{
    proto_forward_t *var_proto = get_proto_forward_var();
    
    // 计算匹配的设备数量
    int match_count = 0;
    for (size_t i = 0; i < var_proto->channels_size; i++)
    {
        for (size_t j = 0; j < var_proto->channels[i]->devs_size; j++)
        {
            device_t *dev = var_proto->channels[i]->devs[j];
            for(int k = 0; k < var->tag_config_count; k++)
            {
                if((strcmp(dev->dev_type,var->tag_config_list[k].type_name) == 0) && (strstr(dev->no,var->tag_config_list[k].dev_name)))
                {
                    match_count++;
                }
            }
        }
    }
    
    if (match_count == 0)
    {
        ems_syslog(LOG_INFO, "没有找到匹配的设备配置项");
        return;
    }
    
    var->device_status_list = malloc(match_count * sizeof(device_status_t));
    if (var->device_status_list == NULL)
    {
        ems_syslog(LOG_ERR, "device_status_list内存分配失败");
        return;
    }
    
    var->device_status_count = 0;
    
    for (size_t i = 0; i < var_proto->channels_size; i++)
    {
        for (size_t j = 0; j < var_proto->channels[i]->devs_size; j++)
        {
            device_t *dev = var_proto->channels[i]->devs[j];
            //遍历当前设备表，使用strstr找出包含配置文件中设备名称的设备，保存到device_status_list中
            for(int k = 0; k < var->tag_config_count; k++)
            {
                if((strcmp(dev->dev_type,var->tag_config_list[k].type_name) == 0) && (strstr(dev->no,var->tag_config_list[k].dev_name)))
                {
                    snprintf(var->device_status_list[var->device_status_count].dev_no,sizeof(var->device_status_list[var->device_status_count].dev_no),"%s",dev->no);
                    snprintf(var->device_status_list[var->device_status_count].tag_name,sizeof(var->device_status_list[var->device_status_count].tag_name),"%s",var->tag_config_list[k].tag_name);
                    var->device_status_list[var->device_status_count].func = var->tag_config_list[k].func;
                    // 初始化last_status为-1（表示未获取过状态）
                    if(var->device_status_list[var->device_status_count].func == 1)
                        dev_get_function_int(dev->no, FUNCTION_GET_CONNECT, &var->device_status_list[var->device_status_count].last_status);
                    else
                        var->device_status_list[var->device_status_count].last_status = dev_get_dev_tag_int(var->device_status_list[var->device_status_count].dev_no,var->device_status_list[var->device_status_count].tag_name);
                    
                    ems_syslog(LOG_DEBUG, "添加设备状态监控: 设备编号=%s, 点位名称=%s, 是否使用适配器=%d,初始值=%d", 
                              var->device_status_list[var->device_status_count].dev_no,
                              var->device_status_list[var->device_status_count].tag_name,
                              var->device_status_list[var->device_status_count].func,
                              var->device_status_list[var->device_status_count].last_status);
                    
                    var->device_status_count++;
                }
            }
        }
    }
    
    ems_syslog(LOG_INFO, "成功初始化设备状态列表，共%d个监控项", var->device_status_count);
}

static int check_tag_status_change(mqtt_emms2_var_t *var)
{
    int status_changed = 0;
    if (var->device_status_list == NULL || var->device_status_count == 0)
    {
        return 0;
    }
    
    // 遍历设备状态列表，检查每个监控点位的状态变化
    for (int i = 0; i < var->device_status_count; i++)
    {
        int current_status = 0;
        device_status_t *status_item = &var->device_status_list[i];
        if(status_item->func == 1)
        {
            dev_get_function_int(status_item->dev_no, FUNCTION_GET_CONNECT, &current_status);
            ems_syslog(LOG_INFO,"适配器获取value:%d",current_status);
        }
        else
            current_status = dev_get_dev_tag_int(status_item->dev_no, status_item->tag_name);
        if (status_item->last_status != current_status)
        {
            ems_syslog(LOG_INFO, "点位状态变化: 设备编号=%s, 点位名称=%s, 旧状态=%d, 新状态=%d", 
                      status_item->dev_no, status_item->tag_name, 
                      status_item->last_status, current_status);
            status_item->last_status = current_status;
            status_changed = 1;
        }
    }
    
    return status_changed;
}


static int mqtt_state_change(void *obj, int state)
{
    ems_syslog(LOG_ERR, "mqtt state change to %d", state);
    BUSINESS_LOG(BLOG_NOTICE, NORTH_ID_0051, ((mqtt_emms2_var_t *)obj)->station_id, "[北向]:[云/站控]: [%s] mqtt state change to [%d] [0:正在连接;1:已连接;2:断开连接]", ((mqtt_emms2_var_t *)obj)->station_id, state);
    if (state != MQTT_CONNECTED)
    {
        server_add_new_log(LOG_EMMS2_NET,"EMS",LOG_MODE_PERIOD,1,NULL);
        ((mqtt_emms2_var_t *)obj)->emms2_status = EMMS2_CONNECTING;
    }
    else
    {
        server_add_new_log(LOG_EMMS2_NET,"EMS",LOG_MODE_PERIOD,2,NULL);
        ((mqtt_emms2_var_t *)obj)->emms2_status = EMMS2_LOGINING;
    }
    return 0;
}

static int sprint_log_info(ONE_LOG_INFO *log_info,void *contex)
{
    char data_tmp[512]={0};
    snprintf(data_tmp,512,"{\"logCode\":%d,\"entityInstance\":\"%s\",\"value\":%d,\"additionalParam\":\"%s\",\"raiseTime\":%ld},",log_info->log_code,log_info->no,log_info->value,log_info->add,log_info->raiseTime);
    strcat((*(char **)contex),data_tmp);
    return 0;
}

static int report_log_info(mqtt_emms2_var_t *var)
{
    char *data_info=calloc(1,1024*10*10);

    sprintf(data_info,"\"messages\":[");

    client_get_log_with_offset(var->log_fd,0,10,sprint_log_info,&data_info);
    
    if (strlen(data_info)>strlen("\"messages\":["))
    {
        data_info[strlen(data_info)-1] = '\0';
    }
    
    strcat(data_info,"]");
    char *data_all = calloc(1,strlen(data_info)+256);
    package_based_field(var,data_all,FUNCTION_LOG,0,data_info);
    char pub_topic[512] = {0};
    sprintf(pub_topic,"emms2/"TYPE_POST"/%s/"FUNCTION_LOG,var->topic_sn);
    ems_syslog(LOG_INFO,"pub topic [%s] data:[%s]",pub_topic,data_all);
    int ret = mqtt_publish_message(var,pub_topic,data_all,strlen(data_all));
    free(data_info);
    free(data_all);
    return ret;
}

void package_pub_data(mqtt_emms2_var_t *var,char **data_all,char **pub_topic)
{
    struct lnxall_buff lbuf;
    unsigned int inlen, len_tmp = 0;
    
    lbuff_init(&lbuf, 4096);
    lbuff_sprintf(&lbuf, "\"messages\":[");
    inlen = lbuf.curlen;

    for (size_t i = 0; i < var->pub_dev->pub_dev_list_num; i++)
    {
        if (0 == strcmp("EMS", var->pub_dev->pub_dev_list[i].device->no))
        {
            continue;   
        }
        lbuff_sprintf(&lbuf, "{\"no\":\"%s\",\"tags\":{", var->pub_dev->pub_dev_list[i].device->no);
        get_std_enum_list_value(var, var->pub_dev->pub_dev_list[i].device->no, &lbuf);
        len_tmp = lbuf.curlen;
        for (size_t j = 0; j < var->pub_dev->pub_dev_list[i].pub_tag_num; j++)
        {
            TAG_INFO_TO_STRING(var->pub_dev->pub_dev_list[i].tag_pub_list[j], lbuf);
        }
        lbuff_sprintf(&lbuf, "\"AlarmCnt\":%d,", dev_get_dev_tag_int(var->pub_dev->pub_dev_list[i].device->no, ALARMCNT)); //alarmCnt没有注册数采点位，单独加入
        if (lbuf.curlen > len_tmp)
            lbuff_setlen(&lbuf, (int) (lbuf.curlen - 1));
        lbuff_append(&lbuf, "}},", 3);
    }

    if (lbuf.curlen > inlen)
        lbuff_setlen(&lbuf, (int) (lbuf.curlen - 1));
    lbuff_append(&lbuf, "]", 1);

    if (*data_all == NULL)
    {
        *data_all=calloc(1, lbuf.curlen + 512);
    }
    if (*pub_topic == NULL)
    {
        *pub_topic=calloc(1,512);
    }

    sprintf(*pub_topic,"emms2/%s/%s/%s",TYPE_POST,var->topic_sn,FUNCTION_SUB_DATA);
    package_based_field(var,*data_all,(char *)FUNCTION_SUB_DATA,0, lbuf.bufptr);
    ems_syslog(LOG_INFO,"package topic [%s] data:[%s]",*pub_topic,*data_all);
    lbuff_free(&lbuf);
}

static int report_pub_data(mqtt_emms2_var_t *var) //上报所有的被订阅的数据
{
    
    int ret = 0;
    char *data_all=NULL,*pub_topic = NULL;
    package_pub_data(var,&data_all,&pub_topic);
    if((ret = mqtt_publish_message(var,pub_topic,data_all,strlen(data_all)))!=0)
        ems_syslog(LOG_WARNING,"send data error!!");

    if (pub_topic!=NULL)
    {
        free(pub_topic);
    }
    if (data_all!=NULL)
    {
        free(data_all);
    }
    return ret;
}
static void *mqtt_emms2_cycle_task(void *data)
{
    int retry = 0;
    mqtt_emms2_var_t *var = data;
    pcs_ctrl_var_t * pcs_var = get_pcs_ctrl_var();
    // time_t report_all_time = 0; //记录上次数据上报时间
    time_t report_heartbeat_time = 0; //记录上次心跳时间
    time_t report_log_time = 0;
    time_t report_alarm_time = 0;
    time_t report_ems_sn_time = 0;
    time_t now;

#if EMMS2_ALARM_BY_CB
    int reflush_ems = 0;
#else
    time_t check_alarm_time = 0;
    unsigned int time_today = 0,last_time_report = 0, last_time_try = 0;
    int reported = 0;
    time_t check_pubTime = 0;
    // 断点续传退避
    time_t last_report_history_time   = 0;
    int    fail_report_history_times  = 0;

    int deleted_data = 1;//开机删除一次 每天12点删除3天前的数据
    sleep(20);//延时启动等一会儿数据刷新
    char sql_del[256] = {0};
#endif
    while (1)
    {
        sleep(1);
        now = time(NULL);
        time_today = get_seconds_of_today();

		int diff = time_today - last_time_report;
        if ((diff < 0) || (diff > SECONDS_IN_DAY))
        {
            reported = 0;
            last_time_report = 0;
        }
        

        unsigned int time_elapsed = time_today + (time_today >= last_time_try ? 0 : SECONDS_IN_DAY) - last_time_try;
        bool can_retry = time_elapsed >= backoff_intervals[retry];
        switch (var->emms2_status)
        {
            case EMMS2_CONNECTING:
                break;
            case EMMS2_LOGINING:
            {
                BUSINESS_LOG(BLOG_NOTICE, NORTH_ID_0050, var->station_id, "[北向]:[云/站控]:%s 状态为EMMS2_LOGINING 登陆中状态", var->station_id);
                if(can_retry)
                {
                    server_add_new_log(LOG_EMMS2_LOGIN,"EMS",LOG_MODE_PERIOD,1,NULL);
                    mqtt_emms2_login(var);
                    ems_syslog(LOG_CRIT, "[云/站控]:%s暂未回复登陆请求,数据无法上报", var->station_id);

                    last_time_try = time_today;
                    if(retry < var->backoff_count) retry++;
                    ems_syslog(LOG_ERR, "next retry in %d seconds", backoff_intervals[retry]);
                }
                // var->emms2_status = EMMS2_DEVREPORTING; //测试用正式需要删掉
            }
                break;
            case EMMS2_DEVREPORTING:
            {
                BUSINESS_LOG(BLOG_NOTICE, NORTH_ID_0050, var->station_id, "[北向]:[云/站控]:%s 状态为EMMS2_DEVREPORTING 设备数据上报状态", var->station_id);
                retry = 0;
                var->offlinecount = 0;
                last_time_try = 0;
                server_add_new_log(LOG_EMMS2_LOGIN,"EMS",LOG_MODE_PERIOD,2,NULL);
                ems_syslog(LOG_WARNING,"mqtt device info reporting ... !");
                ems_syslog(LOG_INFO, "MQTT start upload templates and alarm, station %s status DEVREPORTING", var->station_id);
                check_and_upload_templates_and_alarm(var);
                mqtt_emms2_devInfo(var);
                report_all_alarm(var,0);
                send_sync_cabinet_info(var);
                // report_all_time = 0;
                report_heartbeat_time = 0;

                var->emms2_status = EMMS2_HEARTBEAT;
                // continue;
            }
                break;
            case EMMS2_HEARTBEAT:
            {
                
//2024-1-13 为了使告警在任何本地上报 移动到下面
/*
#if EMMS2_ALARM_BY_CB
                if (reflush_ems == 0)
                {
                    reflush_ems = 1;
                    reflush_all_ems_data();
                }
#else
                if (check_timeout_second(now, &check_alarm_time, 1)) 
                {
                    check_alarm_info(var);
                }
#endif
*/
                BUSINESS_LOG(BLOG_NOTICE, NORTH_ID_0050, var->station_id, "[北向]:[云/站控]:%s 状态为EMMS2_HEARTBEAT 心跳正常状态", var->station_id);
                if (check_timeout_second(now, &check_alarm_time, 1))
                {
                    check_alarm_info(var);
                }
                
                if (check_timeout_second(now, &var->last_check_time, var->base_param.change_deadband))
                {
                    if (check_tag_status_change(var))
                    {
                        ems_syslog(LOG_INFO, "检测到状态变化，触发数据上报");
                        report_all_data(var);
                    }
                }
                
                if ((var->base_param.changeReportInterval!=0) && (time_today%var->base_param.changeReportInterval == 0)) 
                {
                    if(check_change_report(var)==0);
                        //report_heartbeat_time = now;
                }
                if (check_timeout_second(now, &report_heartbeat_time, HEARTBEAT_INTERVAL)) 
                {
                    mqtt_emms2_heartBeat(var);
                    report_heartbeat_time = now;
                    //  continue;
                }
                if(check_timeout_second(now, &report_alarm_time, HEARTBEAT_INTERVAL))
                {
                    report_all_alarm(var,0);
                }
                if (check_timeout_second(now, &report_log_time, 60*3)||(pcs_var->restart_code)||(pcs_var->reboot_code == REBOOT_VALUE)) 
                {
                    if(report_log_info(var)==0);
                        //report_heartbeat_time = now;
                }
                if (check_timeout_second(now, &report_ems_sn_time, 300))
                {
                    check_ems_num_change(var);
                }
                
                if (var->pub_dev != NULL)
                {
                    if (now <= var->pub_dev->end) //超过超时时间停止
                    {
                        if ((var->pub_dev->period != 0) && check_timeout_second(now, &check_pubTime, var->pub_dev->period))
                        {
                            report_pub_data(var);
                        }
                    }
                }
                if ((!var->base_param.saveDisconnectData)&&(var->base_param.reportInterval!=0) && ((reported == 0)||(time_today - last_time_report)>=var->base_param.reportInterval))
                {
                    reported = 1;
                    last_time_report = time_today - (time_today%var->base_param.reportInterval);
                    report_all_data(var);
                }

                if (var->base_param.saveDisconnectData && var->offlinecount>3) //三次心跳不回复自动离线
                {
                    var->emms2_status = EMMS2_LOGINING;
                    set_alarm_need_report(false);
                }
                
                if (get_cabinet_info()->is_slave == 0)
                {
                    
                    mqtt_dp_lc_info_task(var->dp, now);
                    mqtt_dp_config_task(var->dp, now);
                }
                if(0 == strcmp(var->station_id, MQTT_EMMS2_STATION_ID ))  // 暂时对云平台支持对站控不支持
                {
#ifdef EN_ACCIDENT
                    if (accident_upload_control.send_accident_create_upload_flag)
                    {
                        if (accident_upload_control.send_accident_create_upload_flag)
                        {
                            if(accident_info_report_upload(var, FUNCTION_ACCIDENT_CREATE, accident_upload_control.send_md5) < 0)
                            {
                                ems_syslog(LOG_ERR, "[zxf]accident info report upload failed");
                            }
                            accident_upload_control.send_accident_create_upload_flag = false;
                        }
                    }
#endif
                }
                mqtt_hardware_monitor_task(var->dp, now);
            }
                break;
            default:
                break;
        }

        if (var->base_param.saveDisconnectData)
        {
            if ((var->base_param.reportInterval!=0) && ((reported == 0)||(time_today - last_time_report)>=var->base_param.reportInterval))
            {
                reported = 1;
                last_time_report = time_today - (time_today%var->base_param.reportInterval);
                
                ems_syslog(LOG_INFO,"report all data with history ... !");
                report_all_ems_data_by_type(var,TYPE_POST,EMMS2_TAG_TYPE_R,0); //ems数据不保存
                report_all_dev_data_with_history(var); //设备数据保存和等待 超时时间10秒  最低存储间隔为10~15秒
            }
            else
            {
                if (var->emms2_status >= EMMS2_HEARTBEAT && now - last_report_history_time >= backoff_intervals[fail_report_history_times]) //断点续传
                {
                    last_report_history_time = now;
                    if (0 == report_all_dev_history_data(var))
                    {
                        fail_report_history_times = 0;
                    }
                    else
                    {
                        if (fail_report_history_times < var->backoff_count)
                            ++fail_report_history_times;
                    }
                }
            }
            if (time_today < 10) 
            {
                if (deleted_data == 1)
                {
                    deleted_data = 0;
                    ems_syslog(LOG_DEBUG,"delete history data!");
                    if (var->offline_db != NULL)
                    {
                        snprintf(sql_del, sizeof(sql_del), "DELETE FROM " EMMS2_TAGBLE_NAME " WHTER ID < %ld;", time(NULL) - SECONDS_IN_DAY * HISTORY_DATA_HOLD_DAYS); //只保留30天的数据
                        sqlite3_exec_command(var->offline_db->db, sql_del);
                    }
                }
            }
            else
            {
                deleted_data = 1;
            }
        }
    }
    return NULL;
}

// static int load_emms2_cfg(void)
// {
//     int i, num;
//     dev_white_list * wlist;
//     ems_syslog(LOG_NOTICE,"MQTT cloud CLIENT START!!");
// 
//     // initialize `white_list structure array
//     num = dev_white_list_num;
//     for (i = 0; i < num; ++i) {
//         wlist = &alarm_define_list[i];
// 
//         if (wlist->deviceType == NULL)
//             continue;
//         if (strcmp(wlist->deviceType, DEV_BMS) == 0) {
//             const char * config = CONFIG_PATH "/" EMMS2_VENDOR_MODEL_BMS;
//             ems_syslog(LOG_ERR, "Loading emms2_vendor_model from %s...", config);
//             emms2_vendor_models_from_file(wlist->device_vendor_model,
//                 DEVICE_VERNDOR_MODEL_NUM, config);
//         } else if (strcmp(wlist->deviceType, DEV_PCS) == 0) {
//             const char * config = CONFIG_PATH "/" EMMS2_VENDOR_MODEL_PCS;
//             ems_syslog(LOG_ERR, "Loading emms2_vendor_model from %s...", config);
//             emms2_vendor_models_from_file(wlist->device_vendor_model,
//                 DEVICE_VERNDOR_MODEL_NUM, config);
//         }
//     }
//     return 0;
// }

//在数采设备加载完成后加载 不然可能无法匹配到点位信息

static void destory_mqtt_emms2_var_t(mqtt_emms2_var_t* var)
{
    if (var == NULL) return;

    if (var->offline_db != NULL)
    {
        sqlite3_close(var->offline_db->db);
        free(var->offline_db);
        var->offline_db = NULL;
    }
    if(var->alarm_report_list != NULL)
    {
        alarm_report_t* alarm_report = NULL;
        list_for_each_entry(alarm_report, &var->alarm_report_list->alarm_report_queue, list)
        {
            list_del(&alarm_report->list);
            free(alarm_report);
        }
    }

    free(var);
}

static inline int push_alarm_report_to_queue(mqtt_emms2_var_t* var, alarm_report_t* report)
{
    int ret = -1;
    pthread_mutex_lock(&var->alarm_report_list->alarm_report_queue_mutex);
    if (var->alarm_report_list->alarm_report_queue_size < var->alarm_report_list->alarm_report_queue_max_size)
    {
        list_add_tail(&report->list, &var->alarm_report_list->alarm_report_queue);
        var->alarm_report_list->alarm_report_queue_size++;
        ret = 0;
    }
    else
    {
        destory_alarm_report(report);
    }
    pthread_mutex_unlock(&var->alarm_report_list->alarm_report_queue_mutex);
    return ret;
}

static alarm_report_t* pop_alarm_from_report_queue(mqtt_emms2_var_t* var)
{
    alarm_report_t* report = NULL;
    pthread_mutex_lock(&var->alarm_report_list->alarm_report_queue_mutex);
    if (!list_empty(&var->alarm_report_list->alarm_report_queue))
    {
        report = list_first_entry(&var->alarm_report_list->alarm_report_queue, typeof(*report), list);
        list_del(&report->list);
        var->alarm_report_list->alarm_report_queue_size--;
    }
    pthread_mutex_unlock(&var->alarm_report_list->alarm_report_queue_mutex);
    return report;
}

static int get_alarm_report_queue_size(mqtt_emms2_var_t* var)
{
    int          num = 0;
    list_head_t* pos = &var->alarm_report_list->alarm_report_queue;
    pthread_mutex_lock(&var->alarm_report_list->alarm_report_queue_mutex);
    while (pos->next != &var->alarm_report_list->alarm_report_queue)
    {
        pos = pos->next;
        num++;
    }
    pthread_mutex_unlock(&var->alarm_report_list->alarm_report_queue_mutex);
    return num;
}

int clear_alarm_report_queue(mqtt_emms2_var_t* var)
{
    alarm_report_t* alarm = NULL;
    pthread_mutex_lock(&var->alarm_report_list->alarm_report_queue_mutex);
    list_for_each_entry(alarm, &var->alarm_report_list->alarm_report_queue, list)
    {
        list_del(&alarm->list);
        free(alarm);
    }
    var->alarm_report_list->alarm_report_queue_size = 0;
    return 0;
}

int mqtt_add_alarm_report_s(alarm_real_t* real)
{
    for(int i = 0; i < mqtt_emms2_var_list_num; i++)
    {
        mqtt_emms2_var_t* mqtt_emms2_var = mqtt_emms2_var_list[i];
        if (mqtt_emms2_var != NULL)
        {
            push_alarm_report_to_queue(mqtt_emms2_var, create_alarm_report_node(real));
        }
    }
    return 0;
}

static int mqtt_emms2_dp_init(mqtt_emms2_var_t* var)
{
    var->dp = malloc(sizeof(mqtt_dp_var_t));
    if (var->dp == NULL)
    {
        return -1;
    }
    memset(var->dp, 0, sizeof(mqtt_dp_var_t));
    mqtt_dp_var_t* dp = var->dp;
    
    dp->station_id = var->station_id;

    dp->cur_seq = 1;

    local_strlcpy(dp->base_param.host, var->base_param.addr, sizeof(dp->base_param.host));
    dp->base_param.port = var->base_param.port;
    local_strlcpy(dp->base_param.user, var->base_param.user, sizeof(dp->base_param.user));
    local_strlcpy(dp->base_param.pass, var->base_param.pass, sizeof(dp->base_param.pass));
    local_strlcpy(dp->base_param.up_url , var->base_param.config_url,sizeof(dp->base_param.up_url));
    dp->base_param.enable_tls    = var->base_param.tls;
    dp->base_param.disenable_lz4 = var->base_param.disable_lz4;
    dp->base_param.operation     = var->base_param.operation;
    char tmp[128] = {0};
    sprintf(tmp, "%s_emms2_dp", var->topic_sn);
    mqtt_session_t* session = mqtt_session_new(tmp, var->dp);
    if (session == NULL)
    {
        ems_syslog(LOG_NOTICE, "creat default!!");
        free(dp);
        var->dp = NULL;
        return -1;
    }
    dp->session = session;
    mqtt_dp_param_t* mqttparam = &var->dp->base_param;


    mqtt_session_set_address(session, mqttparam->host, mqttparam->port, mqttparam->user, mqttparam->pass);
    mqtt_session_set_opts(session, EMMS2_DEFAULT_IPC_QOS, EMMS2_KEEP_ALIVE_MAX);
    extern int mqtt_dp_handle_recv_msg(void* obj, mqtt_message_t* mqtt_msg);
    mqtt_session_set_callbacks(session, mqtt_dp_handle_recv_msg, NULL);

    if (mqttparam->enable_tls == 1)
        mqtt_session_set_tls(session, MQTT_DP_TLS_CAFILE, "", "");
    else if (mqttparam->enable_tls == 2)
        mqtt_session_set_tls(session, MQTT_DP_TLS_CAFILE, MQTT_DP_TLS_CERTFILE, MQTT_DP_TLS_KEYFILE);

    dp->log_fd              = log_module_client_init();
    dp->login.backoff_count = 0;
    dp->login.last_time     = 0;

    time_t now = time(NULL);

    dp->hardware.interval  = DEFAULT_MONITOR_INTERVAL;
    dp->hardware.last_time = now;

    dp->heartbeat.offline_count = 0;
    dp->heartbeat.last_time     = now;
    dp->heartbeat.interval = DEFAULT_HEARTBEAT_INTERVAL;

    dp->lc_info.interval  = DEFAULT_LC_INTERVAL;
    dp->lc_info.last_time = 0;

    dp->config.last_time = 0;
    dp->config.interval    = DEFAULT_CONFIG_INTERVAL;
    dp->config.operation = mqttparam->operation; // 投运状态
    
    char            sub_topic[EMMS2_MAX_BASE_LEN] = {0};
    sprintf(sub_topic, "emms2/" FUNCTION_TYPE_POSTRESP "/%s/+", var->topic_sn); //登陆应答
    mqtt_session_subscribe(session, sub_topic);
    ems_syslog(LOG_INFO, "sub topic:%s", sub_topic);

    sprintf(sub_topic, "emms2/" FUNCTION_TYPE_GET "/%s/+", var->topic_sn); //平台召测
    mqtt_session_subscribe(session, sub_topic);
    ems_syslog(LOG_INFO, "sub topic:%s", sub_topic);

    sprintf(sub_topic, "emms2/" FUNCTION_TYPE_POSTRESP "/%s/+/lz4", var->topic_sn); //登陆应答
    mqtt_session_subscribe(session, sub_topic);
    ems_syslog(LOG_INFO, "sub topic:%s", sub_topic);

    sprintf(sub_topic, "emms2/" FUNCTION_TYPE_GET "/%s/+/lz4", var->topic_sn); //平台召测
    mqtt_session_subscribe(session, sub_topic);
    ems_syslog(LOG_INFO, "sub topic:%s", sub_topic);

    sprintf(sub_topic, "emms2/" FUNCTION_TYPE_SET "/%s/+", var->topic_sn); //平台设置
    mqtt_session_subscribe(session, sub_topic);
    ems_syslog(LOG_INFO, "sub topic:%s", sub_topic);

    sprintf(sub_topic, "emms2/" FUNCTION_TYPE_SET "/%s/+/lz4", var->topic_sn); //平台设置
    mqtt_session_subscribe(session, sub_topic);
    ems_syslog(LOG_INFO, "sub topic:%s", sub_topic);

    get_board_sn(dp->topic_sn);
    init_mqtt_dp_lc_list(&dp->lc_info);

    return 0;
}

mqtt_emms2_var_t* mqtt_emms2_client_init(char* station_id)
{
    if (station_id == NULL || strlen(station_id) <= 0 || mqtt_emms2_var_list_num >= MQTT_EMMS2_MAX_THREAD_NUM) return NULL;

    mqtt_emms2_var_t* var = calloc(sizeof(mqtt_emms2_var_t), 1);
    if (var == NULL)
    {
        return NULL;
    }
    memset(var, 0, sizeof(mqtt_emms2_var_t));

    mqtt_emms2_var_list[mqtt_emms2_var_list_num++] = var;
    var->station_id = station_id;
    var->seq = 1; //默认从1开始
    if (load_mqtt_emms2_info(var, station_id) < 0)
    {
        ems_syslog(LOG_ERR,"load_mqtt_emms2_info error !!");
    }
    get_board_sn(var->topic_sn);
    ems_syslog(LOG_NOTICE,"client_id:%s",strlen(var->topic_sn)==0?"no_client_id":var->topic_sn);

    mqtt_emms2_param* mqttparam = &var->base_param;
    char tmp[128] = {0};
    sprintf(tmp, "%s_%s", strlen(var->topic_sn)==0?"no_client_id":var->topic_sn, var->station_id);
    mqtt_session_t* session = mqtt_session_new(tmp, var);
    if (session == NULL)
    {
        ems_syslog(LOG_NOTICE,"creat default!!");
        if(var!=NULL)
        {
            destory_mqtt_emms2_var_t(var);
            var = NULL;
        }
        session = NULL;
        return NULL;
    }
    var->session = session;

    mqtt_session_set_address(session, (mqttparam->addr==NULL)?EMMS2_BROKER_ADDR:mqttparam->addr,(mqttparam->port == 0)?1883:mqttparam->port,(mqttparam->user==NULL)?EMMS2_BROKER_USER:mqttparam->user , (mqttparam->pass==NULL)?EMMS2_BROKER_PASS:mqttparam->pass);
    mqtt_session_set_opts(session, EMMS2_DEFAULT_IPC_QOS, EMMS2_KEEP_ALIVE_MAX);
    mqtt_session_set_callbacks(session, mqtt_handle_recv_msg, mqtt_state_change);

    if (var->base_param.tls == 1)
        mqtt_session_set_tls(session, EMMS2_TLS_CAFILE, "", "");
    else if (var->base_param.tls == 2)
        mqtt_session_set_tls(session, EMMS2_TLS_CAFILE, EMMS2_TLS_CERTFILE, EMMS2_TLS_KEYFILE);

    var->last_check_time = 0;          // 初始化为0
    var->device_status_list = NULL;    // 设备状态列表初始化为空
    var->device_status_count = 0;      // 设备状态数量初始化为0
    var->tag_config_list = NULL;       // 点位配置列表初始化为空
    var->tag_config_count = 0;         // 点位配置数量初始化为0
    
    if (load_changetag_config(var) != 0)
    {
        ems_syslog(LOG_WARNING, "变化上报点位配置加载失败,将使用默认配置");
    }
    init_change_report_list(var);

    ems_syslog(LOG_NOTICE, "mqtt client creat success!!");
    ems_syslog(LOG_NOTICE, "client use :[%s],[%d],[%s],[%s]!!", (strlen(mqttparam->addr) == 0) ? EMMS2_BROKER_ADDR : mqttparam->addr, (mqttparam->port == 0) ? 1883 : mqttparam->port, mqttparam->user, mqttparam->pass);

    
    var->log_fd = log_module_client_init();

    init_offline_db_table(var, station_id);
    init_mqtt_alarm_report_list(var);
    init_pub_dev_list(var);

    pthread_mutex_init(&var->mutex, NULL);
    pthread_cond_init(&var->cond, NULL);

    if (strcmp(var->station_id, MQTT_STACTRL_STATION_ID) == 0)
    {
        var->cond_wait_sec = 1;
        var->backoff_count = sizeof(backoff_intervals) / sizeof(backoff_intervals[0]) - 4;
        // 计算用标准点位
    }
    else
    {
        var->cond_wait_sec = 10;
        var->backoff_count = sizeof(backoff_intervals) / sizeof(backoff_intervals[0]) - 1;
    }

    var->emms2_status = EMMS2_LOGINING;
    mqtt_emms2_subscribe_all(var);

    mqtt_emms2_dp_init(var);
    init_sync_cabinet_info_record();

    ems_syslog(LOG_NOTICE, "take new mqtt var!!");
    return var;
}


void mqtt_emms2_start(mqtt_emms2_var_t *var)
{
    if((var!=NULL)&&(var->session!=NULL))
    {
        ems_syslog(LOG_NOTICE,"MQTT cloud CLIENT START!!");
        mqtt_session_start(var->session);
        mqtt_session_start(var->dp->session);

        pthread_create(&var->mqtt_cycle_tid, NULL, mqtt_emms2_cycle_task, var);

        char task_name[128] = {0};
        snprintf(task_name, sizeof(task_name), "mqtt_task_%s", var->station_id);
        pthread_setname_np(var->mqtt_cycle_tid, task_name);
    }
}


/*************************************/
