#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 "cloud_emmsv2_common.h"
#include "dy_utils/cJSON.h"
#include "dy_utils/dy_common.h"
#include "mqtt_session.h"
#include "dy_utils/protocol.h"

void login_rsp_msg(cloud_emmsv2_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;
    }

    int loginresult = -1;
    GET_JSON_VALUE_INT(root, "result", loginresult);
    if (0 == loginresult)
    {
        dbg_syslog(LOG_NOTICE, "login success");
        var->comm_status = COMM_STATUS_LOGIN;
        cloud_emmsv2_pubdevinfo(var);
    }
    dbg_syslog(LOG_DEBUG, "login finished");

    cJSON_Delete(root);
}

static void gen_set_rsp(cloud_emmsv2_var_t *var, int result, char *funcid, int seq)
{
    char topic[TOPIC_MAX_LEN];
    cJSON *rsp_obj = cJSON_CreateObject();
    if (!rsp_obj)
    {
        dbg_syslog(LOG_ERR, "Create rsp object failed");
        return;
    }

    cJSON_AddStringToObject(rsp_obj, "funcId", funcid);
    cJSON_AddStringToObject(rsp_obj, "lcSN", var->sn_str);
    cJSON_AddNumberToObject(rsp_obj, "seq", 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, EMMSSETRSP, var->sn_str, funcid);
    
    mqtt_session_publish(var->session_cloud, topic, data_tmp, strlen(data_tmp));
    free(data_tmp);
    cJSON_Delete(rsp_obj);
}

static void gen_get_rsp(cloud_emmsv2_var_t *var, char *funcid, char *dataname, cJSON *data, int seq)
{
    char topic[TOPIC_MAX_LEN] = {0};
    cJSON *rsp_obj = cJSON_CreateObject();
    if (!rsp_obj)
    {
        dbg_syslog(LOG_ERR, "Create rsp object failed");
        return;
    }

    cJSON_AddStringToObject(rsp_obj, "funcId", funcid);
    cJSON_AddStringToObject(rsp_obj, "lcSN", var->sn_str);
    cJSON_AddNumberToObject(rsp_obj, "seq", seq);
    cJSON_AddNumberToObject(rsp_obj, "time", time(NULL));
    cJSON_AddItemToObject(rsp_obj, dataname, data);

    char *data_tmp = cJSON_PrintUnformatted(rsp_obj);
    snprintf(topic, TOPIC_MAX_LEN, "%s/%s/%s/%s", PLATFORM, EMMSGETRSP, var->sn_str, funcid);
    
    mqtt_session_publish(var->session_cloud, topic, data_tmp, strlen(data_tmp));
    free(data_tmp);
    cJSON_Delete(rsp_obj);
}

void gen_set_dev_service(cloud_emmsv2_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);
    cJSON_Delete(rt_data);
}

void gen_set_rglt_service(cloud_emmsv2_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 emsctrl_telm_msg(cloud_emmsv2_var_t *var, mqtt_message_t *mqtt_msg)
{
    int seq = 0;
    cJSON *root = cJSON_Parse(mqtt_msg->payload);
    if (root == NULL)
    {
        dbg_syslog(LOG_ERR, "parse mqtt msg failed");
        return;
    }

    GET_JSON_VALUE_INT(root, "seq", seq);
    cJSON *tags  = cJSON_CreateObject();
    /* collect emsctrl telemetry tags */

    gen_get_rsp(var, FUNC_TELM, "tags", tags, seq);
    cJSON_Delete(root);
    return;
}

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

    GET_JSON_VALUE_INT(root, "seq", seq);
    cJSON *tags = cJSON_CreateObject();
    /* collect emsctrl telecontrol tags */

    gen_get_rsp(var, FUNC_TELC, "tags", tags, seq);
    cJSON_Delete(root);
    return;
}

void subdev_telm_msg(cloud_emmsv2_var_t *var, mqtt_message_t *mqtt_msg)
{
    int seq = 0;
    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;
    }

    GET_JSON_VALUE_INT(root, "seq", seq);
    cJSON *down_messages = cJSON_GetObjectItem(root, "messages");
    if (!down_messages || !cJSON_IsArray(down_messages))
    {
        dbg_syslog(LOG_ERR, "down messages is not an array");
        goto fail_out;
    }

    cJSON *upmessage = cJSON_CreateArray();
    int node_cnt = cJSON_GetArraySize(down_messages);
    for(int i = 0; i < node_cnt; i++)
    {
        cJSON *subdev = cJSON_GetArrayItem(down_messages, i);
        char devno[32] = {0};
        GET_JSON_VALUE_STRING(subdev, "no", devno);

        /* obtain all telemetry data of devno */
        snprintf(cmd, CMD_MAX_LENGTH, "zrange %s 0 -1", devno);
        redisReply *reply_data = (redisReply *)redisCommand(var->context, cmd);
        if (NULL == reply_data || (reply_data->type != REDIS_REPLY_ARRAY) || (reply_data->elements <= 0))
        {
            dbg_syslog(LOG_ERR, "get dev %s data failed", devno);
            continue;
        }
        cJSON *curdevdata = cJSON_CreateObject();
        cJSON_AddItemToArray(upmessage, curdevdata);
        cJSON_AddStringToObject(curdevdata, "no", devno);
        cJSON *tags = cJSON_CreateObject();
        cJSON_AddItemToObject(curdevdata, "tags", tags);
        for (int j = 0; j < reply_data->elements; j++)
        {
            char *tok = strtok(reply_data->element[j]->str, "|");
            char *tok2 = strtok(0, "|");
            if (strstr(tok2, "."))
            {
                cJSON_AddNumberToObject(tags, tok, atof(tok2));
            }
            else
            {
                cJSON_AddNumberToObject(tags, tok, atoi(tok2));
            }
        }
        freeReplyObject(reply_data);
    }

    gen_get_rsp(var, FUNC_SUBTELM, "messages", upmessage, seq);
    cJSON_Delete(root);
    return;

fail_out:
    cJSON_Delete(root);
    return;
}

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

    GET_JSON_VALUE_INT(root, "seq", seq);
    cJSON *messages = cJSON_CreateObject();
    /* collect alarm and gen message */


    gen_get_rsp(var, FUNC_ALLALARM, "messages", messages, seq);
    cJSON_Delete(root);
    return;
}

void emsctrl_sched_msg(cloud_emmsv2_var_t *var, mqtt_message_t *mqtt_msg)
{
    int seq = 0;
    char date[MAX_DATE_LEN] = {0};
    char starttime[TIME_STR_LEN] = {0};
    char endtime[TIME_STR_LEN] = {0};
    char cmd[CMD_MAX_LENGTH + MAX_STRA_LEN] = {0};
    char stra[MAX_STRA_LEN] = {0};
    char tempstr[20] = {0};
    
    int power = 0;
    int wggl = 0;
    cJSON *root = cJSON_Parse(mqtt_msg->payload);
    if (root == NULL)
    {
        dbg_syslog(LOG_ERR, "parse mqtt msg failed");
        return;
    }

    GET_JSON_VALUE_INT(root, "seq", seq);
    cJSON *messages = cJSON_GetObjectItem(root, "messages");
    int node_cnt = cJSON_GetArraySize(messages);
    for (int  i = 0; i < node_cnt; i++)
    {
        unsigned short segpower[MAX_TIME_SEG] = {0}; //记录每一段的有功功率
        unsigned char segmode[MAX_TIME_SEG] = {0};   //记录每一段的充放电模式
        short segwggl[MAX_TIME_SEG] = {0};           //记录每一段的无功功率
        cJSON *sche_day = cJSON_GetArrayItem(messages, i);
        GET_JSON_VALUE_STRING(sche_day, "Date", date); //yyyy-mm-dd，形如2023-02-01
        cJSON *schedule = cJSON_GetObjectItem(sche_day, "Schedule");
        int sche_cnt = cJSON_GetArraySize(schedule);
        for (int j = 0; j < sche_cnt; j++)
        {
            cJSON *sche_segm = cJSON_GetArrayItem(schedule, j);
            GET_JSON_VALUE_STRING(sche_segm, "StartTime", starttime); //hh:mm，形如01:00
            GET_JSON_VALUE_STRING(sche_segm, "EndTime", endtime);   //hh:mm，形如01:30
            GET_JSON_VALUE_INT(sche_segm, "Power", power);
            GET_JSON_VALUE_INT(sche_segm, "Wggl", wggl);
            unsigned char curmode = (power == 0)? 0 : ((power < 0)? 1 : 2); //0-静置 1-充电 2-放电
            unsigned short curpower = (power < 0)? (0-power) : power;
            int startindex = (((starttime[0]-'0')<<5)+((starttime[0]-'0')<<3))+((starttime[1]-'0')<<2)+(10 *(starttime[3]-'0') + starttime[4]-'0')/15; //(时*60+分)/15=时*4+分/15=时的十位数*40+时的个位数*4+分/15，x*40=x*32+x*8=(x<<5)+(x<<3)
            int endindex = (((endtime[0]-'0')<<5)+((endtime[0]-'0')<<3))+((endtime[1]-'0')<<2)+(10 *(endtime[3]-'0') + endtime[4]-'0')/15;
            for (int k = startindex; k < endindex; k++)
            {
                segpower[k] = curpower;
                segmode[k] = curmode;
                segwggl[k] = wggl;
            }
        }
        
        uint tscore = atoi(date)*10000 + atoi(date+5)*100+atoi(date+8); //年*10000+月*100+日，如2023-02-01转换成20230201作为存在数据库realstra表中策略的分数
        snprintf(cmd, CMD_SHORT_LENGTH, "zremrangebyscore "REDIS_STRA_BASE" %u %u", tscore, tscore); //先删除该日期下原有的策略
        redisReply *reply = (redisReply *)redisCommand(var->context, cmd);
        if (NULL == reply)
        {
            dbg_syslog(LOG_ERR, "remove tscore %u strategy failed", tscore);
            continue;
        }
        freeReplyObject(reply);

        snprintf(tempstr, 20, "%u|", tscore);
        strncat(stra,tempstr,strlen(tempstr));
        for(int n = 0; n < MAX_TIME_SEG; n++)
        {
            snprintf(tempstr, 20, "%u|%u|%d|", segpower[n], segmode[n], segwggl[n]);
            strncat(stra,tempstr,strlen(tempstr));
        }
        snprintf(cmd, CMD_MAX_LENGTH + MAX_STRA_LEN, "zadd "REDIS_STRA_BASE" %u %s", tscore, stra);
        reply = (redisReply*)redisCommand(var->context,cmd);
        if (NULL == reply)
        {
            dbg_syslog(LOG_ERR, "write schedule execcmd fail");
            continue;
        } 
        freeReplyObject(reply);
        memset(stra,0,MAX_STRA_LEN);
    }
    
    gen_set_rsp(var, 0, FUNC_SCHED, seq);
    cJSON_Delete(root);

    if (var->controlMode == 1)  //当前在自动模式下则立刻下发一次策略
    {
        time_t now;
        time(&now);
        struct tm *t;
        t = localtime(&now);
        uint tmscore = (t->tm_year + 1900) * 10000 + (t->tm_mon + 1) * 100 + t->tm_mday;
        cloud_emmsv2_pubstra(var, tmscore);
    }

    return;

//fail_out:
//    gen_set_rsp(var, 1, FUNC_SCHED);
//    cJSON_Delete(root);
//    return;
}

void emsctrl_telc_set_msg(cloud_emmsv2_var_t *var, mqtt_message_t *mqtt_msg)
{
    int seq = 0;
    char cmd_setparam[CMD_MAX_LENGTH] = {0};
    char cmd_proparam[CMD_MAX_LENGTH] = {0};
    char db_data[12] = {0};
    bool powerflag = false;

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

    GET_JSON_VALUE_INT(root, "seq", seq);
    strcat(cmd_setparam, "hset "REDIS_SETPARAM_BASE" ");
    strcat(cmd_proparam, "hset "REDIS_PROPARAM_BASE" ");
    cJSON *tags = cJSON_GetObjectItem(root, "tags");
    cJSON *curtag = cJSON_GetArrayItem(tags, 0);
    do
    {
        if (strstr(curtag->string, SET_ONOFF))  //设置开关机
        {
            cJSON *rt_data = cJSON_CreateObject();
            cJSON_AddNumberToObject(rt_data, "runState", curtag->valueint);
            gen_set_rglt_service(var, rt_data, "openpcs_set");

            strcat(cmd_setparam, curtag->string);
            strcat(cmd_setparam, " ");
            snprintf(db_data, 12, "%d ", curtag->valueint);
            strcat(cmd_setparam, db_data);
        }
        else if (strstr(curtag->string, SET_POWER) || strstr(curtag->string, SET_REPOWER))  //设置有功和无功功率
        {
            powerflag = true;
            strcat(cmd_setparam, curtag->string);
            strcat(cmd_setparam, " ");
            snprintf(db_data, 12, "%d ", curtag->valueint);
            strcat(cmd_setparam, db_data);
        }
        else if (0 == strncmp(curtag->string, SET_CONNECT, strlen(SET_CONNECT)))    //PCS并离网
        {
            cJSON *rt_data = cJSON_CreateObject();
            cJSON_AddNumberToObject(rt_data, "connState", curtag->valueint);
            gen_set_rglt_service(var, rt_data, "pcsconnect_set");

            strcat(cmd_setparam, curtag->string);
            strcat(cmd_setparam, " ");
            snprintf(db_data, 12, "%d ", curtag->valueint);
            strcat(cmd_setparam, db_data);
        }
        else if (strstr(curtag->string, SET_CLR))  //清除保护
        {            
            strcat(cmd_setparam, curtag->string);
            strcat(cmd_setparam, " ");
            snprintf(db_data, 12, "%d ", curtag->valueint);
            strcat(cmd_setparam, db_data);
        }
        else if (strstr(curtag->string, SET_CTRLMODE))
        {
            if ((curtag->valueint == 1) && (var->controlMode != 1)) //削峰填谷
            {
                time_t now;
                time(&now);
                struct tm *t;
                t = localtime(&now);
                uint tmscore = (t->tm_year + 1900) * 10000 + (t->tm_mon + 1) * 100 + t->tm_mday;
                cloud_emmsv2_pubstra(var,tmscore);
                var->lastmday = t->tm_mday;
            }
            //var->controlMode = curtag->valueint;
            cJSON *rt_data = cJSON_CreateObject();
            cJSON_AddNumberToObject(rt_data, "ControlMode", curtag->valueint);
            gen_set_rglt_service(var, rt_data, "emms_modeset");

            strcat(cmd_setparam, curtag->string);
            strcat(cmd_setparam, " ");
            snprintf(db_data, 12, "%d ", curtag->valueint);
            strcat(cmd_setparam, db_data);
        }
        else if (strstr(curtag->string, SET_BMSCONN))   //BMS上下高压
        {
            cJSON *rt_data = cJSON_CreateObject();
            char devsn[32] = {0};
            GET_JSON_VALUE_STRING(root, "deviceSN", devsn);
            cJSON_AddNumberToObject(rt_data, "setconnect", curtag->valueint);
            cJSON_AddStringToObject(rt_data, "deviceSN", devsn);
            gen_set_rglt_service(var, rt_data, "BmsSetConnect");

            strcat(cmd_setparam, curtag->string);
            strcat(cmd_setparam, " ");
            snprintf(db_data, 12, "%d ", curtag->valueint);
            strcat(cmd_setparam, db_data);
        }
        else   //保护参数
        {
            strcat(cmd_proparam, curtag->string);
            strcat(cmd_proparam, " ");
            snprintf(db_data, 12, "%.4lf ", curtag->valuedouble);
            strcat(cmd_proparam, db_data);
        }
        curtag = curtag->next;
    } while (curtag!=NULL);

    if (true == powerflag)
    {
        double yggl,wggl;
        GET_JSON_VALUE_DOUBLE(tags, SET_POWER, yggl);
        GET_JSON_VALUE_DOUBLE(tags, SET_REPOWER, wggl);
        cJSON *rt_data = cJSON_CreateObject();
        cJSON_AddNumberToObject(rt_data, "pcs_yggl", yggl);
        cJSON_AddNumberToObject(rt_data, "pcs_wggl", wggl);
        gen_set_rglt_service(var, rt_data, "pcsgl_set");
    }

    if (strlen(cmd_setparam) > strlen("hset "REDIS_SETPARAM_BASE" "))
    {
        redisReply *reply_set = redisCommand(var->context, cmd_setparam);
        if (NULL == reply_set)
        {
            dbg_syslog(LOG_ERR, "set telecontrol param failed");
        }
        freeReplyObject(reply_set);
    }

    if (strlen(cmd_proparam) > strlen("hset "REDIS_PROPARAM_BASE" "))
    {
        redisReply *reply_pro = redisCommand(var->context, cmd_proparam);
        if (NULL == reply_pro)
        {
            dbg_syslog(LOG_ERR, "set telecontrol param failed");
        }
        freeReplyObject(reply_pro);
    }

    /* parse tags and set ems_ctrl */
    if (root)
    {
        goto succ_out;
    }
    else
    {
        goto fail_out;
    }

succ_out:
    gen_set_rsp(var, 0, FUNC_TELC, seq);
    cJSON_Delete(root);
    return;

fail_out:
    gen_set_rsp(var, 1, FUNC_TELC, seq);
    cJSON_Delete(root);
    return;
}

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

    GET_JSON_VALUE_INT(root, "seq", seq);
    cJSON *messages = cJSON_GetObjectItem(root, "messages");
    int node_cnt = cJSON_GetArraySize(messages);
    for (int i = 0; i < node_cnt; i++)
    {
        cJSON *msg = cJSON_GetArrayItem(messages, i);
        char subdevid[32] = {0};
        GET_JSON_VALUE_STRING(msg,"no", subdevid);
        cJSON *tags = cJSON_GetObjectItem(msg, "tags");
        cJSON *curtag = cJSON_GetArrayItem(tags, 0);
        do
        {
            if (strstr(curtag->string, SET_BMSCONN))
            {
                char clustersn[32] = {0};
                char clustersn2[32] = {0};
                char *rack = strstr(subdevid, "Rack");
                if (rack != NULL)
                {
                    strncpy(clustersn, subdevid, strlen(subdevid)-strlen(rack));
                    strncpy(clustersn2, subdevid, strlen(subdevid)-strlen(rack));
                    strcat(clustersn, "Cluster");
                    strcat(clustersn2, "Cluster_");
                    strcat(clustersn, rack+4);
                    strcat(clustersn2, rack+4);
                    int clusteridx = strlen(clustersn) - 1;
                    while(clustersn[clusteridx] != '_')
                    {
                        clustersn[clusteridx--]='\0';
                    }
                    clustersn[clusteridx] = '\0';
                    int clusteridx2 = strlen(clustersn2) - 1;
                    while(clustersn2[clusteridx2] != '_')
                    {
                        clustersn2[clusteridx2--]='\0';
                    }
                    clustersn2[clusteridx2] = '\0';
                    cJSON *rt_data = cJSON_CreateObject();
                    if (curtag->valueint > 0)
                        cJSON_AddNumberToObject(rt_data, "TAG_7530", curtag->valueint);
                    else 
                        cJSON_AddNumberToObject(rt_data, "TAG_7530", 2);

                    cJSON *rt_data2 = cJSON_CreateObject();
                    if (curtag->valueint > 0)
                        cJSON_AddNumberToObject(rt_data2, "TAG_7530", curtag->valueint);
                    else 
                        cJSON_AddNumberToObject(rt_data2, "TAG_7530", 2);
                    gen_set_dev_service(var, rt_data, "Start_Connection_set", clustersn);
                    gen_set_dev_service(var, rt_data2, "Start_Connection_set", clustersn2);
                }
            }
            else if (strstr(curtag->string, SET_SOCSET))
            {
                cJSON *rt_data = cJSON_CreateObject();
                cJSON_AddNumberToObject(rt_data, "soc", curtag->valueint);
                gen_set_rglt_service(var, rt_data, "charge_soc_set");
            }
            else if (strstr(curtag->string, SET_DO5))
            {
                cJSON *rt_data = cJSON_CreateObject();
                cJSON_AddNumberToObject(rt_data, SET_DO5, curtag->valueint);
                gen_set_dev_service(var, rt_data, SET_DO5, subdevid);
            }
            else if (strstr(curtag->string, SET_DO6))
            {
                cJSON *rt_data = cJSON_CreateObject();
                cJSON_AddNumberToObject(rt_data, SET_DO6, curtag->valueint);
                gen_set_dev_service(var, rt_data, SET_DO6, subdevid);
            }
            else if (strstr(curtag->string, SET_PCSRST))
            {
                cJSON *rt_data = cJSON_CreateObject();
                cJSON_AddNumberToObject(rt_data, "TAG_4EB0", 1);
                gen_set_dev_service(var, rt_data, "pcs_reset", subdevid);
            }
            curtag = curtag->next;
        } while (curtag != NULL);
    }
    
    /* sub tel control */
    if (root)
    {
        goto succ_out;
    }
    else
    {
        goto fail_out;
    }

succ_out:
    gen_set_rsp(var, 0, FUNC_SUBTELC, seq);
    cJSON_Delete(root);
    return;

fail_out:
    gen_set_rsp(var, 1, FUNC_SUBTELC, seq);
    cJSON_Delete(root);
    return;
}

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

    GET_JSON_VALUE_INT(root, "seq", seq);
    cJSON *messages = cJSON_GetObjectItem(root, "messages");
    cJSON *alarmcodes = cJSON_GetObjectItem(messages, "alarmCodes");
    if (!alarmcodes || !cJSON_IsArray(alarmcodes))
    {
        dbg_syslog(LOG_ERR, "no alarmCodes found in AlarmDisable");
        goto fail_out;
    }

    /* disable given alarm */
    if (root)
    {
        goto succ_out;
    }
    
succ_out:
    gen_set_rsp(var, 0, FUNC_ALARMDIS, seq);
    cJSON_Delete(root);
    return;

fail_out:
    gen_set_rsp(var, 1, FUNC_ALARMDIS, seq);
    cJSON_Delete(root);
    return;
}

void dashboard_msg(cloud_emmsv2_var_t *var, mqtt_message_t *mqtt_msg)
{
    int seq = 0; 
    int timeout = DEF_DB_TIMEOUT;
    int period = DEF_DB_PERIOD;
    int action = -1;
    cJSON *root = cJSON_Parse(mqtt_msg->payload);
    if (root == NULL)
    {
        dbg_syslog(LOG_ERR, "parse mqtt msg failed");
        return;
    }

    GET_JSON_VALUE_INT(root, "seq", seq);
    cJSON *info = cJSON_GetObjectItem(root, "info");
    if (cJSON_HasObjectItem(info, "timeout"))
    {
        GET_JSON_VALUE_INT(info, "timeout", timeout);
    }
    if (cJSON_HasObjectItem(info, "period"))
    {
        GET_JSON_VALUE_INT(info, "period", period);
    }
    if (!cJSON_HasObjectItem(info, "action"))
    {
        dbg_syslog(LOG_ERR, "no action segment in dashboardctrl");
        goto fail_out;
    }
    GET_JSON_VALUE_INT(info, "action", action);
    var->db_endtime = time(NULL);
    if (action == 1)
    {
        var->db_period = period;
        var->db_endtime += timeout;
    }

    gen_set_rsp(var, 0, FUNC_DASHBOARDCTRL, seq);
    cJSON_Delete(root);
    return;

fail_out:
    gen_set_rsp(var, 1, FUNC_DASHBOARDCTRL, seq);
    cJSON_Delete(root);
    return;
}

/*
* 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)
{
    cloud_emmsv2_var_t *var = (cloud_emmsv2_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))
            login_rsp_msg(var, mqtt_msg);
    }
    else if(!strcmp(action, EMMSGET))
    {
        if (!strcmp(func, FUNC_TELM))
            emsctrl_telm_msg(var, mqtt_msg);
        else if (!strcmp(func, FUNC_TELC))
            emsctrl_telc_msg(var, mqtt_msg);
        else if (!strcmp(func, FUNC_SUBTELM))
            subdev_telm_msg(var, mqtt_msg);
        else if (!strcmp(func, FUNC_ALLALARM))
            allalarm_msg(var, mqtt_msg);
    } 
    else if (!strcmp(action, EMMSSET))
    {
        if (!strcmp(func, FUNC_SCHED))
            emsctrl_sched_msg(var, mqtt_msg);
        else if (!strcmp(func, FUNC_TELC))
            emsctrl_telc_set_msg(var, mqtt_msg);
        else if (!strcmp(func, FUNC_SUBTELC))
            subdev_telc_msg(var, mqtt_msg);
        else if (!strcmp(func, FUNC_ALARMDIS))
            alarmdis_msg(var, mqtt_msg);
        else if (!strcmp(func, FUNC_DASHBOARDCTRL))
            dashboard_msg(var, mqtt_msg);
    }

    free(buf);
    return 0;
}
