#include "mqtt_dp.h"
#include "../mqtt_emms2/mqtt_emms2.h"
#include "../ems/jsonct.h"
#include "../common.h"
#include "../proto_forward.h"
#include "collector-api.h"
#include "lz4.h"
#include "../log_module/log_module.h"
#include "sqlite3.h"
#include <stdbool.h>
#include <stdint.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <sys/syslog.h>
#include <sys/types.h>
#include <sys/ipc.h>
#include <sys/msg.h>
#include "../opcua/ua_client.h"
#include "../cc_lc_com/pub_fun.h"
#include "lnxall_buffer.h"
#include <dirent.h>
#include <fnmatch.h>
#include <time.h>
#include <unistd.h>
#include "../ems/ems_activition.h"

extern int                cabinet_info_tab_count_g;
extern cabinet_info_tab_t cabinet_info_tab[];
extern cabinet_info_t     cabinet_info;
extern SYSTEM_INFO        system_info;
extern pthread_mutex_t    cab_info_mutex;

static int load_mqtt_dp_config(mqtt_dp_var_t* var, char* file_path)
{
    int    ret  = -1;
    char*  cfg  = NULL;
    cJSON* root = NULL;
    cJSON* item = NULL;
    cfg         = read_file_data(file_path);
    if (cfg == NULL)
    {
        ems_syslog(LOG_ERR, "read file %s failed", file_path);
        goto ERROR;
    }

    root = cJSON_Parse(cfg);
    if (root == NULL)
    {
        ems_syslog(LOG_ERR, "parse file %s failed", file_path);
        goto ERROR;
    }

    item = cJSON_GetObjectItem(root, "host");
    if (item == NULL)
    {
        ems_syslog(LOG_ERR, "parse file %s failed, no host", file_path);
        goto ERROR;
    }
    local_strlcpy(var->base_param.host, item->valuestring, sizeof(var->base_param.host));

    item = cJSON_GetObjectItem(root, "port");
    if (item == NULL)
    {
        ems_syslog(LOG_ERR, "parse file %s failed no port", file_path);
        goto ERROR;
    }
    var->base_param.port = item->valueint;

    item = cJSON_GetObjectItem(root, "user");
    if (item == NULL)
    {
        ems_syslog(LOG_ERR, "parse file %s failed, no user", file_path);
        goto ERROR;
    }
    local_strlcpy(var->base_param.user, item->valuestring, sizeof(var->base_param.user));
    item = cJSON_GetObjectItem(root, "pass");
    if (item == NULL)
    {
        ems_syslog(LOG_ERR, "parse file %s failed, no pass", file_path);
        goto ERROR;
    }
    local_strlcpy(var->base_param.pass, item->valuestring, sizeof(var->base_param.pass));

    item = cJSON_GetObjectItem(root, "uploadConfigUrl");
    if (item == NULL)
    {
        char tmp[256] = {0};
        snprintf(tmp, sizeof(tmp), "http://%s/%s", var->base_param.host, MQTT_DP_UPLOAD_CONFIG_API);
        local_strlcpy(var->base_param.up_url, tmp, sizeof(var->base_param.up_url));
    }
    else{
        local_strlcpy(var->base_param.up_url, item->valuestring, sizeof(var->base_param.up_url));
    }

    item = cJSON_GetObjectItem(root, "enable_tls");
    if (item == NULL || item->valueint > 2)
    {
        ems_syslog(LOG_NOTICE, "parse file %s failed no enable_tls", file_path);
        var->base_param.enable_tls = 0;
    }
    else
    {
        var->base_param.enable_tls = item->valueint;
    }

    item = cJSON_GetObjectItem(root, "disable_lz4");
    if (item == NULL) 
    {
        ems_syslog(LOG_ERR, "parse file %s failed no disable_lz4", file_path);
        var->base_param.disenable_lz4 = 0;
    }
    else
    {
        var->base_param.disenable_lz4 = item->valueint > 0 ? 1 : 0;
    }

    // var->base_param.operation = get_opeartion_status(var->station_id);
    var->base_param.operation = 1 ; 
    ret =  0;
ERROR:
    if (cfg != NULL) free(cfg);
    if (root != NULL) cJSON_Delete(root);
    return ret;
}

int get_opeartion_status(char* station_id)
{
    int ret = 0;
    cJSON* root = NULL;
    cJSON* obj  = NULL;
    char* file = read_file_data(MQTT_OPERATION_FILE);
    if(file == NULL)
    {
        ems_syslog(LOG_NOTICE, "read file %s failed", MQTT_OPERATION_FILE);
        goto END;
    }
    root = cJSON_Parse(file);
    if(root == NULL)
    {
        ems_syslog(LOG_NOTICE, "parse file %s failed", MQTT_OPERATION_FILE);
        goto END;
    }

    obj = cJSON_GetObjectItem(root, station_id);
    if (obj != NULL)
    {
        ret = obj->valueint;
    }

END:
    if (file != NULL) free(file);
    if (root != NULL) cJSON_Delete(root);
    return ret;
}

static void destory_mqtt_dp_var_t(mqtt_dp_var_t* var)
{
    if (var->session != NULL)
    {
        mqtt_session_disconn(var->session);
    }
}
static int mqtt_process_recv_msg(mqtt_dp_var_t* var, char* type, char* function_id, char* data)
{
    if (type == NULL || function_id == NULL || data == NULL || type[0] == '\0' || function_id[0] == '\0' || data[0] == '\0')
    {
        ems_syslog(LOG_ERR, "mqtt_process_recv_msg param error");
        return -1;
    }

    int    ret  = -1;
    cJSON* root = cJSON_Parse(data);
    if (root == NULL)
    {
        ems_syslog(LOG_ERR, "mqtt_process_recv_msg parse data failed");
        goto END;
    }

    if (strcmp(function_id, FUNCTION_ID_LOGIN) == 0)
    {
        extern int mqtt_process_login_message(mqtt_dp_var_t * var, char* type, cJSON* root);
        ret = mqtt_process_login_message(var, type, root);
    }
    else if (strcmp(function_id, FUNCTION_ID_HEARTBEAT) == 0)
    {
        extern int mqtt_process_heartbeat_message(mqtt_dp_var_t * var, char* type, cJSON* root);
        ret = mqtt_process_heartbeat_message(var, type, root);
    }
    else if (strcmp(function_id, FUNCTION_ID_LC_INFO) == 0)
    {
        extern int mqtt_process_lcinfo_message(mqtt_dp_var_t * var, char* type, cJSON* root);
        ret = mqtt_process_lcinfo_message(var, type, root);
    }
    else if (strcmp(function_id, FUNCTION_ID_CONFIG) == 0)
    {
        extern int mqtt_process_config_message(mqtt_dp_var_t * var, char* type, cJSON* root);
        ret = mqtt_process_config_message(var, type, root);
    }
    else if (strcmp(function_id, FUNCTION_ID_MONITOR) == 0)
    {
        extern int mqtt_process_monitor_message(mqtt_dp_var_t * var, char* type, cJSON* root);
        ret = mqtt_process_monitor_message(var, type, root);
    }
    else if (strcmp(function_id, FUNCTION_ID_UPGRADE) == 0)
    {
        extern int mqtt_process_upgrade_message(mqtt_dp_var_t * var, char* type, cJSON* root);
        ret = mqtt_process_upgrade_message(var, type, root);
    }
    else if (strcmp(function_id, FUNCTION_ID_PROGRESS) == 0)
    {
        extern int mqtt_process_progress_message(mqtt_dp_var_t * var, char* type, cJSON* root);
        ret = mqtt_process_progress_message(var, type, root);
    }
    else if (strcmp(function_id, FUNCTION_ID_OPERATION) == 0)
    {
        void mqtt_dp_report_activition_result(mqtt_dp_var_t * var, char* type, cJSON* root);
        mqtt_dp_report_activition_result(var, type, root);
    }
    // else
    // {
    //     ems_syslog(LOG_ERR, "Unkown function_id:%s", function_id);
    // }

END:
    if (root != NULL) cJSON_Delete(root);
    return ret;
}

//MQTT消息处理函数，负责接收并解析MQTT消息，根据消息内容执行相应操作
int mqtt_dp_handle_recv_msg(void* obj, mqtt_message_t* mqtt_msg)
{
    mqtt_dp_var_t* var = (mqtt_dp_var_t*)obj;

    size_t data_len = mqtt_msg->payloadLen;
    char*  get_data = NULL;

    int  ret_dec          = 0;
    char message_type[64] = {0}, tmp_c[64] = {0}, function_id[64] = {0};

    ems_syslog(LOG_INFO, "%s, received MQTT topic:%s payload length:%d", var->station_id, mqtt_msg->topic, mqtt_msg->payloadLen);
    if (strstr(mqtt_msg->topic, "lz4"))
    {
        data_len *= 15;
        get_data = (char*)malloc(data_len);
        if (!get_data) goto END;
        ret_dec = LZ4_decompress_safe(mqtt_msg->payload, get_data, mqtt_msg->payloadLen, data_len - 1);
        if (ret_dec < 0)
        {
            ems_syslog(LOG_ERR, " LZ4_decompress not enough space,ret:%d, data_len:%zu, payloadLen:%d", ret_dec, data_len, mqtt_msg->payloadLen);
            goto END;
        }
        get_data[ret_dec] = '\0';
        ems_syslog(LOG_DEBUG, "decompress  payload %s", get_data);
        sscanf(mqtt_msg->topic, "emms2/%[^/]/%[^/]/%s/lz4", message_type, tmp_c, function_id);
    }
    else
    {
        get_data = (char*)malloc(data_len + 1);
        if (!get_data) goto END;
        memcpy(get_data, mqtt_msg->payload, data_len);
        get_data[data_len] = '\0';
        sscanf(mqtt_msg->topic, "emms2/%[^/]/%[^/]/%s", message_type, tmp_c, function_id);
    }
    if (message_type[0] == '\0' || function_id[0] == '\0')
    {
        ems_syslog(LOG_ERR, "mqtt topic format error:%s", mqtt_msg->topic);
        goto END;
    }
    //var设备状态，message_type 消息类型，function_id 消息ID，get_data 消息内容
    ret_dec = mqtt_process_recv_msg(var, message_type, function_id, get_data);

END:
    if (get_data != NULL) free(get_data);
    return ret_dec;
}

static int mqtt_state_change(void* obj, int state)
{
    if (state != MQTT_CONNECTED)
    {
        server_add_new_log(LOG_EMMS2_NET, "EMS", LOG_MODE_PERIOD, 1, NULL);
        ((mqtt_dp_var_t*)obj)->status = MQTT_DP_CONNECTING;
    }
    else
    {
        server_add_new_log(LOG_EMMS2_NET, "EMS", LOG_MODE_PERIOD, 2, NULL);
        ((mqtt_dp_var_t*)obj)->status = MQTT_DP_LOGINING;
    }
    return 0;
}

static int mqtt_dp_subscribe_all(mqtt_dp_var_t* var)
{
    mqtt_session_t* session                 = var->session;
    char            sub_topic[MQTT_DP_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);
    return 0;
}

int mqtt_dp_publish_message(mqtt_dp_var_t* var, char* pub_topic, char* data, int data_len)
{
    if (var->base_param.disenable_lz4)
    {
        return mqtt_session_publish(var->session, pub_topic, data, data_len);
    }
    else
    {
        int   ret             = 0;
        int   src_size        = data_len + 1;
        int   max_dst_size    = LZ4_compressBound(src_size);
        char* compressed_data = malloc(max_dst_size);

        const int compressed_data_size = LZ4_compress_default(data, compressed_data, src_size, max_dst_size);
        if (compressed_data_size <= 0)
        {
            ems_syslog(LOG_ERR, "LZ4_compress err");
            free(compressed_data);
            return -1;
        }
        // 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);
        ems_syslog(LOG_INFO, "lz4 change topic [%s]", topic_lz4);
        ret = mqtt_session_publish(var->session, topic_lz4, compressed_data, compressed_data_size);
        free(compressed_data);
        return ret;
    }
}

int init_mqtt_dp_lc_list(mqtt_dp_lc_list_t* lc_list)
{
    // 本机信息
    struct lc_info_t* lc = &lc_list->ems_info;
    get_board_sn(lc->sn);
    sprintf(lc->name, "EMS");
    sprintf(lc->no, "EMS");
    lc->is_lc = 0;
    sprintf(lc->ua_url, "opc.tcp://127.0.0.1:4840");

    if (cabinet_info.en_lc_ctrl)
    {
        pcs_ctrl_var_t* var    = get_pcs_ctrl_var();
        lc_it*          list   = var->lc_list;
        int             lc_cnt = var->lc_list_len;
        if (lc_cnt > 0)
        {
            lc_list->lc_list = calloc(sizeof(struct lc_info_t), lc_cnt);
            lc_list->lc_num  = 0;
            for (int j = 0; j < lc_cnt; j++)
            {
                lc_it*            lc      = &list[j];
                if (lc->name[0] == '\0' || lc->no[0] == '\0' || lc->sn[0] == '\0' || lc->ipaddr[0] == '\0' || lc->type == eLC_TYPE_VLC)
                    continue;

                struct lc_info_t* lc_info = &lc_list->lc_list[lc_list->lc_num++];
                sprintf(lc_info->sn, "%s", lc->sn);
                sprintf(lc_info->name, "%s", lc->name);
                sprintf(lc_info->no, "%s", lc->no);
                lc_info->is_lc = 1;
                if (lc->type == eLC_TYPE_LC)
                    sprintf(lc_info->ua_url, "opc.tcp://%s:4841", lc->ipaddr);
                else if (lc->type == eLC_TYPE_EMS)
                    sprintf(lc_info->ua_url, "opc.tcp://%s:4840", lc->ipaddr);
            }
        }
    }
    else
    {
        lc_list->lc_list = NULL;
        lc_list->lc_num  = 0;
    }
    return 0;
}

int mqtt_dp_client_init(mqtt_dp_var_t* var)
{
    if (var == NULL) { return -1; }

    memset(var, 0, sizeof(mqtt_dp_var_t));
    var->station_id = MQTT_DP_STATION_ID;

    var->cur_seq = 1; //默认从1开始
    if (load_mqtt_dp_config(var, MQTT_DB_CONFIG_PATH) < 0)
    {
        ems_syslog(LOG_ERR, "load_mqtt_emms2_info error !!");
        goto ERROR;
    }

    get_board_sn(var->topic_sn);

    mqtt_dp_param_t* mqttparam = &var->base_param;
    char tmp[129] = {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!!");
        goto ERROR;
    }
    var->session = session;

    mqtt_session_set_address(session, mqttparam->host, mqttparam->port, mqttparam->user, mqttparam->pass);
    mqtt_session_set_opts(session, MQTT_DP_DEFAULT_IPC_QOS, MQTT_DP_KEEP_ALIVE_MAX);
    mqtt_session_set_callbacks(session, mqtt_dp_handle_recv_msg, mqtt_state_change);

    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);

    ems_syslog(LOG_NOTICE, "mqtt client creat success, client use :[%s],[%d],[%s],[%s]!!", mqttparam->host, mqttparam->port, mqttparam->user, mqttparam->pass);

    time_t now               = time(NULL);
    var->log_fd              = log_module_client_init();
    var->login.backoff_count = 0;
    var->login.last_time     = now;

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

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

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

    var->config.last_time = 0;
    var->config.interval    = DEFAULT_CONFIG_INTERVAL;
    var->config.operation = 1; // 投运状态(已屏蔽)

    mqtt_dp_subscribe_all(var);
    init_mqtt_dp_lc_list(&var->lc_info);

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

    return 0;
ERROR:
    if (var != NULL) destory_mqtt_dp_var_t(var);
    return -1;
}

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

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

static const int backoff[]   = {1, 10, 30, 60, 600, 1800, 3600};
static const int backoff_len = sizeof(backoff) / sizeof(backoff[0]);

static int mqtt_dp_publish_login_message(mqtt_dp_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);
    }
    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}",
                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);
    pthread_mutex_unlock(&cab_info_mutex);

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

    char pub_topic[512] = {0};
    sprintf(pub_topic, "emms2/" FUNCTION_TYPE_POST "/%s/" FUNCTION_ID_LOGIN, var->topic_sn);
    ems_syslog(LOG_INFO, "pub topic [%s] data:[%s]", pub_topic, data_all);
    int ret = mqtt_dp_publish_message(var, pub_topic, data_all, strlen(data_all));
    free(data_info);
    free(data_all);
    return ret;
}

static void mqtt_dp_login_task(mqtt_dp_var_t* var, time_t now) //登录
{
    if (now - var->login.last_time > backoff[var->login.backoff_count])
    {
        if (mqtt_dp_publish_login_message(var) < 0)
        {
            ems_syslog(LOG_WARNING, "mqtt_dp_publish_login_message failed!!");
        }
        var->login.backoff_count++;
        if (var->login.backoff_count >= backoff_len)
        {
            var->login.backoff_count = backoff_len - 1;
        }
        var->login.last_time = now;
    }
}

static int mqtt_dp_publish_heartBeat_message(mqtt_dp_var_t* var, time_t now) //心跳
{
    char  pub_topic[256] = {0};
    char* data_all       = malloc(1024);

    package_based_field(var, data_all, FUNCTION_HEARTBEAT, 0, NULL);
    snprintf(pub_topic, sizeof(pub_topic), "emms2/" FUNCTION_TYPE_POST "/%s/" FUNCTION_HEARTBEAT, var->topic_sn);

    ems_syslog(LOG_INFO, "pub topic [%s] data:[%s]", pub_topic, data_all);

    int ret = mqtt_dp_publish_message(var, pub_topic, data_all, strlen(data_all));
    free(data_all);
    return ret;
}

void mqtt_dp_heartbeat_task(mqtt_dp_var_t* var, time_t now)
{
    if (now - var->heartbeat.last_time > var->heartbeat.interval)
    {
        if (mqtt_dp_publish_heartBeat_message(var, now) < 0)
        {
            ems_syslog(LOG_WARNING, "mqtt_dp_publish_heartBeat_message failed!!");
        }
        else
        {
            var->heartbeat.last_time = now;
            var->heartbeat.offline_count++;
        }
    }
}


void mqtt_dp_check_ems_config(mqtt_dp_var_t* var)
{
    char md5path[256]    = {0};
    char md5[33]          = {0};
    char filepath[280] = {0};

    snprintf(md5path, sizeof(md5path), "%s%s", MQTT_DP_CONFIG_MD5_FILE, var->station_id);

    char* his = read_file_data(md5path);
    cJSON* md5_root = NULL;
    char* md5str = NULL;
    DIR* dir = NULL;
    struct dirent* entry = NULL;

    if(his != NULL)
    {
        md5_root = cJSON_Parse(his);
        if (md5_root == NULL)
        {
            md5_root = cJSON_CreateObject();
        }
    }
    else{
        md5_root = cJSON_CreateObject();
    }

    // 打开目录
    dir = opendir(MQTT_DP_CONFIG_TAR_DIR);
    if (dir == NULL)
    {
        goto END;
    }

    // 遍历目录中的文件
    while ((entry = readdir(dir)) != NULL)
    {
        if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0)
        {
            continue;
        }
        if (fnmatch("*.tar.gz", entry->d_name, 0) != 0) continue;
        char* filename = entry->d_name;
        snprintf(filepath, sizeof(filepath), "%s/%s", MQTT_DP_CONFIG_TAR_DIR, filename);
        if (access(filepath, F_OK) != 0)
        {
            ems_syslog(LOG_ERR, "file %s not exist", filepath);
            continue;
        }
        get_file_md5(filepath, md5);
        cJSON* item = cJSON_GetObjectItem(md5_root, filename);
        if (item == NULL || strcmp(item->valuestring, md5) != 0)
        {
            ems_syslog(LOG_INFO, "file %s md5 changed, old: %s, new: %s, start upload", filepath, item ? item->valuestring : "NULL", md5);
            if (http_post_upload_file(var->base_param.up_url, filename, filepath, NULL, NULL) == 0)
            {
                if (item == NULL)
                    cJSON_AddItemToObject(md5_root, filename, cJSON_CreateString(md5));
                else
                    cJSON_ReplaceItemInObject(md5_root, filename, cJSON_CreateString(md5));
            }
            else
            {
                ems_syslog(LOG_ERR, "upload config file %s failed, md5: %s", filepath, md5);
            }
        }
    }

    md5str = cJSON_Print(md5_root);
    write_file_data_safe(md5path, md5str, strlen(md5str));

END:
    if (md5str) free(md5str);
    if (his) free(his);
    if (md5_root) cJSON_Delete(md5_root);
    if (dir) closedir(dir);
}

void mqtt_dp_config_task(mqtt_dp_var_t* var, time_t now) //配置
{
    if (var->config.operation == false) return; // 未投运直接退出
    if (now - var->config.last_time > var->config.interval)
    {
        mqtt_dp_check_ems_config(var);
        var->config.last_time = now;
    }
}

void get_mem_info(uint32_t* mem_used, uint32_t* mem_total)
{
    FILE* fd;
    fd = fopen("/proc/meminfo", "r");
    if (fd == NULL)
    {
        perror("open /proc/meminfo failed\n");
        exit(0);
    }
    size_t   bytes_read;
    size_t   read;
    char*    line   = NULL;
    uint32_t avimem = 0;
    uint32_t total  = 0;

    while ((read = getline(&line, &bytes_read, fd)) != -1)
    {
        if (strstr(line, "MemTotal") != NULL)
        {
            sscanf(line, "%*s%d%*s", &total);
        }
        if (strstr(line, "MemAvailable") != NULL)
        {
            sscanf(line, "%*s%d%*s", &avimem);
        }
    }
    if (line != NULL)
    {
        free(line);
    }
    *mem_used  = (total - avimem) / 1024;
    *mem_total = total / 1024;

    fclose(fd);
}

extern void get_cpu_usage(double* cpu_usage);
extern int  cmd_call(char* arg, int timeout, char** dist_arg, int getreturn);
static int  mqtt_db_get_hardware_info(hardware_info_t* var)
{
    unsigned long long disk_used, disk_total;
    uint32_t           mem_used, mem_total;
    double             cpu_usage;
    char               tmp_c[64] = {0};
    double             tmp       = 0;
    get_cpu_usage(&cpu_usage);
    var->cpu_usage = (int)(cpu_usage);
    get_mem_info(&mem_used, &mem_total);
    var->mem_used  = mem_used;
    var->mem_total = mem_total;

    get_disk_usage("/", NULL, &disk_used, &disk_total);
    var->disk_info_num = 1;
    snprintf(var->disk_info[0].name, sizeof(var->disk_info[0].name), "硬盘");
    var->disk_info[0].used  = disk_used / 1024/1024;
    var->disk_info[0].total = disk_total/1024/1024;

    const char* disk_path = get_disk_path();
    snprintf(tmp_c, sizeof(tmp_c), "mountpoint -q %s", disk_path);
    if (cmd_call(tmp_c, 10, NULL, 1) == 0)
    {
        get_disk_usage(disk_path, &tmp, &disk_used, &disk_total);
        var->disk_info_num = 2;
        snprintf(var->disk_info[1].name, sizeof(var->disk_info[1].name), "扩展磁盘");
        var->disk_info[1].used  = disk_used / 1024/1024;
        var->disk_info[1].total = disk_total / 1024/1024;
    }

    char* sf = read_file_data("/etc/signal_4g.txt");
    if (sf != NULL)
    {
        int signal  = atoi(sf);
        var->signal = signal;
        free(sf);
    }
    if (var->signal <= 0)
        var->signal = 99;
    return 0;
}

static void mqtt_dp_publish_hardware_monitor_message(mqtt_dp_var_t* var) //硬件监控
{
    char  pub_topic[512] = {0};
	char  data_all[2048];

    hardware_info_t info;

    if (mqtt_db_get_hardware_info(&info) < 0)
    {
        ems_syslog(LOG_ERR, "mqtt_db_get_hardware_info failed!!");
		return;
    }

    char data_info[1024] = {0};
    int  size            = sizeof(data_info);

    int index = snprintf(data_info,
                         size,
                         "\"info\":{\"cpuUsed\":%d,\"memUsed\":%d,\"memTotal\":%d,",
                         info.cpu_usage,
                         info.mem_used,
                         info.mem_total);

    index += snprintf(data_info + index, size - index, "\"diskInfo\":[");
    if (info.disk_info_num > 0)
    {
        for (int i = 0; i < info.disk_info_num; i++)
        {
            index += snprintf(data_info + index,
                              size - index,
                              "{\"diskName\":\"%s\",\"diskUsed\":%d,\"diskTotal\":%d},",
                              info.disk_info[i].name,
                              info.disk_info[i].used,
                              info.disk_info[i].total);
        }
        index -= 1;
    }
    index += snprintf(data_info + index, size - index, "],");
    index += snprintf(data_info + index, size - index, "\"signal\":%d}", info.signal);

    package_based_field(var, data_all, FUNCTION_ID_MONITOR, 0, data_info);
    snprintf(pub_topic, sizeof(pub_topic), "emms2/" FUNCTION_TYPE_POST "/%s/" FUNCTION_ID_MONITOR, var->topic_sn);

    mqtt_dp_publish_message(var, pub_topic, data_all, strlen(data_all));
    ems_syslog(LOG_INFO, "pub topic [%s] data:[%s]", pub_topic, data_all);
}

void mqtt_hardware_monitor_task(mqtt_dp_var_t* var, time_t now) //硬件监控
{
    if (now - var->hardware.last_time > var->hardware.interval)
    {
        mqtt_dp_publish_hardware_monitor_message(var);
        var->hardware.last_time = now;
    }
}

void mqtt_report_activition_task(mqtt_dp_var_t* var)
{
    ems_activition_record_t* recort = ems_activition_get_record();
    if (atomic_load(&recort->is_activition) == true && recort->is_report == false)
    {
        char* data_all       = malloc(1024);
        char  pub_topic[512] = {0};
        char  data_info[256] = {0};
        snprintf(data_info, sizeof(data_info), "\"activationState\":1");

        package_based_field(var, data_all, FUNCTION_ID_OPERATION, 0, data_info);
        snprintf(pub_topic, sizeof(pub_topic), "emms2/" FUNCTION_TYPE_POST "/%s/" FUNCTION_ID_OPERATION, var->topic_sn);
        mqtt_dp_publish_message(var, pub_topic, data_all, strlen(data_all));
        ems_syslog(LOG_INFO, "pub topic [%s] data:[%s]", pub_topic, data_all);
        recort->is_report = true;
        free(data_all);
    }
}

void mqtt_dp_report_activition_result(mqtt_dp_var_t* var, char* type, cJSON* root)
{
    cJSON* result = cJSON_GetObjectItem(root, "result");
    cJSON* time   = cJSON_GetObjectItem(root, "activationTime");
    if (result == NULL)
    {
        return;
    }

    if (result && result->valueint == 0)
    {
        ems_activition_set_time(0, 0);
    }
    else if (time != NULL && time->valuestring != NULL && time->valuestring[0] != '\0')
    {
        char*     time_str = time->valuestring;
        struct tm tm;
        memset(&tm, 0, sizeof(struct tm));
        if (strptime(time_str, "%Y-%m-%d %H:%M:%S", &tm) == NULL)
        {
            ems_syslog(LOG_ERR, "时间字符串格式不正确");
            return;
        }
        time_t t = mktime(&tm);
        ems_activition_set_time(t, 1);
    }
    else
    {
        // 异常处理
        ems_syslog(LOG_ERR, "激活结果获取异常");
        ems_activition_record_t* recort = ems_activition_get_record();
        recort->is_report = false;
    }
}

static int lc_get_product(UA_Client* client, char* product)
{
    if (client == NULL || product == NULL) return -1;
    int        ret = -1;
    UA_Variant val;

    UA_Variant_init(&val);
    UA_StatusCode retval = UA_Client_readValueAttribute(client, system_info.aboutInfo.productNodeId, &val);
    if (retval == UA_STATUSCODE_GOOD)
    {
        if (val.type == &UA_TYPES[UA_TYPES_STRING])
        {
            UA_String* str = (UA_String*)val.data;
            if(str->length > 0)
            {
                memcpy(product, (char*)str->data, str->length);
                product[str->length] = '\0';
                ret                  = 0;
            }
        }
    }
    UA_Variant_clear(&val);
    return ret;
}

static int lc_get_model(UA_Client* client, char* model)
{
    if (client == NULL || model == NULL) return -1;
    int        ret = -1;
    UA_Variant val;

    UA_Variant_init(&val);
    UA_StatusCode retval = UA_Client_readValueAttribute(client, system_info.aboutInfo.modelNodeId, &val);
    if (retval == UA_STATUSCODE_GOOD)
    {
        if (val.type == &UA_TYPES[UA_TYPES_STRING])
        {
            UA_String* str = (UA_String*)val.data;
            if(str->length > 0)
            {
                memcpy(model, (char*)str->data, str->length);
                model[str->length] = '\0';
                ret                = 0;
            }
        }
    }
    UA_Variant_clear(&val);
    return ret;
}

static int lc_get_version(UA_Client* client, char* version)
{
    if (client == NULL || version == NULL) return -1;
    int        ret = -1;
    UA_Variant val;

    UA_Variant_init(&val);
    UA_StatusCode retval = UA_Client_readValueAttribute(client, system_info.aboutInfo.versionNodeId, &val);
    if (retval == UA_STATUSCODE_GOOD)
    {
        if (val.type == &UA_TYPES[UA_TYPES_STRING])
        {
            UA_String* str = (UA_String*)val.data;
            if (str->length > 0)
            {
                memcpy(version, (char*)str->data, str->length);
                version[str->length] = '\0';
                ret                  = 0;
            }
        }
    }
    UA_Variant_clear(&val);
    return ret;
}

static int get_lc_info(struct lc_info_t* lc_info)
{
    char* url      = NULL;
    int   update   = 0;
    char  tmp[256] = {0};
    url = lc_info->ua_url;
    if (url[0] == '\0') return 0;

    UA_Client* client = ua_connect_opcua(url);
    if (client == NULL) return 0;

    if (lc_get_product(client, tmp) < 0)
    {
        ems_syslog(LOG_ERR, "lc_get_model failed");
    }
    else
    {
        if (strcmp(tmp, lc_info->product) != 0)
        {
            update = 1;
            local_strlcpy(lc_info->product, tmp, sizeof(lc_info->product));
        }
    }
    if (lc_get_version(client, tmp) < 0)
    {
        ems_syslog(LOG_ERR, "lc_get_version failed");
    }
    else
    {
        if (strcmp(tmp, lc_info->version) != 0)
        {
            update = 1;
            local_strlcpy(lc_info->version, tmp, sizeof(lc_info->version));
        }
    }
    if (lc_get_model(client, tmp) < 0)
    {
        ems_syslog(LOG_ERR, "lc_get_model failed");
    }
    else
    {
        if (strcmp(tmp, lc_info->model) != 0)
        {
            update = 1;
            local_strlcpy(lc_info->model, tmp, sizeof(lc_info->model));
            ems_syslog(LOG_ERR, "lc_get_model %s, product:%s",lc_info->model, tmp);
        }
    }
    ua_client_delete(client);
    return update;
}

static int mqtt_dp_publish_lc_info_message(mqtt_dp_var_t* var)
{
    int                ret = -1;
    struct lnxall_buff buf;
    lbuff_init(&buf, 1024);
    const int size     = 1024;
    char*     data_all = NULL;
    char*     topic    = NULL;
    char*     info     = malloc(size);

    lbuff_sprintf(&buf, "\"info\":[");
    snprintf(info,
             size,
             "{\"lc\":%d,\"name\":\"%s\",\"sn\":\"%s\",\"model\":\"%s\", \"productId\":\"%s\",\"firmware\":\"%s\"},",
             var->lc_info.ems_info.is_lc ? 2 : 1,
             var->lc_info.ems_info.name,
             var->lc_info.ems_info.sn,
             var->lc_info.ems_info.model,
             var->lc_info.ems_info.product,
             var->lc_info.ems_info.version);
    lbuff_append(&buf, info, strlen(info));

    for (int i = 0; i < var->lc_info.lc_num; i++)
    {
        snprintf(info,
                 size,
                 "{\"lc\":%d,\"name\":\"%s\",\"sn\":\"%s\",\"model\":\"%s\", \"productId\":\"%s\",\"firmware\":\"%s\"},",
                 var->lc_info.lc_list[i].is_lc ? 2 : 1,
                 var->lc_info.lc_list[i].name,
                 var->lc_info.lc_list[i].sn,
                 var->lc_info.lc_list[i].model,
                 var->lc_info.lc_list[i].product,
                 var->lc_info.lc_list[i].version);

        lbuff_append(&buf, info, strlen(info));
    }
    buf.bufptr[buf.curlen - 1] = ']';

    topic    = malloc(256);
    data_all = malloc(buf.curlen + 512);
    package_based_field(var, data_all, FUNCTION_ID_LC_INFO, 0, buf.bufptr);
    snprintf(topic, 256, "emms2/" FUNCTION_TYPE_POST "/%s/" FUNCTION_ID_LC_INFO, var->topic_sn);

    ret = mqtt_dp_publish_message(var, topic, data_all, strlen(data_all));

    if (data_all) free(data_all);
    if (info) free(info);
    if (topic) free(topic);
    lbuff_free(&buf);
    
    return ret;
}

void mqtt_dp_lc_info_task(mqtt_dp_var_t* var, time_t now) //上报从机信息
{
    if (now - var->lc_info.last_time < var->lc_info.interval)
        return;

    int update = get_lc_info(&var->lc_info.ems_info);
    for( int i = 0; i < var->lc_info.lc_num; i++)
    {
        update += get_lc_info(&var->lc_info.lc_list[i]);
    }
    if (update > 0)
    {
        if (mqtt_dp_publish_lc_info_message(var) < 0)
        {
            ems_syslog(LOG_ERR, "mqtt_dp_publish_lc_info_message failed");
        }
    }
    var->lc_info.last_time = now;
}

static void* mqtt_dp_cycle_task(void* para)
{
    mqtt_dp_var_t* var = (mqtt_dp_var_t*)para;
    if (var == NULL) return NULL;
    sleep(5);
    time_t now = 0;
    while (1)
    {
        sleep(1);
        now = time(NULL);
        switch (var->status)
        {
            case MQTT_DP_CONNECTING: break;
            case MQTT_DP_LOGINING:
            {
                mqtt_dp_login_task(var, now);
                break;
            }
            case MQTT_DP_DEVREPORTING:
            {
                var->status = MQTT_DP_HEARTBEAT;
                break;
            }
            case MQTT_DP_HEARTBEAT:
            {
                if (get_cabinet_info()->is_slave == 0)
                {
                    mqtt_dp_lc_info_task(var, now);
                    mqtt_dp_config_task(var, now);
                }
                mqtt_dp_heartbeat_task(var, now);
                mqtt_hardware_monitor_task(var, now);
                mqtt_report_activition_task(var);
                break;
            }
        }
    }
    return NULL;
}

void mqtt_dp_start(void)
{
    mqtt_dp_var_t* var = malloc(sizeof(mqtt_dp_var_t));
    if (var == NULL) return;

    if (mqtt_dp_client_init(var) < 0) goto ERROR;

    ems_syslog(LOG_NOTICE, "MQTT cloud CLIENT START!!");
    if (mqtt_session_start(var->session) != 0) goto ERROR;

    if (pthread_create(&var->mqtt_cycle_tid, NULL, mqtt_dp_cycle_task, var) < 0)
    {
        ems_syslog(LOG_ERR, "mqtt_dp_cycle_task create failed!!");
        goto ERROR;
    }
    else
    {
        pthread_setname_np(var->mqtt_cycle_tid, "mqtt_task_dp");
    }

    return;
ERROR:
    ems_syslog(LOG_ERR, "MQTT DP CLIENT START FAILED!!");
    if (var) destory_mqtt_dp_var_t(var);
}
