#include "mqtt_north.h"
#include "jsonct.h"
#include "../common.h"

void resp_read_dev_info(north_mqtt_var_t *var,mqtt_dev_info_t *dev_info,base_frame_head_t *head_info,analsys_dev_info_t *analsys_dev_info) ;
int creat_head_info(north_mqtt_var_t *var,base_frame_head_t *head_info,char *message);


char *mqtt_type_string[] = {
    "U16",
    "S16",
    "U32",
    "S32",
    "FLOAT",
    "BI16"
};


MQTT_DATA_TYPE get_mqtt_type_by_string(char *typeString)
{
    for (MQTT_DATA_TYPE i = MQTT_TYPE_U16 ; i < MQTT_TYPE_MAX ; i++)
    {
        if(strcmp(typeString,mqtt_type_string[i])==0)
        return i;
    }
    ems_syslog(LOG_WARNING, "MQTT_TYPE_ERROR");
    return TYPE_ERROR;
}


// {
//     "dev":[
//         {
//             "dev_id":"pcs456", //转成mqtt topic里面的字段
//             "module":"ems" ,//属于哪个系统吧
//             "type":"pcsdata",//设备类别
//             "blocks":
//             [
//                 //操作写寄存器，如果数采配置了只读，返回错误
//                 {
//                     "block_id":"0x8021",
//                     "cycle":1000,//ms
//                     //"start_sub":1, 永远从1开始
//                     "sub_id":
//                     [
//                         {"type":"S16","tag":"Power", "scale":10},//默认从地址1开始，自动向后累加
//                         {"type":"S16", "tag":"Current", "scale":10},//如果需要空开，必须指定占位长度
//                         {"type":"S16","tag":"Status","scale":10}
//                     ]
//                 },
//                 {
//                     "block_id":"0x8021",
//                     "cycle":1000,//ms
//                     //"start_sub":1, 永远从1开始
//                     "sub_id":
//                     [
//                         {"type":"S16","tag":"Power1", "scale":10}//默认从地址1开始，自动向后累加
//                     ]
//                 }    
//             ]
//         }
//     ]
// }

int paser_north_mqtt_sub(mqtt_block_info_t *mqtt_block_info,cJSON *nodes,char *dev_id)
{
    int node_cnt = cJSON_GetArraySize(nodes);
    char tmp_c[50] = {0},tmp_c2[64] = {0};
    for(int i = 0;i<node_cnt;i++)
    {
        cJSON *node = cJSON_GetArrayItem(nodes, i);
        if (node)
        {
            mqtt_block_info->mqtt_sub_info[mqtt_block_info->mqtt_sub_num] = calloc(1,sizeof(mqtt_sub_info_t));
            memset(tmp_c,0,sizeof(tmp_c));
            GET_JSON_STRING(node,tmp_c,"type");
            mqtt_block_info->mqtt_sub_info[mqtt_block_info->mqtt_sub_num]->data_type = get_mqtt_type_by_string(tmp_c);
            GET_JSON_FLOAT(node,mqtt_block_info->mqtt_sub_info[mqtt_block_info->mqtt_sub_num]->scale,"scale");
            memset(tmp_c,0,sizeof(tmp_c));
            GET_JSON_STRING(node,tmp_c,"tag");
            sprintf(tmp_c2,"%s.%s",dev_id,tmp_c);
            mqtt_block_info->mqtt_sub_info[mqtt_block_info->mqtt_sub_num]->tag = find_tag_by_name(get_proto_forward_var(),tmp_c2);
            strcpy(mqtt_block_info->mqtt_sub_info[mqtt_block_info->mqtt_sub_num]->name,tmp_c2);
            ems_syslog(LOG_INFO, "[subnum:%d]:tagName:%s data_type:%d,scale:%f",i,tmp_c2,
                mqtt_block_info->mqtt_sub_info[mqtt_block_info->mqtt_sub_num]->data_type,
                mqtt_block_info->mqtt_sub_info[mqtt_block_info->mqtt_sub_num]->scale);

            mqtt_block_info->mqtt_sub_num++;
        }
    }
    return 0;
}

int paser_north_mqtt_blocks(mqtt_dev_info_t *mqtt_dev_info,cJSON *nodes,char *dev_id)
{
    int node_cnt = cJSON_GetArraySize(nodes);
    char tmp_c[64] = {0};
    cJSON *tmp = NULL;
    for(int i = 0;i<node_cnt;i++)
    {
        cJSON *node = cJSON_GetArrayItem(nodes, i);
        if (node)
        {
            mqtt_dev_info->mqtt_block_info[mqtt_dev_info->mqtt_block_num] = calloc(1,sizeof(mqtt_block_info_t));

            memset(tmp_c,0,sizeof(tmp_c));
            GET_JSON_STRING(node, tmp_c, "block_id");
            mqtt_dev_info->mqtt_block_info[mqtt_dev_info->mqtt_block_num]->block_id=analysis_addr(tmp_c) ;
            GET_JSON_INT(node,mqtt_dev_info->mqtt_block_info[mqtt_dev_info->mqtt_block_num]->cycle,"cycle");
            mqtt_dev_info->mqtt_block_info[mqtt_dev_info->mqtt_block_num]->dev_info = mqtt_dev_info;

            ems_syslog(LOG_INFO, "[subnum:%d]:block_id:%X,cycle:%d",i,
             mqtt_dev_info->mqtt_block_info[mqtt_dev_info->mqtt_block_num]->block_id,
             mqtt_dev_info->mqtt_block_info[mqtt_dev_info->mqtt_block_num]->cycle);
             
            tmp = cJSON_GetObjectItemCaseSensitive(node, "sub_id");
            if (tmp)
                paser_north_mqtt_sub(mqtt_dev_info->mqtt_block_info[mqtt_dev_info->mqtt_block_num],tmp,dev_id);
            
            mqtt_dev_info->mqtt_block_num++;
        }
    }
    return 0;
}
int paser_north_mqtt_file(north_mqtt_var_t *pcfg, const char *pdata)
{
    cJSON *root = NULL;
    cJSON *nodes = NULL;
    cJSON *tmp = NULL;
    int ret = 0;
    root = json_parse_string_with_comments(pdata);
    if (!root)
    {
        ems_syslog(LOG_ERR, "parse cfg file error");
        ret = -1;
        goto out;
    }

    nodes = cJSON_GetObjectItemCaseSensitive(root, "dev");
    if (!nodes)
    {
        ems_syslog(LOG_ERR, "get nodes_cfg failed");
        ret = -1;
        goto out;
    }
    int node_cnt = cJSON_GetArraySize(nodes);
    for(int i = 0;i<node_cnt;i++)
    {
        cJSON *node = cJSON_GetArrayItem(nodes, i);
        if (node)
        {
            pcfg->mqtt_dev_info[pcfg->mqtt_dev_num] = calloc(1,sizeof(mqtt_dev_info_t));
            GET_JSON_VALUE_STRING(node, "dev_id", pcfg->mqtt_dev_info[pcfg->mqtt_dev_num]->dev_id);
            GET_JSON_VALUE_STRING(node, "module", pcfg->mqtt_dev_info[pcfg->mqtt_dev_num]->module);
            GET_JSON_VALUE_STRING(node, "type", pcfg->mqtt_dev_info[pcfg->mqtt_dev_num]->type);

            ems_syslog(LOG_NOTICE, "[node_cnt:%d]dev_id:%s,module:%s,type:%s ",i,
            pcfg->mqtt_dev_info[pcfg->mqtt_dev_num]->dev_id,
            pcfg->mqtt_dev_info[pcfg->mqtt_dev_num]->module,
            pcfg->mqtt_dev_info[pcfg->mqtt_dev_num]->type);

            tmp = cJSON_GetObjectItemCaseSensitive(node, "blocks");
            if (tmp)
            {
                paser_north_mqtt_blocks(pcfg->mqtt_dev_info[pcfg->mqtt_dev_num],tmp,pcfg->mqtt_dev_info[pcfg->mqtt_dev_num]->dev_id);
            }      
            pcfg->mqtt_dev_num++;
        }
    }
out:
    cJSON_Delete(root);
    return ret;
}




int load_mqtt_north_info(north_mqtt_var_t *pcfg)
{
    char *pdata = NULL;
    int ret = 0;
    pdata = read_file_data(NORTH_MQTT_INFO_FILE);

    if (pdata == NULL)
    {
        ems_syslog(LOG_ERR, "read cfg file %s error", NORTH_MQTT_INFO_FILE);
        return -1;
    }
    ret = paser_north_mqtt_file(pcfg, pdata);

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

/**
 * set
 * 
 * ipc/E4B853F333/ems/pcs456/pcsdata/Set/binaryFlow
 * 
 * 0022D31C
 * 00
 * 00
 * 001B   //27
 * 01
 * 00008031
 * 0001
 * 0006
 * 0010
 * 00030000000000000000000000000001
 * 
*/

/**
 * 
 * setresp
 * 
 * ipc/E4B853F333/ems/pcs456/pcsdata/SetResp/binaryFlow
 * 
 * 
 * 00230098
 * 00
 * 00
 * 0015
 * 01
 * 00000000
 * 0001
 * 000A
 * 000A
 * 00000000000000000000
 * 
*/

int get_head_info(north_mqtt_var_t *var,base_frame_head_t *head_info,char *message,int len)
{
    int message_index = 0;
    head_info->TransId.v_c[3] = message[message_index++];
    head_info->TransId.v_c[2] = message[message_index++];
    head_info->TransId.v_c[1] = message[message_index++];
    head_info->TransId.v_c[0] = message[message_index++];
    head_info->DevLen = message[message_index++];
    memcpy(head_info->DevMes,&message[message_index],head_info->DevLen);
    message_index+=head_info->DevLen;

    head_info->CmdLen = message[message_index++];
    memcpy(head_info->CMD,&message[message_index],head_info->CmdLen);
    message_index+=head_info->CmdLen;

    head_info->PayloadLen.v_c[1] = message[message_index++];
    head_info->PayloadLen.v_c[0] = message[message_index++];

    ems_syslog(LOG_INFO,"TransId:%X,devmesLen:%d,cmdLen:%d,payLoadLen:%d,message_index:%d",
    head_info->TransId.v_ui,head_info->DevLen,head_info->CmdLen,head_info->PayloadLen.v_us,message_index);
    if(head_info->PayloadLen.v_us!=(len-message_index))
    {
        ems_syslog(LOG_NOTICE,"length error!");
        return -1;
    }
    return message_index;
}

mqtt_dev_info_t *find_dev_info_p_by_id(north_mqtt_var_t *var,char *deviceId)
{
    for (size_t i = 0; i < var->mqtt_dev_num; i++)
    {
        if (strcmp(var->mqtt_dev_info[i]->dev_id,deviceId)==0 )
        {
            return var->mqtt_dev_info[i];
        }
    }
    return NULL;
}

int analsys_payload_info(north_mqtt_var_t *var,analsys_dev_info_t *analsys_dev_info,char *message,int payloadLen)
{
    int index = 0;
    analsys_dev_info->block_num = message[index++];
    ems_syslog(LOG_INFO,"block_num:%d ",analsys_dev_info->block_num);
    for (size_t i = 0; i < analsys_dev_info->block_num; i++)
    {
        analsys_dev_info->block_info[i].blockId.v_c[3]=message[index++];
        analsys_dev_info->block_info[i].blockId.v_c[2]=message[index++];
        analsys_dev_info->block_info[i].blockId.v_c[1]=message[index++];
        analsys_dev_info->block_info[i].blockId.v_c[0]=message[index++];

        analsys_dev_info->block_info[i].StartSubId.v_c[1]=message[index++];
        analsys_dev_info->block_info[i].StartSubId.v_c[0]=message[index++];

        analsys_dev_info->block_info[i].SubNum.v_c[1]=message[index++];
        analsys_dev_info->block_info[i].SubNum.v_c[0]=message[index++];

        analsys_dev_info->block_info[i].Bytes.v_c[1]=message[index++];
        analsys_dev_info->block_info[i].Bytes.v_c[0]=message[index++];

        analsys_dev_info->block_info[i].data_pr=&message[index];
        index+=analsys_dev_info->block_info[i].Bytes.v_us;

        ems_syslog(LOG_INFO,"id_i:%X,StartSubId:%d,SubNum:%d,Bytes:%d!",analsys_dev_info->block_info[i].blockId.v_ui,analsys_dev_info->block_info[i].StartSubId.v_us,\
        analsys_dev_info->block_info[i].SubNum.v_us,analsys_dev_info->block_info[i].Bytes.v_us);
        if(payloadLen<index)
        {
            ems_syslog(LOG_WARNING,"index>payloadLen error!");
            return -1;
        }
    }
    return index;
}


int write_tag_data_by_type(MQTT_DATA_TYPE type,tag_t *tag,char *data,float scale,bool clearbit)
{
    int index = 0;
    // static int takeBit = false;
    // static int bitIndex = 0;
    // if(clearbit)
    // {
    //     if(takeBit)
    //     {
    //         index+=2;
    //         takeBit = false;
    //         bitIndex = 0;
    //         return index;
    //     }
    //     return 0;
    // }
    // if((type!=MQTT_TYPE_BI16)&&(takeBit))
    // {
    //     index+=2;
    //     takeBit = false;
    //     bitIndex = 0;
    // }
    // else if((type==MQTT_TYPE_BI16)&&(takeBit)&&(bitIndex==16))
    // {
    //     index+=2;
    //     takeBit = false;
    //     bitIndex = 0;
    // }
    ems_syslog(LOG_INFO,"type:%d!!",type);
    switch (type)
    {
    case MQTT_TYPE_U16:
        {
            UNION_USHORT data_t;
            data_t.v_c[1]=data[index++];
            data_t.v_c[0]=data[index++];
            
            tag_value_t tag_val;
            tag_val.value.to_int = data_t.v_us;
            tag_val.type = TYPE_TAG_INT;

            write_data_to_tag_by_p(tag,&tag_val,scale);
        }
        break;
    case MQTT_TYPE_S16:
        {

            UNION_SHORT data_t;
            data_t.v_c[1]=data[index++];
            data_t.v_c[0]=data[index++];
            
            tag_value_t tag_val;
            tag_val.value.to_int = data_t.v_s;
            tag_val.type = TYPE_TAG_INT;

            write_data_to_tag_by_p(tag,&tag_val,scale);
        }
        break;
    case MQTT_TYPE_U32:
        {

            UNION_UINT data_t;
            data_t.v_c[3]=data[index++];
            data_t.v_c[2]=data[index++];
            data_t.v_c[1]=data[index++];
            data_t.v_c[0]=data[index++];
            tag_value_t tag_val;
            tag_val.value.to_int = data_t.v_ui;
            tag_val.type = TYPE_TAG_INT;

            write_data_to_tag_by_p(tag,&tag_val,scale);
        }
        break;
    case MQTT_TYPE_S32:
        {

            UNION_INT data_t;
            data_t.v_c[3]=data[index++];
            data_t.v_c[2]=data[index++];
            data_t.v_c[1]=data[index++];
            data_t.v_c[0]=data[index++];
            tag_value_t tag_val;
            tag_val.value.to_int = data_t.v_i;
            tag_val.type = TYPE_TAG_INT;

            write_data_to_tag_by_p(tag,&tag_val,scale);
        }
        break;
    case MQTT_TYPE_FLOAT:
        {

            UNION_FLOAT data_t;
            data_t.v_c[3]=data[index++];
            data_t.v_c[2]=data[index++];
            data_t.v_c[1]=data[index++];
            data_t.v_c[0]=data[index++];
            tag_value_t tag_val;
            tag_val.value.to_float = data_t.v_f;
            tag_val.type = TYPE_TAG_FLOAT;

            write_data_to_tag_by_p(tag,&tag_val,scale);
        }
        break;
        
    // case MQTT_TYPE_BI16:
    //     {  
    //         UNION_SHORT data_t;
    //         data_t.v_c[1]=data[index];
    //         data_t.v_c[0]=data[index+1];
            
    //         tag_value_t tag_val;
    //         tag_val.value.to_int = 0;
    //         tag_val.type = TYPE_TAG_INT;
    //         if (CHK_AMASK_BIT(data_t.v_s, bitIndex))
    //         {
    //             tag_val.value.to_int = 0;
    //         }
    //         write_data_to_tag_by_p(tag,&tag_val,scale);
    //         bitIndex++;
    //     }
        break;
    default:
        break;
    }
    for (size_t i = 0; i < index; i++)
    {
        ems_syslog(LOG_INFO,"data:%x!!",data[i]);
    }
    return index;
}
int get_write_error_resp_Conten_data(char *message)
{
    UNION_UINT errorcode;
    
    unsigned char index = 0;

    message[index++] = 0; //bits

    message[index++] = 0; //CMD
    message[index++] = 0;


    errorcode.v_ui = 1;

    message[index++] = errorcode.v_c[3]; //error code
    message[index++] = errorcode.v_c[2];
    message[index++] = errorcode.v_c[1];
    message[index++] = errorcode.v_c[0];

    message[index++] = 0; //reseved
    message[index++] = 0;
    message[index++] = 0;
    message[index++] = 0;

    message[0] = index-1;//bits
    return index;
}

int get_write_resp_Conten_data(north_mqtt_var_t *var,mqtt_block_info_t *mqtt_block_info,unsigned short data_len,char *data,char *message,int start_sub_id,int num)
{
    UNION_UINT errorcode;
    
    unsigned short index = 0,data_index = 0;

    message[index++] = 0; //bits
    message[index++] = 0;

    message[index++] = 0; //CMD
    message[index++] = 0;


    errorcode.v_ui = 0;
    for (int i = start_sub_id-SUB_INDEX_ADDR_OFFSET; i < MIN(start_sub_id-SUB_INDEX_ADDR_OFFSET+num,mqtt_block_info->mqtt_sub_num); i++)
    {
        ems_syslog(LOG_INFO,"sub_id:%d name:%s!!",i,mqtt_block_info->mqtt_sub_info[i]->name);
        data_index+=write_tag_data_by_type(mqtt_block_info->mqtt_sub_info[i]->data_type,mqtt_block_info->mqtt_sub_info[i]->tag,&data[data_index],mqtt_block_info->mqtt_sub_info[i]->scale,false);
        if(data_index > data_len) 
        {
            ems_syslog(LOG_WARNING," data_index > data_len error!");
            break;
        }
    }
    // data_index+=write_tag_data_by_type(0,0,0,0,true);
    if(data_index > data_len) 
    {
        ems_syslog(LOG_WARNING," data_index > data_len error!");
    }

    message[index++] = errorcode.v_c[3]; //error code
    message[index++] = errorcode.v_c[2];
    message[index++] = errorcode.v_c[1];
    message[index++] = errorcode.v_c[0];

    message[index++] = 0; //reseved
    message[index++] = 0;
    message[index++] = 0;
    message[index++] = 0;

    unsigned short payloar_len = index -2;
    message[0] = (payloar_len >> 8)&0x00ff;//负载长度
    message[1] = (payloar_len     )&0x00ff;
    return index;
}

int framing_write_blocks_data(north_mqtt_var_t *var,mqtt_dev_info_t *dev_info,analsys_dev_info_t *analsys_dev_info,char *message)
{
    unsigned short message_index = 0;

    message[message_index++] = 0;//负载长度
    message[message_index++] = 0;

    message[message_index++] = analsys_dev_info->block_num;
    ems_syslog(LOG_INFO,"block_num:%d ",analsys_dev_info->block_num);
    for (size_t i = 0; i < analsys_dev_info->block_num; i++)
    {
        int finded = 0;
        for (size_t j = 0; j < dev_info->mqtt_block_num; j++)
        {
            if(dev_info->mqtt_block_info[j]->block_id==analsys_dev_info->block_info[i].blockId.v_ui)
            {
                finded = 1;
                message[message_index++] = analsys_dev_info->block_info[i].blockId.v_c[3]; //blockid
                message[message_index++] = analsys_dev_info->block_info[i].blockId.v_c[2];
                message[message_index++] = analsys_dev_info->block_info[i].blockId.v_c[1];
                message[message_index++] = analsys_dev_info->block_info[i].blockId.v_c[0];

                message[message_index++] = analsys_dev_info->block_info[i].StartSubId.v_c[1]; //startsubid
                message[message_index++] = analsys_dev_info->block_info[i].StartSubId.v_c[0];

                message[message_index++] = analsys_dev_info->block_info[i].SubNum.v_c[1];
                message[message_index++] = analsys_dev_info->block_info[i].SubNum.v_c[0]; //subnum

                message_index+= get_write_resp_Conten_data(
                    var,
                    dev_info->mqtt_block_info[j],
                    analsys_dev_info->block_info[i].Bytes.v_us,
                    analsys_dev_info->block_info[i].data_pr,
                    &message[message_index],
                    analsys_dev_info->block_info[i].StartSubId.v_us,
                    analsys_dev_info->block_info[i].SubNum.v_us
                    );

                break;
            }
        }
        if (finded == 0)
        {
            message_index+= get_write_error_resp_Conten_data(message);
        }
    }
    unsigned short payloar_len = message_index -2;
    message[0] = (payloar_len >> 8)&0x00ff;//负载长度
    message[1] = (payloar_len     )&0x00ff;
    return message_index;
}

void resp_write_dev_info(north_mqtt_var_t *var,mqtt_dev_info_t *dev_info,base_frame_head_t *head_info,analsys_dev_info_t *analsys_dev_info)  //读响应
{
    char pub_topic[MAX_BASE_SHORT_LEN * 5] = {0};
    sprintf(pub_topic,"ipc/%s/%s/%s/%s/SetResp/binaryFlow",var->topic_sn,dev_info->module,dev_info->dev_id,dev_info->type);
    char message[1024*2] = {0};
    int index = 0;
    index+=creat_head_info(var,head_info,message);
    index+=framing_write_blocks_data(var,dev_info,analsys_dev_info,&message[index]);

    ems_syslog_hex(LOG_INFO,message,index,"resp_write_dev_info:");

    if(mqtt_session_publish(var->session, pub_topic,message, index)!=0)
        ems_syslog(LOG_WARNING,"send data error!!");
}


//ipc/E4B853F333/bms/bau123/baudata/GetResp/binaryFlow
/**
 * 0013D4E8
 * 00
 * 00
 * 0029
 * 01
 * 00000002
 * 1001
 * 001E
 * 001E
 * 02C00619049807A80620054E060903F2055D05FD03560761035A06450555
 * 
 * 
 * demo
 * 
 * get
 * 0000 0337
 * 00
 * 00
 * 000b
 * 01
 * 00008021
 * 0001
 * 0022
 * 0000
 * 
 * getresp
 * 0000 0337
 * 00
 * 00
 * 007d
 * 01
 * 00008021
 * 0001
 * 0022
 * 0072
 * 0000000000000000000000000000000000000000000OOD000OOO0O000O0000000O00O000000000000000000000000000000000000000000000000000000000000000000000000000 0000 00 0000 0000 0000 0000 0000 0000 0000 000 0000 0000 0000oo00 0000 0000 00000000000000
 * 
 * 
 * set
 * 0000406a
 * 00
 * 00
 * 000d
 * 01
 * 00008001
 * 0002
 * 0001
 * 0002
 * 0001
 * 
 * setresp
 * 0000406a
 * 00
 * 00
 * 0015
 * 01
 * 00008001
 * 0002
 * 0001
 * 000a
 * 00000000000000000000
 * 
 * 实际的setresp
 * 0033547c
 * 00
 * 00
 * 0015
 * 01
 * 00008031
 * 0001
 * 0006
 * 000a
 * 00 00 00 00 00 00 00 00 00 00
 * 
 * set test
 * 0000406a
 * 00
 * 00
 * 000d
 * 01
 * 00008031
 * 0001
 * 0001
 * 0002
 * 0005
 * 
 * 
 * 0000406a0000000d01000080310001000100020005
 * 
 * gettest
 * 
 * 00000337
 * 00
 * 00
 * 000b
 * 01
 * 00008031
 * 0001
 * 000A
 * 0000
 * 
 * 000003370000000b01000080310001000A0000
 * 
 * 实际返回
 * 00000337
 * 00
 * 00
 * 0023
 * 01
 * 00008031
 * 0001
 * 000A
 * 0018
 * 000300000000000000002774000000000000000000000000
 * 
 * 000003370000002301000080310001000A0018000300000000000000002774000000000000000000000000
 * 
 * 发送000003370000000b0100008031000200090000
 * 返回
 * 
*/
#define MAX_MESSAGE_LEN 1024*2
static unsigned int TransId = 1;

int get_tag_data_by_type(MQTT_DATA_TYPE type,tag_t *tag,char *dist_data,float scale,bool clearbit)
{
    int index = 0;
    ems_syslog(LOG_INFO,"type:%d!!",type);
    // static int takeBit = false;
    // static int bitIndex = 0;
    // if(clearbit)
    // {
    //     if(takeBit)
    //     {
    //         index+=2;
    //         takeBit = false;
    //         bitIndex = 0;
    //         return index;
    //     }
    //     return 0;
    // }
    // if((type!=MQTT_TYPE_BI16)&&(takeBit))
    // {
    //     index+=2;
    //     takeBit = false;
    //     bitIndex = 0;
    // }
    // else if((type==MQTT_TYPE_BI16)&&(takeBit)&&(bitIndex==16))
    // {
    //     index+=2;
    //     takeBit = false;
    //     bitIndex = 0;
    // }
    switch (type)
    {
    case MQTT_TYPE_U16:
        {
            
            unsigned short data_t = 0;
            tag_value_t tag_val;
            read_data_form_tag_by_p(tag,TYPE_TAG_INT,&tag_val,scale);
            data_t = tag_val.value.to_int;
            dist_data[index++] = (data_t>>8)&0x00ff;
            dist_data[index++] = data_t&0x00ff;
        }
        break;
    case MQTT_TYPE_S16:
        {

            short data_t = 0;
            tag_value_t tag_val;
            read_data_form_tag_by_p(tag,TYPE_TAG_INT,&tag_val,scale);
            data_t = tag_val.value.to_int;
            dist_data[index++] = (data_t>>8)&0x00ff;
            dist_data[index++] = data_t&0x00ff;
        }
        break;
    case MQTT_TYPE_U32:
        {

            unsigned int data_t = 0;
            tag_value_t tag_val;
            read_data_form_tag_by_p(tag,TYPE_TAG_INT,&tag_val,scale);
            data_t = tag_val.value.to_int;
            dist_data[index++] = (data_t>>24)&0x00ff;
            dist_data[index++] = (data_t>>16)&0x00ff;
            dist_data[index++] = (data_t>>8)&0x00ff;
            dist_data[index++] = data_t&0x00ff;
        }
        break;
    case MQTT_TYPE_S32:
        {

            int data_t = 0;
            tag_value_t tag_val;
            read_data_form_tag_by_p(tag,TYPE_TAG_INT,&tag_val,scale);
            data_t = tag_val.value.to_int;
            dist_data[index++] = (data_t>>24)&0x00ff;
            dist_data[index++] = (data_t>>16)&0x00ff;
            dist_data[index++] = (data_t>>8)&0x00ff;
            dist_data[index++] = data_t&0x00ff;
        }
        break;
    case MQTT_TYPE_FLOAT:
        {

            float data_t = 0;
            tag_value_t tag_val;
            read_data_form_tag_by_p(tag,TYPE_TAG_FLOAT,&tag_val,scale);
            data_t = tag_val.value.to_float;
            dist_data[index++] = ((unsigned char *)&data_t)[3];
            dist_data[index++] = ((unsigned char *)&data_t)[2];
            dist_data[index++] = ((unsigned char *)&data_t)[1];
            dist_data[index++] = ((unsigned char *)&data_t)[0];
        }break;
    // case MQTT_TYPE_BI16:
    //     {  
    //         unsigned short data_t = 0;
    //         tag_value_t tag_val;
    //         read_data_form_tag_by_p(tag,TYPE_TAG_INT,&tag_val,scale);

    //         ((unsigned char *)&data_t)[0] = dist_data[index+1];
    //         ((unsigned char *)&data_t)[1] = dist_data[index];
    //         if (tag_val.value.to_int == 1)
    //         {
    //             SET_AMASK_BIT(data_t, bitIndex);
    //         }
    //         dist_data[index] = (data_t>>8)&0x00ff;
    //         dist_data[index+1] = data_t&0x00ff;

    //         bitIndex++;
    //     }
    //     break;
    default:
        break;
    }
    return index;
}
int get_resp_Conten_data(north_mqtt_var_t *var,mqtt_block_info_t *mqtt_block_info,char *message,unsigned short start_sub_id,unsigned short num)
{
    unsigned short index = 0;
    message[index++] = 0; //bytes
    message[index++] = 0;
    ems_syslog(LOG_INFO,"start_sub_id:%d ,%d!!",start_sub_id,num);
    for (int i = start_sub_id-SUB_INDEX_ADDR_OFFSET; i < MIN(start_sub_id-SUB_INDEX_ADDR_OFFSET+num,mqtt_block_info->mqtt_sub_num); i++)
    {
        ems_syslog(LOG_INFO,"sub_id:%d name:%s!!",i,mqtt_block_info->mqtt_sub_info[i]->name);
        index+=get_tag_data_by_type(mqtt_block_info->mqtt_sub_info[i]->data_type,mqtt_block_info->mqtt_sub_info[i]->tag,&message[index],mqtt_block_info->mqtt_sub_info[i]->scale,false);
    }
    // index+=get_tag_data_by_type(0,0,0,0,true);
    unsigned short payloar_len = index -2;
    message[0] = (payloar_len >> 8)&0x00ff;//负载长度
    message[1] = (payloar_len     )&0x00ff;
    return index;
}


int creat_head_info(north_mqtt_var_t *var,base_frame_head_t *head_info,char *message)
{
    int message_index = 0;
    unsigned int transid_t = 0;
    if (head_info->TransId.v_ui > 0)
    {
        transid_t = head_info->TransId.v_ui;
    }
    else
    {
        transid_t = TransId;
        TransId++;
    }
    // ems_syslog(LOG_WARNING,"transid_t:%X!!",transid_t);

    message[message_index++] = (transid_t>>24)&0x000000ff;
    message[message_index++] = (transid_t>>16)&0x000000ff;
    message[message_index++] = (transid_t>>8 )&0x000000ff;
    message[message_index++] = (transid_t    )&0x000000ff;

    message[message_index++] = head_info->DevLen;
    memcpy(&message[message_index],head_info->DevMes,head_info->DevLen);
    message_index+=head_info->DevLen;

    message[message_index++] = head_info->CmdLen;//预留指令长度
    memcpy(&message[message_index],head_info->CMD,head_info->CmdLen);
    message_index+=head_info->CmdLen;
    ems_syslog(LOG_INFO,"head len:%d!!",message_index);
    return message_index;
}
int framing_blocks_data(north_mqtt_var_t *var,mqtt_dev_info_t *dev_info,analsys_dev_info_t *analsys_dev_info,char *message)
{
    unsigned short message_index = 0;

    message[message_index++] = 0;//负载长度
    message[message_index++] = 0;

    message[message_index++] = analsys_dev_info->block_num;
    for (size_t i = 0; i < analsys_dev_info->block_num; i++)
    {
        for (size_t j = 0; j < dev_info->mqtt_block_num; j++)
        {   
            if(dev_info->mqtt_block_info[j]->block_id==analsys_dev_info->block_info[i].blockId.v_ui)
            {
                ems_syslog(LOG_NOTICE,"block_id:[%X]",dev_info->mqtt_block_info[j]->block_id);
                message[message_index++] = analsys_dev_info->block_info[i].blockId.v_c[3]; //blockid
                message[message_index++] = analsys_dev_info->block_info[i].blockId.v_c[2];
                message[message_index++] = analsys_dev_info->block_info[i].blockId.v_c[1];
                message[message_index++] = analsys_dev_info->block_info[i].blockId.v_c[0];

                message[message_index++] = analsys_dev_info->block_info[i].StartSubId.v_c[1]; //startsubid
                message[message_index++] = analsys_dev_info->block_info[i].StartSubId.v_c[0];

                message[message_index++] = analsys_dev_info->block_info[i].SubNum.v_c[1];
                message[message_index++] = analsys_dev_info->block_info[i].SubNum.v_c[0]; //subnum

                message_index+= get_resp_Conten_data(
                    var,
                    dev_info->mqtt_block_info[j],
                    &message[message_index],
                    analsys_dev_info->block_info[i].StartSubId.v_us,
                    analsys_dev_info->block_info[i].SubNum.v_us
                    );
                break;
            }
        }
    }
    unsigned short payloar_len = message_index -2;
    message[0] = (payloar_len >> 8)&0x00ff;//负载长度
    message[1] = (payloar_len     )&0x00ff;

    ems_syslog(LOG_INFO,"payload len:%d!!",message_index);
    return message_index;
}
/**
 * 
 * 00000003
 * 06
 * 504353343536
 * 00
 * 0020
 * 01
 * 00008013
 * 0001
 * 000A
 * 1600
 * 000000000000000000000000000000000000000000
 * 
*/
void resp_read_dev_info(north_mqtt_var_t *var,mqtt_dev_info_t *dev_info,base_frame_head_t *head_info,analsys_dev_info_t *analsys_dev_info)  //读响应
{
    char pub_topic[MAX_BASE_SHORT_LEN * 5] = {0};
    sprintf(pub_topic,"ipc/%s/%s/%s/%s/GetResp/binaryFlow",var->topic_sn,dev_info->module,dev_info->dev_id,dev_info->type);
    char message[1024*2] = {0};
    int index = 0;
    index+=creat_head_info(var,head_info,message);
    index+=framing_blocks_data(var,dev_info,analsys_dev_info,&message[index]);

    ems_syslog_hex(LOG_INFO,message,index,"resp_read_dev_info:");

    if(mqtt_session_publish(var->session, pub_topic,message, index)!=0)
        ems_syslog(LOG_WARNING,"send data error!!");
}


void *mqtt_cycle_task(void *data)
{
    north_mqtt_var_t *var = data;
    long long time = 0;
    while (1)
    {
        sleep(1);
        if(mqtt_session_get_state(var->session)!=MQTT_CONNECTED)
        {
            ems_syslog(LOG_WARNING,"mqtt no connected");
            continue;
        }
        time++;
        ems_syslog(LOG_INFO,"mqtt_cycle_task is running time:%lld!",time);
        for (size_t i = 0; i < var->mqtt_dev_num; i++)
        {
            for (size_t j = 0; j < var->mqtt_dev_info[i]->mqtt_block_num; j++)
            {
                // ems_syslog(LOG_NOTICE,"block_id:%X,cycle:%d",
                // var->mqtt_dev_info[i]->mqtt_block_info[j]->block_id,
                // var->mqtt_dev_info[i]->mqtt_block_info[j]->cycle
                // );

                if((var->mqtt_dev_info[i]->mqtt_block_info[j]->cycle!=0)&&(time%(var->mqtt_dev_info[i]->mqtt_block_info[j]->cycle/1000)==0))
                {
                    base_frame_head_t head_info;
                    head_info.TransId.v_ui = 0;
                    // head_info.DevLen = strlen(var->mqtt_dev_info[i]->dev_id);
                    // strcpy(head_info.DevMes,var->mqtt_dev_info[i]->dev_id);
                    head_info.CmdLen = 0;

                    analsys_dev_info_t analsys_dev_info;
                    analsys_dev_info.block_num = 1;
                    analsys_dev_info.block_info[0].blockId.v_ui = var->mqtt_dev_info[i]->mqtt_block_info[j]->block_id;
                    analsys_dev_info.block_info[0].StartSubId.v_us = SUB_INDEX_ADDR_OFFSET;
                    analsys_dev_info.block_info[0].SubNum.v_us = var->mqtt_dev_info[i]->mqtt_block_info[j]->mqtt_sub_num;
                    resp_read_dev_info(var,var->mqtt_dev_info[i],&head_info,&analsys_dev_info);
                }
            }
        }
    }
}

static int mqtt_handle_recv_msg(void *obj, mqtt_message_t *mqtt_msg)
{
    north_mqtt_var_t *var = obj;
    ems_syslog(LOG_INFO, "received MQTT topic:%s payload length:%d", mqtt_msg->topic, mqtt_msg->payloadLen);
    char *get_data = mqtt_msg->payload;
    int len = mqtt_msg->payloadLen;
    char device_id[64] = {0},tmp_c[64] = {0};
    //                  ipc/21881FFF0028/ems/PCS456/pcsdata/Set/binaryFlow
    sscanf(mqtt_msg->topic,"ipc/%[^/]/%[^/]/%[^/]/%[^/]/%[^/]/binaryFlow",tmp_c,tmp_c,device_id,tmp_c,tmp_c);
    ems_syslog(LOG_NOTICE,"device:[%s]",device_id);
    mqtt_dev_info_t *device_info = find_dev_info_p_by_id(var,device_id);

    if(device_info != NULL)
    {
        base_frame_head_t head_info;
        // int index = 0;
        ems_syslog_hex(LOG_INFO,get_data,len,"recv_hex:");
        len=get_head_info(var,&head_info,get_data,len);
        if(len<0)
            return 0;
        analsys_dev_info_t analsys_dev_info;
        analsys_payload_info(var,&analsys_dev_info,get_data+len,head_info.PayloadLen.v_us);

        if (strstr(mqtt_msg->topic, "/Set/binaryFlow")) //设置参数
        {
            resp_write_dev_info(var,device_info,&head_info,&analsys_dev_info);
        }
        else if (strstr(mqtt_msg->topic, "/Get/binaryFlow")) //获取数据
        {            
            resp_read_dev_info(var,device_info,&head_info,&analsys_dev_info);
        }
    }
    else //返回失败 不响应
    {
        ems_syslog(LOG_WARNING,"can't find device %s!",device_id);
    }
    return 0;
    
}



static void north_mqtt_subscribe_all(north_mqtt_var_t *var)
{
    mqtt_session_t *session = var->session;
    char sub_topic[MAX_BASE_SHORT_LEN * 5] = {0};
    for (size_t i = 0; i < var->mqtt_dev_num; i++)
    {
        mqtt_dev_info_t *dev_info= var->mqtt_dev_info[i];
        sprintf(sub_topic,"ipc/%s/%s/%s/%s/Set/binaryFlow",var->topic_sn,dev_info->module,dev_info->dev_id,dev_info->type);
        mqtt_session_subscribe(session, sub_topic);
        ems_syslog(LOG_INFO, "sub topic:%s", sub_topic);

        sprintf(sub_topic,"ipc/%s/%s/%s/%s/Get/binaryFlow",var->topic_sn,dev_info->module,dev_info->dev_id,dev_info->type);
        mqtt_session_subscribe(session, sub_topic);
        ems_syslog(LOG_INFO, "sub topic:%s", sub_topic);
    }
}

//在数采设备加载完成后加载 不然无法匹配到点位和订阅信息
north_mqtt_var_t *north_mqtt_client_init(north_mqtt_param *mqttparam)
{
    north_mqtt_var_t *var= calloc(sizeof(north_mqtt_var_t),1);

    if(load_mqtt_north_info(var)<0)
    {
        if(var!=NULL)
        {
            free(var);
            var = NULL;
        }
        return NULL;
    }

    ems_syslog(LOG_NOTICE,"client_id:%s",strlen(mqttparam->client_id)==0?"no_client_id":mqttparam->client_id);
    mqtt_session_t *session =  mqtt_session_new(strlen(mqttparam->client_id)==0?"no_client_id":mqttparam->client_id,var);
    if (session == NULL)
    {
        ems_syslog(LOG_NOTICE,"creat default!!");
        if(var!=NULL)
        {
            free(var);
            var = NULL;
        }
        session = NULL;
        return NULL;
    }
    var->session = session;
    ems_syslog(LOG_NOTICE,"mqtt client creat success!!");
    strcpy(var->topic_sn,strlen(mqttparam->topic_sn)==0?"no_sn":mqttparam->topic_sn);

    ems_syslog(LOG_NOTICE,"topic_sn:%s",strlen(mqttparam->topic_sn)==0?"no_sn":mqttparam->topic_sn);

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

    mqtt_session_set_address(session, (strlen(mqttparam->addr)==0)?INTERNAL_BROKER_ADDR:mqttparam->addr,(mqttparam->port == 0)?1883:mqttparam->port,mqttparam->user, mqttparam->pass);
    mqtt_session_set_opts(session, DEFAULT_IPC_QOS, KEEP_ALIVE_MAX);

    mqtt_session_set_callbacks(session, mqtt_handle_recv_msg,NULL);

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

void mqtt_north_start(north_mqtt_var_t *var)
{
    if((var!=NULL)&&(var->session!=NULL))
    {
        ems_syslog(LOG_NOTICE,"MQTT CLIENT START!!");
        mqtt_session_start(var->session);
        if (0 == pthread_create(&var->mqtt_cycle_tid,NULL,mqtt_cycle_task,var))
        {
            pthread_setname_np(var->mqtt_cycle_tid, "mqttCycleTask");
        }
    }
}