#include <dirent.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <sys/stat.h>
#include <sys/time.h>
#include <sys/types.h>
#include <unistd.h>

#include "ems_smu_common.h"
#include "dy_utils/cJSON.h"
#include "dy_utils/dy_common.h"
#include "mqtt_session.h"
#include "dy_utils/protocol.h"

#define EMS_UPGRADE_FILE "/tmp/ems.update"

static void gen_set_rsp(ems_smu_var_t *var, cJSON *rsp_obj, int result, char *funcid)
{
    char topic[TOPIC_MAX_LEN];
    cJSON_AddStringToObject(rsp_obj, "funcId", funcid);
    cJSON_AddStringToObject(rsp_obj, "lcSN", var->sn_str);
    cJSON_AddNumberToObject(rsp_obj, "seq", var->seq++);
    cJSON_AddNumberToObject(rsp_obj, "time", time(NULL));
    cJSON_AddNumberToObject(rsp_obj, "result", result);

    char *data_tmp = cJSON_PrintUnformatted(rsp_obj);
    snprintf(topic, TOPIC_MAX_LEN, "%s/%s/%s/%s", PLATFORM, SETRSP, var->sn_str, funcid);
    ipc_session_publish(var->session, topic, data_tmp, strlen(data_tmp));
    free(data_tmp);
}

static void gen_get_rsp(ems_smu_var_t *var, char *funcid, cJSON *rsp_obj)
{
    char topic[TOPIC_MAX_LEN] = {0};
    cJSON_AddStringToObject(rsp_obj, "funcId", funcid);
    cJSON_AddStringToObject(rsp_obj, "lcSN", var->sn_str);
    cJSON_AddNumberToObject(rsp_obj, "seq", var->seq++);
    cJSON_AddNumberToObject(rsp_obj, "time", time(NULL));

    char *data_tmp = cJSON_PrintUnformatted(rsp_obj);
    snprintf(topic, TOPIC_MAX_LEN, "%s/%s/%s/%s", PLATFORM, GETRSP, var->sn_str, funcid);
    ipc_session_publish(var->session, topic, data_tmp, strlen(data_tmp));
    free(data_tmp);
}

void gen_set_dev_service(ems_smu_var_t *var, cJSON *rt_data, char *identifier, char *devsn)
{
    char topic[TOPIC_MAX_LEN * 2] = {0};
    cJSON_AddStringToObject(rt_data, "identifier", identifier);
    cJSON_AddNumberToObject(rt_data, "mi",var->mi++);
    cJSON_AddNumberToObject(rt_data, "timestamp", time(NULL));
    cJSON_AddStringToObject(rt_data, "sn", devsn);

    char *data_tmp = cJSON_PrintUnformatted(rt_data);
    snprintf(topic, TOPIC_MAX_LEN * 2, "ipc/%s/NULL/device/%s/data/Set_Rglt", var->sn_str, devsn);

    ipc_session_publish(var->session, topic, data_tmp, strlen(data_tmp));
    free(data_tmp);
}

void gen_set_rglt_service(ems_smu_var_t *var, cJSON *rt_data, char *identifier)
{
    char topic[TOPIC_MAX_LEN * 2] = {0};
    cJSON_AddStringToObject(rt_data, "identifier", identifier);
    cJSON_AddNumberToObject(rt_data, "mi",var->mi++);
    cJSON_AddNumberToObject(rt_data, "timestamp", time(NULL));
    cJSON_AddStringToObject(rt_data, "sn", var->emsctrl_sn_str);
    cJSON_AddNumberToObject(rt_data, "emmsv2", 1);

    char *data_tmp = cJSON_PrintUnformatted(rt_data);
    snprintf(topic, TOPIC_MAX_LEN * 2, "ipc/%s/NULL/device/%s/data/Set_Rglt", var->sn_str, var->emsctrl_sn_str);

    ipc_session_publish(var->session, topic, data_tmp, strlen(data_tmp));
    free(data_tmp);
    cJSON_Delete(rt_data);
}

void ems_smu_loginrsp(ems_smu_var_t *var, mqtt_message_t *mqtt_msg)
{
    int result = 0;
    cJSON *root = cJSON_Parse(mqtt_msg->payload);
    if (root == NULL)
    {
        dbg_syslog(LOG_ERR, "parse mqtt payload failed");
        return;
    }
    
    GET_JSON_VALUE_INT(root, "result", result);
    if (result == STAT_CONNECTED)
        var->state = STAT_CONNECTED;

    cJSON_Delete(root);
}

void ems_smu_nodestructure(ems_smu_var_t *var, mqtt_message_t *mqtt_msg)
{
    char cmd[CMD_SHORT_LENGTH] = {0};
    cJSON *nodes = cJSON_CreateArray();
    cJSON *rsp_obj = cJSON_CreateObject();
    cJSON_AddItemToObject(rsp_obj, "nodes", nodes);
    
    dbg_syslog(LOG_ERR, "return nodestruc");
    snprintf(cmd, CMD_SHORT_LENGTH, "hgetall devs");
    redisReply *reply = (redisReply *)redisCommand(var->context, cmd);
    if (NULL == reply || reply->type != REDIS_REPLY_ARRAY)
    {
        dbg_syslog(LOG_ERR, "Get all devs failed");
    }
    else 
    {
        dbg_syslog(LOG_ERR, "elems %lu", reply->elements);
        for (int i = 0; i <= reply->elements; i+=2)
        {
            dbg_syslog(LOG_ERR, "i %d", i);
            char devno[MAX_DEVTYPE_LEN] = {0};
            if (i < reply->elements)
            {
                dbg_syslog(LOG_ERR, "iterating dev %s %s", reply->element[i]->str, reply->element[i+1]->str);
                char *tok = strtok(reply->element[i+1]->str, "|");
                strncpy(devno, tok, strlen(tok));
            }
            else 
            {
                dbg_syslog(LOG_ERR, "iterating EMS");
                strcpy(devno, "EMS");
            }
            dev_nodestruct_t *curdev;
            HASH_FIND_STR(var->devtypenode, devno, curdev);
            if (NULL == curdev)
            {
                dbg_syslog(LOG_ERR, "devno %s not found in devtypenode", devno);
                continue;
            }

            cJSON *node = cJSON_CreateObject();
            cJSON_AddItemToArray(nodes, node);
            cJSON_AddStringToObject(node, "TYPE", curdev->devtype);
            cJSON_AddStringToObject(node, "DEV_NO", (i < reply->elements) ? (reply->element[i]->str):devno);
            cJSON_AddStringToObject(node, "DISP", curdev->devdisp);
            cJSON *alarmArray = cJSON_CreateArray();
            cJSON_AddItemToObject(node, "ALARM", alarmArray);
            cJSON *statArray = cJSON_CreateArray();
            cJSON_AddItemToObject(node, "STATUS", statArray);
            cJSON *setArray = cJSON_CreateArray();
            cJSON_AddItemToObject(node, "SETTING", setArray);
            cJSON *dataArray = cJSON_CreateArray();
            cJSON_AddItemToObject(node, "DATA", dataArray);

            alarm_entry_t *alarmpos;
            for(alarmpos = curdev->alarmnode; alarmpos != NULL; alarmpos = alarmpos->hh.next)
            {
                cJSON *alarmentry = cJSON_CreateObject();
                cJSON_AddStringToObject(alarmentry, "NAME", alarmpos->entryname);
                cJSON_AddStringToObject(alarmentry, "DISP", alarmpos->entryDisp);
                cJSON_AddNumberToObject(alarmentry, "TYPE", alarmpos->type);
                cJSON_AddStringToObject(alarmentry, "DESC", alarmpos->entryDesc);
                cJSON_AddItemToArray(alarmArray, alarmentry); 
            }            

            set_entry_t *setpos;
            for(setpos = curdev->setnode; setpos!= NULL; setpos = setpos->hh.next)
            {
                cJSON *setentry = cJSON_CreateObject();
                cJSON_AddStringToObject(setentry, "NAME", setpos->entryname);
                cJSON_AddStringToObject(setentry, "DISP", setpos->entryDisp);
                cJSON_AddNumberToObject(setentry, "TYPE", setpos->type);
                cJSON_AddStringToObject(setentry, "DESC", setpos->entryDesc);
                cJSON_AddItemToArray(setArray, setentry);
            }

            stat_entry_t *statpos;
            for(statpos = curdev->statnode; statpos != NULL; statpos = statpos->hh.next)
            {
                cJSON *statentry = cJSON_CreateObject();
                cJSON_AddStringToObject(statentry, "NAME", statpos->entryname);
                cJSON_AddStringToObject(statentry, "DISP", statpos->entryDisp);
                cJSON_AddNumberToObject(statentry, "TYPE", statpos->type);
                cJSON_AddStringToObject(statentry, "DESC", statpos->entryDesc);
                cJSON_AddItemToArray(statArray, statentry);
            }

            data_entry_t *datapos;
            for(datapos = curdev->datanode; datapos != NULL; datapos = datapos->hh.next)
            {
                cJSON *dataentry = cJSON_CreateObject();
                cJSON_AddStringToObject(dataentry, "NAME", datapos->entryname);
                cJSON_AddStringToObject(dataentry, "DISP", datapos->entryDisp);
                cJSON_AddNumberToObject(dataentry, "TYPE", datapos->type);
                cJSON_AddStringToObject(dataentry, "DESC", datapos->entryDesc);
                cJSON_AddItemToArray(dataArray, dataentry);
            }
        }
    }    

    gen_get_rsp(var, FUNC_NODESTRUCTURE, rsp_obj);
    cJSON_Delete(rsp_obj);
    freeReplyObject(reply);
}

void ems_smu_historydata(ems_smu_var_t *var, mqtt_message_t *mqtt_msg)
{
    char name[MAX_NAME_LEN] = {0};
    char stime[MAX_TIME_LEN] = {0};
    char etime[MAX_TIME_LEN] = {0};
    char cmd[CMD_SHORT_LENGTH] = {0};
    char sdate[12] = {0}, edate[12] = {0};
    uint startsec, endsec = 0;
    cJSON *historyData = cJSON_CreateArray();
    cJSON *root = cJSON_Parse(mqtt_msg->payload);
    if (root == NULL)
    {
        dbg_syslog(LOG_ERR, "parse mqtt msg failed");
        return;
    }

    GET_JSON_VALUE_STRING(root, "name", name);
    GET_JSON_VALUE_STRING(root, "startTime", stime);
    GET_JSON_VALUE_STRING(root, "endTime", etime);
    char *tok = strtok(stime, " ");
    strncpy(sdate, tok, strlen(tok));
    char *tok2 = strtok(0, " ");
    startsec = atoi(tok2) * 10000 + atoi(tok2 + 3) * 100 + atoi(tok2 + 6);
    char *tok3 = strtok(etime, " ");
    strncpy(edate, tok3, strlen(tok3));
    char *tok4 = strtok(0, " ");
    endsec = atoi(tok4) * 10000 + atoi(tok4 + 3) * 100 + atoi(tok4 + 6);
    dbg_syslog(LOG_ERR, "endsec %u %s %s %s %s", endsec, tok, tok2, tok3, tok4);
    bool loop_exitflag = false;
    uint esec = 240000;
    if (strncmp(sdate, edate, strlen(edate)) == 0)
    {
        loop_exitflag = true;
        esec = endsec;
    }
    while(true)
    {
        time_t dayts = 0;
        dbg_syslog(LOG_ERR, "name %s sdate %s ssec %u esec %u", name, sdate, startsec, esec);
        snprintf(cmd, CMD_SHORT_LENGTH, "zrangebyscore %s-%s %u %u", name, sdate, startsec, esec);
        redisReply *reply = (redisReply *)redisCommand(var->context, cmd);
        if (NULL == reply || reply->type != REDIS_REPLY_ARRAY || reply->elements <= 0)
        {
            dbg_syslog(LOG_NOTICE, "no data in %s", cmd);
            freeReplyObject(reply);
            break;
        }
        for (int i = 0; i < reply->elements; i++)
        {
            cJSON *dvalue = cJSON_CreateObject();
            char *datatok = strtok(reply->element[i]->str, "|");
            char *tstok = strtok(0, "|");
            cJSON_AddNumberToObject(dvalue, "value", atof(datatok));
            cJSON_AddNumberToObject(dvalue, "time", atoi(tstok));
            cJSON_AddItemToArray(historyData, dvalue);
            if (dayts <= 0)
                dayts = atoi(tstok);
        }
        if (loop_exitflag == true)
            break;
        
        dbg_syslog(LOG_ERR, "dayts is %lu %lu", dayts, sizeof(time_t));
        struct tm* tm_new = localtime(&dayts);
        dbg_syslog(LOG_ERR, "dayts is %lu", dayts);
        snprintf(sdate, 12, "%04d-%02d-%02d", tm_new->tm_year + 1900, tm_new->tm_mon + 1, tm_new->tm_mday + 1);
        dbg_syslog(LOG_ERR, "sdate %s", sdate);
        if (strncmp(sdate, edate, strlen(edate)) == 0)
        {
            esec = endsec;
            loop_exitflag = true;
        }
        startsec = 0;

        freeReplyObject(reply);
    }
    

    cJSON *rsp_obj = cJSON_CreateObject();
    cJSON_AddItemToObject(rsp_obj, "historydata", historyData);
    cJSON_AddStringToObject(rsp_obj, "name", name);
    cJSON_AddStringToObject(rsp_obj, "startTime", stime);
    cJSON_AddStringToObject(rsp_obj, "endTime", etime);
    gen_get_rsp(var, FUNC_HISTORYDATA, rsp_obj);
    cJSON_Delete(rsp_obj);
}

void ems_smu_allalarm(ems_smu_var_t *var, mqtt_message_t *mqtt_msg)
{
    cJSON *messages = cJSON_CreateArray();
    char cmd[CMD_MAX_LENGTH] = {0};
    snprintf(cmd, CMD_MAX_LENGTH, "hgetall "REDIS_REAL_ALARM_TABLE);
    redisReply *reply = (redisReply *)redisCommand(var->context, cmd);
    if (reply == NULL || reply->type != REDIS_REPLY_ARRAY || reply->elements < 0)
    {
        dbg_syslog(LOG_ERR, "get allalarm failed");
        return;
    }
    else
    {  
        for(int i = 0; i < reply->elements; i+=2)
        {
            cJSON *alarmmsg = cJSON_CreateObject();
            char *tok1 = strtok(reply->element[i]->str, "|");
            cJSON_AddStringToObject(alarmmsg, "alarmCode", tok1);
            tok1 = strtok(0, "|");
            cJSON_AddNumberToObject(alarmmsg, "severity", atoi(tok1));
            tok1 = strtok(0, "|");
            cJSON_AddStringToObject(alarmmsg, "entityInstance", tok1);

            char *tok2 = strtok(reply->element[i+1]->str, "|");
            cJSON_AddNumberToObject(alarmmsg, "reasonCode", atoi(tok2));
            tok2 = strtok(0, "|");
            cJSON_AddNumberToObject(alarmmsg, "raiseTime", atoi(tok2));
            tok2 = strtok(0, "|");
            if (!tok2)
                cJSON_AddStringToObject(alarmmsg, "additionalParam", "");
            else 
                cJSON_AddStringToObject(alarmmsg, "additionalParam", tok2);
                       
            cJSON_AddItemToArray(messages, alarmmsg); 
        }
    }
    freeReplyObject(reply);
    cJSON *rsp_obj = cJSON_CreateObject();
    cJSON_AddItemToObject(rsp_obj, "messages", messages);
    gen_get_rsp(var, FUNC_ALLALARM, rsp_obj);
    cJSON_Delete(rsp_obj);
}

void ems_smu_getsysparam(ems_smu_var_t *var, mqtt_message_t *mqtt_msg)
{
    cJSON *messages = cJSON_CreateArray();

    cJSON *lcid = cJSON_CreateObject();
    cJSON_AddStringToObject(lcid, "name", "id");
    cJSON_AddStringToObject(lcid, "value", var->lcid);
    cJSON_AddStringToObject(lcid, "desc", "站点编号");
    cJSON_AddItemToArray(messages, lcid);

    cJSON *lccap = cJSON_CreateObject();
    cJSON_AddStringToObject(lccap, "name", "capacity");
    cJSON_AddNumberToObject(lccap, "value", var->capacity);
    cJSON_AddStringToObject(lccap, "desc", "容量");
    cJSON_AddItemToArray(messages, lccap);

    cJSON *lcpower = cJSON_CreateObject();
    cJSON_AddStringToObject(lcpower, "name", "power");
    cJSON_AddNumberToObject(lcpower, "value", var->power);
    cJSON_AddStringToObject(lcpower, "desc", "功率");
    cJSON_AddItemToArray(messages, lcpower);

    cJSON *lcsta = cJSON_CreateObject();
    cJSON_AddStringToObject(lcsta, "name", "station");
    cJSON_AddStringToObject(lcsta, "value", var->lcsta);
    cJSON_AddStringToObject(lcsta, "desc", "坐标");
    cJSON_AddItemToArray(messages, lcsta);

    cJSON *rsp_obj = cJSON_CreateObject();
    cJSON_AddItemToObject(rsp_obj, "param", messages);
    gen_get_rsp(var, FUNC_GETSYSPARAM, rsp_obj);
    cJSON_Delete(rsp_obj);
}

void ems_smu_sysuse(ems_smu_var_t *var, mqtt_message_t *mqtt_msg)
{
    cJSON *messages = cJSON_CreateArray();

    cJSON *cpuuse = cJSON_CreateObject();
    cJSON_AddStringToObject(cpuuse, "name", "cpu");
    cJSON_AddStringToObject(cpuuse, "desc", "处理器");
    cJSON_AddStringToObject(cpuuse, "value", var->sysuse.cpuuse);
    cJSON_AddItemToArray(messages, cpuuse);

    cJSON *memuse = cJSON_CreateObject();
    cJSON_AddStringToObject(memuse, "name", "mem");
    cJSON_AddStringToObject(memuse, "desc", "内存");
    cJSON_AddStringToObject(memuse, "value", var->sysuse.memuse);
    cJSON_AddItemToArray(messages, memuse);

    cJSON *diskuse = cJSON_CreateObject();
    cJSON_AddStringToObject(diskuse, "name", "disk");
    cJSON_AddStringToObject(diskuse, "desc", "硬盘(已用/总kb)");
    cJSON_AddStringToObject(diskuse, "value", var->sysuse.diskuse);
    cJSON_AddItemToArray(messages, diskuse);

    cJSON *extdisk = cJSON_CreateObject();
    cJSON_AddStringToObject(extdisk, "name", "extdisk");
    cJSON_AddStringToObject(extdisk, "desc", "扩展磁盘(已用/总kb)");
    cJSON_AddStringToObject(extdisk, "value", var->sysuse.extdisk);
    cJSON_AddItemToArray(messages, extdisk);

    cJSON *inetconn = cJSON_CreateObject();
    cJSON_AddStringToObject(inetconn, "name", "inetconn");
    cJSON_AddStringToObject(inetconn, "desc", "公网链接");
    cJSON_AddNumberToObject(inetconn, "value", var->sysuse.inetcnt);
    cJSON_AddItemToArray(messages, inetconn);

    cJSON *rsp_obj = cJSON_CreateObject();
    cJSON_AddItemToObject(rsp_obj, "messages", messages);
    gen_get_rsp(var, FUNC_SYSUSE, rsp_obj);
    cJSON_Delete(rsp_obj);
}

void ems_smu_netstat(ems_smu_var_t *var, mqtt_message_t *mqtt_msg)
{
    cJSON *rsp_obj = cJSON_CreateObject();
    cJSON *messages = cJSON_CreateArray();
    cJSON_AddItemToObject(rsp_obj, "messages", messages);
    
    char data_dist[1024*16] = {0};
    sprintf(data_dist, "{");
    char *data_tmp = NULL;
    cmd_call(PATH_NET_STATUS, 3, &data_tmp, 1);

    if (data_tmp != NULL)
    {
        dbg_syslog(LOG_ERR, "cmd return %s", data_tmp);
        cJSON *netparse = cJSON_Parse(data_tmp);
        if (NULL == netparse)
        {
            dbg_syslog(LOG_ERR, "parse netstat failed");
            goto out;
        }
        char *devs[] = {"wan", "wwan", "lan"};
        char *items[] = {"proto", "l3_device", "txpackets", "rxpackets", "subnet", "status"};
        char *itemsout[] = {"proto", "intf", "txpacket", "rxpacket", "subnet", "connstat"};
        for (int i = 0; i < 3; i++)
        {
            cJSON *rawobj = cJSON_GetObjectItem(netparse, devs[i]);
            if (NULL == rawobj)
            {
                dbg_syslog(LOG_ERR, "no %s netdev", devs[i]);
            }
            cJSON *addObj = cJSON_CreateObject();
            for (int j = 0; j < 6; j++)
            {
                char val[64] = {0};
                GET_JSON_VALUE_STRING(rawobj, items[j], val);
                if (j < 5)
                    cJSON_AddStringToObject(addObj, itemsout[j], val);
                else 
                    cJSON_AddStringToObject(addObj, itemsout[j], atoi(val) > 0? "已连接":"未连接");
            }
            cJSON_AddItemToArray(messages, addObj);
        }
        cJSON_Delete(netparse);
        free(data_tmp);
    }

out:
    gen_get_rsp(var, FUNC_NETSTAT, rsp_obj);
    cJSON_Delete(rsp_obj);
}

void ems_smu_portcfg(ems_smu_var_t *var, mqtt_message_t *mqtt_msg)
{
    cJSON *port = NULL;

    char *pdata = NULL;
    pdata = read_file_data(PORT_CFG_PATH);
    if (pdata != NULL)
    {
        cJSON *root = cJSON_Parse(pdata);
        if (NULL != root)
        {
            cJSON *devcfg = cJSON_GetObjectItem(root, "dev_cfg");
            if (devcfg != NULL)
            {
                port = cJSON_Duplicate(devcfg, 1);
            }
            cJSON_Delete(root);
        }
    }

    if (port == NULL)
        port = cJSON_CreateObject();
    cJSON *rsp_obj = cJSON_CreateObject();
    cJSON_AddItemToObject(rsp_obj, "dev_cfg", port);
    gen_get_rsp(var, FUNC_GETPORTCFG, rsp_obj);
    cJSON_Delete(rsp_obj);
}

void ems_smu_about(ems_smu_var_t *var, mqtt_message_t *mqtt_msg)
{
    cJSON *messages = cJSON_CreateArray();

    cJSON *stano = cJSON_CreateObject();
    cJSON_AddStringToObject(stano, "name", "StationNo");
    cJSON_AddStringToObject(stano, "desc", "站点编号");
    cJSON_AddStringToObject(stano, "value", var->lcid);
    cJSON_AddItemToArray(messages, stano);

    cJSON *lcsn = cJSON_CreateObject();
    cJSON_AddStringToObject(lcsn, "name", "DevNo");
    cJSON_AddStringToObject(lcsn, "desc", "控制器编号");
    cJSON_AddStringToObject(lcsn, "value", var->sn_str);
    cJSON_AddItemToArray(messages, lcsn);

    cJSON *capa = cJSON_CreateObject();
    cJSON_AddStringToObject(capa, "name", "Capacity");
    cJSON_AddStringToObject(capa, "desc", "装机容量");
    cJSON_AddNumberToObject(capa, "value", var->capacity);
    cJSON_AddItemToArray(messages, capa);
    
    cJSON *emsver = cJSON_CreateObject();
    cJSON_AddStringToObject(emsver, "name", "EmsVer");
    cJSON_AddStringToObject(emsver, "desc", "EMS版本");
    cJSON_AddStringToObject(emsver, "value", var->softver);
    cJSON_AddItemToArray(messages, emsver);

    cJSON *rsp_obj = cJSON_CreateObject();
    cJSON_AddItemToObject(rsp_obj, "messages", messages);
    gen_get_rsp(var, FUNC_ABOUT, rsp_obj);
    cJSON_Delete(rsp_obj);
}

void ems_smu_setpoint(ems_smu_var_t *var, mqtt_message_t *mqtt_msg)
{
    char cmd[CMD_MAX_LENGTH] = {0};
    cJSON *root = cJSON_Parse(mqtt_msg->payload);
    if (root == NULL)
    {
        dbg_syslog(LOG_ERR, "parse mqtt msg failed");
        return;
    }

    cJSON *param = cJSON_GetObjectItem(root, "param");
    int paracnt = cJSON_GetArraySize(param);
    char setpointname[MAX_NAME_LEN] = {0};
    int setpointval = 0;
    for (int i = 0; i < paracnt; i++)
    {
        char devsn[SN_MAX_LEN] = {0};
        char devno[MAX_DEVTYPE_LEN] = {0};
        char entryname[MAX_PNAME_LEN] = {0};
        cJSON *para = cJSON_GetArrayItem(param, i);
        GET_JSON_VALUE_STRING(para, "name", setpointname);
        GET_JSON_VALUE_INT(para, "value", setpointval);
        char *tok = strtok(setpointname, ".");
        strncpy(devsn, tok, strlen(tok));
        snprintf(cmd, CMD_MAX_LENGTH, "hget devs %s", tok);
        redisReply *reply = (redisReply *)redisCommand(var->context, cmd);
        if (reply == NULL || REDIS_REPLY_STRING != reply->type)
        {
            dbg_syslog(LOG_ERR, "execcmd error");
            freeReplyObject(reply);
            continue;
        }
        tok = strtok(0, ".");
        strncpy(entryname, tok, strlen(tok));
        char *dtok = strtok(reply->str, "|");
        strncpy(devno, dtok, strlen(dtok));
        freeReplyObject(reply);

        dev_nodestruct_t *curdev;
        HASH_FIND_STR(var->devtypenode, devno, curdev);
        if (NULL == curdev)
        {
            dbg_syslog(LOG_ERR, "curdev not in devs");
            continue;
        }

        set_entry_t *curset;
        HASH_FIND_STR(curdev->setnode, entryname, curset);
        if (NULL == curset)
        {
            dbg_syslog(LOG_ERR, "curset not found in setnode");
            continue;
        }

        cJSON *setdev_obj = cJSON_CreateObject();
        cJSON_AddNumberToObject(setdev_obj, curset->settag, setpointval);
        gen_set_dev_service(var, setdev_obj, curset->sevname, devsn);
        cJSON_Delete(setdev_obj);
    }
    cJSON *rsp_obj = cJSON_CreateObject();
    cJSON *newparam = cJSON_Duplicate(param, 1); 
    cJSON_AddItemToObject(rsp_obj, "param", newparam);
    gen_set_rsp(var, rsp_obj, 1, FUNC_SETPOINT);
    cJSON_Delete(rsp_obj);
    cJSON_Delete(root);
}

void ems_smu_setsysparam(ems_smu_var_t *var, mqtt_message_t *mqtt_msg)
{
    cJSON *root = cJSON_Parse(mqtt_msg->payload);
    if (root == NULL)
    {
        dbg_syslog(LOG_ERR, "parse mqtt msg failed");
        return;
    }

    cJSON *param = cJSON_GetObjectItem(root, "param");
    int paracnt = cJSON_GetArraySize(param);
    dbg_syslog(LOG_ERR, "paracnt %d", paracnt);
    for(int i = 0; i < paracnt; i++)
    {
        cJSON *para = cJSON_GetArrayItem(param, i);
        if (para)
        {
            char paraname[MAX_ID_LEN] = {0};
            GET_JSON_VALUE_STRING(para, "name", paraname);
            dbg_syslog(LOG_ERR, "paraname %s", paraname);
            if (!strncmp(paraname, SYSPARA_ID, strlen(SYSPARA_ID)))
            {
                char paraid[MAX_ID_LEN] = {0};
                GET_JSON_VALUE_STRING(para, "value", paraid);
                memset(var->lcid, 0, MAX_ID_LEN);
                strncpy(var->lcid, paraid, strlen(paraid));
                cJSON *iditem = cJSON_GetObjectItem(var->staroot, SYSPARA_ID);
                memset(iditem->valuestring, 0, strlen(iditem->valuestring));
                strncpy(iditem->valuestring, paraid, strlen(paraid));
            }
            else if (!strncmp(paraname, SYSPARA_CAP, strlen(SYSPARA_CAP)))
            {
                int capacity = 0;
                GET_JSON_VALUE_INT(para, "value", capacity);
                var->capacity = capacity;
                cJSON *capitem = cJSON_GetObjectItem(var->staroot, SYSPARA_CAP);
                capitem->valueint = capacity;
            }
            else if (!strncmp(paraname, SYSPARA_POWER, strlen(SYSPARA_POWER)))
            {
                int power = 0; 
                GET_JSON_VALUE_INT(para, "value", power);
                var->power = power;
                dbg_syslog(LOG_ERR, "power %d var->power %d", power, var->power);
                cJSON *poweritem = cJSON_GetObjectItem(var->staroot, SYSPARA_POWER);
                poweritem->valueint = power;
            }
            else if (!strncmp(paraname, SYSPARA_STAN, strlen(SYSPARA_STAN)))
            {
                char parasta[MAX_STAN_LEN] = {0};
                GET_JSON_VALUE_STRING(para, "value", parasta);
                memset(var->lcsta, 0, MAX_STAN_LEN);
                strncpy(var->lcsta, parasta, strlen(parasta));
                cJSON *staitem = cJSON_GetObjectItem(var->staroot, SYSPARA_STAN);
                memset(staitem->valuestring, 0, strlen(staitem->valuestring));
                strncpy(staitem->valuestring, parasta, strlen(parasta));
            }
        }
    }
    if (paracnt > 0)
    {
        char *data_tmp = cJSON_PrintUnformatted(var->staroot);
        write_file_data(EMMSV2_STATION_CONFIG, data_tmp, strlen(data_tmp));
        free(data_tmp);
    }

    cJSON *rsp_obj = cJSON_CreateObject();
    cJSON_AddItemToObject(rsp_obj, "param", cJSON_Duplicate(param, 1));
    gen_set_rsp(var, rsp_obj, 1, FUNC_SETSYSPARAM);
    
    char topic[TOPIC_MAX_LEN] = {0};
    snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/lnxall/device/sysparam", var->sn_str);
    char *data_tmp2 = cJSON_PrintUnformatted(rsp_obj);
    ipc_session_publish(var->session, topic, data_tmp2, strlen(data_tmp2));
    free(data_tmp2);

    cJSON_Delete(rsp_obj);
    cJSON_Delete(root);
}

void ems_smu_setportcfg(ems_smu_var_t *var, mqtt_message_t *mqtt_msg)
{
    cJSON *rsp_obj = cJSON_CreateObject();

    cJSON *root = cJSON_Parse(mqtt_msg->payload);
    if (root == NULL)
    {
        dbg_syslog(LOG_ERR, "parse mqtt msg failed");
        goto fail_out;
    }
    cJSON *dstjson = cJSON_CreateObject();
    cJSON *port = cJSON_GetObjectItem(root, "port");
    cJSON_AddItemToObject(dstjson, "dev_cfg", cJSON_Duplicate(port, 1));
    cJSON_AddStringToObject(dstjson, "ver", "0.1");
    char *portstr = cJSON_Print(dstjson);
    if (portstr)
    {
        dbg_syslog(LOG_ERR, "portstr %s", portstr);
        write_file_data(PORT_CFG_PATH, portstr, strlen(portstr));
        free(portstr);
        portstr = NULL;
    }
    cJSON_Delete(dstjson);
/*
    char *pdata = NULL;
    pdata = read_file_data(PORT_CFG_PATH);
    cJSON *curportconf = cJSON_Parse(pdata);
    if (curportconf == NULL)
    {
        dbg_syslog(LOG_ERR, "parse port cfg failed");
        goto fail_out;
    }
    cJSON *curport = cJSON_GetObjectItem(curportconf, "dev_cfg");

    cJSON *port = cJSON_GetObjectItem(root, "port");
    char *porttype[] = {"uart", "lan", "wan", "can"};
    for (int i = 0; i < 4; i++)
    {
        dbg_syslog(LOG_ERR, "i %d porttype %s", i, porttype[i]);
        if (cJSON_HasObjectItem(port, porttype[i]))
        {
            dbg_syslog(LOG_ERR, "i %d porttype %s found", i, porttype[i]);
            cJSON *portt = cJSON_GetObjectItem(port, porttype[i]);
            int portcnt = cJSON_GetArraySize(portt);
            dbg_syslog(LOG_ERR, "portcnt %d", portcnt);
            for(int j = 0; j < portcnt; j++)
            {
                char curname[PORT_NAME_LEN] = {0};
                cJSON *singleport = cJSON_GetArrayItem(portt, j);
                GET_JSON_VALUE_STRING(singleport, "name", curname);
                dbg_syslog(LOG_ERR, "port name %s", curname);

                cJSON *portt_old = cJSON_GetObjectItem(curport, porttype[i]);
                cJSON *attr_old = cJSON_GetObjectItem(portt_old, "attr");
                int old_cnt = cJSON_GetArraySize(attr_old);
                dbg_syslog(LOG_ERR, "old_cnt %d ", old_cnt);
                for(int k = 0; k < old_cnt; k++)
                {
                    char oldname[PORT_NAME_LEN] = {0};
                    cJSON *oldsingleport = cJSON_GetArrayItem(attr_old, k);
                    GET_JSON_VALUE_STRING(oldsingleport, "name", oldname);
                    dbg_syslog(LOG_ERR, "oldname %s curname %s", oldname, curname);
                    if (strncmp(oldname, curname, strlen(oldname)) == 0)
                    {
                        cJSON_DeleteItemFromArray(attr_old, k);
                        cJSON_InsertItemInArray(attr_old, k, cJSON_Duplicate(singleport, 1));
                        break;
                    }
                }
            }
        }
    }
    
    char *portstr = cJSON_Print(curportconf);
    if (portstr)
    {
        dbg_syslog(LOG_ERR, "portstr %s", portstr);
        write_file_data(PORT_CFG_PATH, portstr, strlen(portstr));
        free(portstr);
        portstr = NULL;
    }
    cJSON_Delete(curportconf); 
*/
    gen_set_rsp(var, rsp_obj, 1, FUNC_SETPORTCFG); 
    cJSON_Delete(rsp_obj);
    cJSON_Delete(root);
    return;

fail_out:
    gen_set_rsp(var, rsp_obj, 0, FUNC_SETPORTCFG);
    cJSON_Delete(rsp_obj);
    cJSON_Delete(root);
    return;
}

void ems_smu_setlcupgrade(ems_smu_var_t *var, mqtt_message_t *mqtt_msg)
{
    cJSON *root = cJSON_Parse(mqtt_msg->payload);
    if (root == NULL)
    {
        dbg_syslog(LOG_ERR, "parse mqtt msg failed");
        return;
    }

    char url[URL_STR_LEN] = {0};
    char databuff[URL_STR_LEN] = {0};
    GET_JSON_VALUE_STRING(root, "url", url);
    cJSON *rsp_obj = cJSON_CreateObject();
    if (strlen(url) <= 0)
        gen_set_rsp(var, rsp_obj, 0, FUNC_LCUPGRADE);
    else
        gen_set_rsp(var, rsp_obj, 1, FUNC_LCUPGRADE);

    cmd_call("rm -f "EMS_UPGRADE_FILE, 10, NULL, 0);
    sprintf(databuff, "cp %s "EMS_UPGRADE_FILE, url);
    if (cmd_call(databuff, 10, NULL, 1) != 0)
        dbg_syslog(LOG_ERR, "upgrade failed");
    else 
    {
        dbg_syslog(LOG_ERR, "upgrade success");
        cmd_call("systemctl restart subemsd.service", 1, NULL, 0);
    }

    cJSON_Delete(rsp_obj);
    cJSON_Delete(root);
}

void ems_smu_setntp(ems_smu_var_t *var, mqtt_message_t *mqtt_msg)
{
    cJSON *root = cJSON_Parse(mqtt_msg->payload);
    if (root == NULL)
    {
        dbg_syslog(LOG_ERR, "parse mqtt msg failed");
        return;
    }

    char smutime[LONG_STR_LEN] = {0};
    GET_JSON_VALUE_STRING(root, "smutime", smutime);
    cJSON *rsp_obj = cJSON_CreateObject();
    if (strlen(smutime) <= 0)
        gen_set_rsp(var, rsp_obj, 0, FUNC_NTP);
    else
        gen_set_rsp(var, rsp_obj, 1, FUNC_NTP);

    cJSON_Delete(rsp_obj);
    cJSON_Delete(root);
}

void ems_smu_setreboot(ems_smu_var_t *var, mqtt_message_t *mqtt_msg)
{
    cJSON *root = cJSON_Parse(mqtt_msg->payload);
    if (root == NULL)
    {
        dbg_syslog(LOG_ERR, "parse mqtt msg failed");
        return;
    }

    cJSON *rsp_obj = cJSON_CreateObject();
    gen_set_rsp(var, rsp_obj, 1, FUNC_REBOOT);
    cmd_call("reboot -d 20 &", 1, NULL, 0);
    cJSON_Delete(rsp_obj);
    cJSON_Delete(root);
}

/*
* split topic into action, func_id
* retval: Compression algorithm
*          -1 - invalid topic; 0 - No Compression; 1 - lz4; ...
*/
static int split_topic(char *buf, char *sep, char *action, char *func)
{
    char *tok;
    int ret = -1;

    dbg_syslog(LOG_INFO, "split topic %s", buf);
    tok = strtok(buf, sep);
    if (!tok || strncmp(tok, PLATFORM, strlen(PLATFORM)))
    {
        dbg_syslog(LOG_ERR, "invalid platform %s", tok);
        return ret;
    }

    tok = strtok(0, sep);
    if (!tok)
    {
        dbg_syslog(LOG_ERR, "invalid action %s", tok);
        return ret;
    }
    snprintf(action, MAX_ACT_LEN, "%s", tok);

    tok = strtok(0, sep);
    if (!tok)
    {
        dbg_syslog(LOG_ERR, "invalid SN %s", tok);
        return ret;
    }

    tok = strtok(0, sep);
    if (!tok)
    {
        dbg_syslog(LOG_ERR, "invalid function id %s", tok);
        return ret;
    }
    snprintf(func, MAX_FUNC_LEN, "%s", tok);

    tok = strtok(0, sep);
    if (!tok)
    {
        dbg_syslog(LOG_DEBUG, "No compress");
        return 0;
    }

    if (!strcmp(tok, CMPLZ4))
    {
        ret = 1;
    }

    return ret;
}

int recv_host_msg(void *obj, mqtt_message_t *mqtt_msg)
{
    ems_smu_var_t *var = (ems_smu_var_t *)obj;
    char *buf = NULL;
    char action[MAX_ACT_LEN] = {0};
    char func[MAX_FUNC_LEN] = {0};

    dbg_syslog(LOG_DEBUG, "External MQTT client received MQTT topic:%s payload:%s",
              mqtt_msg->topic, mqtt_msg->payload);

    buf = strdup(mqtt_msg->topic);
    int compress = split_topic(buf, "/", action, func); 
    if (compress < 0)
        return -1;
    else if (compress == 1)
    {
        //uncompress
    }

    if (!strcmp(action, LCPOSTRSP))
    {
        if (!strcmp(func, FUNC_LOGIN))
            ems_smu_loginrsp(var, mqtt_msg);
    }
    else if(!strcmp(action, GET))
    {
        if (!strcmp(func, FUNC_NODESTRUCTURE))
            ems_smu_nodestructure(var, mqtt_msg);
        else if (!strcmp(func, FUNC_HISTORYDATA))
            ems_smu_historydata(var, mqtt_msg);
        else if (!strcmp(func, FUNC_ALLALARM))
            ems_smu_allalarm(var, mqtt_msg);
        else if (!strcmp(func, FUNC_GETSYSPARAM))
            ems_smu_getsysparam(var, mqtt_msg);
        else if (!strcmp(func, FUNC_SYSUSE))
            ems_smu_sysuse(var, mqtt_msg);
        else if (!strcmp(func, FUNC_NETSTAT))
            ems_smu_netstat(var, mqtt_msg);
        else if (!strcmp(func, FUNC_GETPORTCFG))
            ems_smu_portcfg(var, mqtt_msg);
        else if (!strcmp(func, FUNC_ABOUT))
            ems_smu_about(var, mqtt_msg);
    } 
    else if (!strcmp(action, SET))
    {
        if (!strcmp(func, FUNC_SETPOINT))
            ems_smu_setpoint(var, mqtt_msg);
        else if (!strcmp(func, FUNC_SETSYSPARAM))
            ems_smu_setsysparam(var, mqtt_msg);
        else if (!strcmp(func, FUNC_SETPORTCFG))
            ems_smu_setportcfg(var, mqtt_msg);
        else if (!strcmp(func, FUNC_LCUPGRADE))
            ems_smu_setlcupgrade(var, mqtt_msg);
        else if (!strcmp(func, FUNC_NTP))
            ems_smu_setntp(var, mqtt_msg);
        else if (!strcmp(func, FUNC_REBOOT))
            ems_smu_setreboot(var, mqtt_msg);
    }

    free(buf);
    return 0;
}
