#include <stdio.h>
#include <stdlib.h>
#include <sys/stat.h>
#include <sys/sysinfo.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"

static void parse_stra(char *strastr, unsigned short *segpower, unsigned char *segmode, short *segwggl)
{
    char *tok;
    tok = strtok(strastr, "|");
    if (!tok)
    {
        dbg_syslog(LOG_ERR, "strastr format error");
    }

    for(int i = 0; i < MAX_TIME_SEG; i++)
    {
        tok = strtok(0, "|");
        segpower[i] = atoi(tok);
        tok = strtok(0, "|");
        segmode[i] = atoi(tok);
        tok = strtok(0, "|");
        segwggl[i] = atoi(tok);
    }
}

static void timer_dev_report(ems_smu_var_t *var)
{
    char topic[TOPIC_MAX_LEN] = {0};
    cJSON *root = cJSON_CreateObject();
    if (!root)
    {
        dbg_syslog(LOG_ERR, "create json object failed");
        return;
    }
    cJSON_AddStringToObject(root, "funcId", FUNC_DATAREPORT);
    cJSON_AddStringToObject(root, "lcSN", var->sn_str);
    cJSON_AddNumberToObject(root, "seq", var->seq++);
    cJSON_AddNumberToObject(root, "time", time(NULL));
    cJSON *data = cJSON_CreateArray();
    cJSON_AddItemToObject(root, "data", data);

    char cmd[CMD_SHORT_LENGTH] = {0};
    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");
        freeReplyObject(reply);
        return;
    } 
    else 
    {
        for(int i = 0; i <= reply->elements; i+=2)
        {
            dev_nodestruct_t *curdevtype;
            char devsn[SN_MAX_LEN] = {0};
            char dtype[MAX_DEVTYPE_LEN] = {0};

            if (i < reply->elements)
            {
                snprintf(devsn, SN_MAX_LEN, "%s", reply->element[i]->str);
                char *tok = strtok(reply->element[i+1]->str, "|");
                snprintf(dtype, MAX_DEVTYPE_LEN, "%s", tok);
            }
            else
            {
                strcpy(devsn, "EMS");
                strcpy(dtype, "EMS");
            }
            HASH_FIND_STR(var->devtypenode, dtype, curdevtype);
            if (curdevtype == NULL)
            {
                dbg_syslog(LOG_ERR, "devtype %s %s not found ", devsn, dtype);
                continue;
            }

            snprintf(cmd, CMD_SHORT_LENGTH, "zrange %s 0 -1", devsn);
            redisReply *dataReply = (redisReply *)redisCommand(var->context, cmd);
            if (NULL == reply || dataReply->type != REDIS_REPLY_ARRAY)
            {
                dbg_syslog(LOG_ERR, "get dev data failed");
                freeReplyObject(dataReply);
                continue;
            }
            for (int j = 0; j < dataReply->elements; j++)
            {
                cJSON *dataItem = cJSON_CreateObject();
                char *dtok = strtok(dataReply->element[j]->str, "|");
                char dataname[MAX_NAME_LEN] = {0};
                snprintf(dataname, MAX_NAME_LEN, "%s.%s", devsn, dtok);
                cJSON_AddStringToObject(dataItem, "name", dataname);
                
                data_entry_t *curdata;
                HASH_FIND_STR(curdevtype->datanode, dtok, curdata);
                if (curdata == NULL)
                {
                    dbg_syslog(LOG_ERR, "data node %s not found", dtok);
                    continue;
                }
                dtok = strtok(0, "|");
                if (curdata->type == 0)
                {
                    cJSON_AddNumberToObject(dataItem, "value", atoi(dtok));
                }
                else
                {
                    cJSON_AddNumberToObject(dataItem, "value", atof(dtok));
                }
                cJSON_AddNumberToObject(dataItem, "quality", 1);
                cJSON_AddItemToArray(data, dataItem);
            }
            freeReplyObject(dataReply);
            
            bool devonline = false;
            int devalarmcnt = 0;
            
            snprintf(cmd, CMD_SHORT_LENGTH, "zrange %s_alarm 0 -1", devsn);
            redisReply *alarmReply = (redisReply *)redisCommand(var->context, cmd);
            if (NULL == reply || alarmReply->type != REDIS_REPLY_ARRAY)
            {
                dbg_syslog(LOG_ERR, "get dev data failed");
                freeReplyObject(alarmReply);
                continue;
            }
            for (int j = 0; j < alarmReply->elements; j++)
            {
                cJSON *dataItem = cJSON_CreateObject();
                char *dtok = strtok(alarmReply->element[j]->str, "|");
                char dataname[MAX_NAME_LEN] = {0};
                snprintf(dataname, MAX_NAME_LEN, "%s.%s", devsn, dtok);
                cJSON_AddStringToObject(dataItem, "name", dataname);
                
                alarm_entry_t *curalarm;
                HASH_FIND_STR(curdevtype->alarmnode, dtok, curalarm);
                if (curalarm == NULL)
                {
                    dbg_syslog(LOG_ERR, "data node %s not found", dtok);
                    continue;
                }
                dtok = strtok(0, "|");
                if (atoi(dtok) > 0)
                    devalarmcnt++;
                if (curalarm->type == 0)
                {
                    cJSON_AddNumberToObject(dataItem, "value", atoi(dtok));
                }
                else
                {
                    cJSON_AddNumberToObject(dataItem, "value", atof(dtok));
                }
                cJSON_AddNumberToObject(dataItem, "quality", 1);
                cJSON_AddItemToArray(data, dataItem);
            }
            freeReplyObject(alarmReply);

            if (strstr(devsn, "EMS"))
            {
                devonline=true;
            }
            else
            {
                snprintf(cmd, CMD_SHORT_LENGTH, "zrange %s 0 0", devsn);
                redisReply *onlineReply = (redisReply *)redisCommand(var->context, cmd);
                if (NULL == reply || onlineReply->type != REDIS_REPLY_ARRAY || onlineReply->elements<1)
                {
                    dbg_syslog(LOG_ERR, "get dev onlinestat failed");
                    freeReplyObject(onlineReply);
                    continue;
                }
                char *onlinetok = strtok(onlineReply->element[0]->str, "|");
                if (strstr(onlinetok, "Online"))
                    devonline=(atoi(strtok(0,"|"))>0) ? true : false;
                freeReplyObject(onlineReply);
            }
            char *fixname[]={"Online", "vendor", "model", "alarmCnt"};
            for (int j = 0; j < 4; j++)
            {
                stat_entry_t *curstat;
                HASH_FIND_STR(curdevtype->statnode, fixname[j], curstat);
                if(curstat!= NULL)
                {
                    cJSON *dataItem = cJSON_CreateObject();
                    char dataname[MAX_NAME_LEN] = {0};
                    snprintf(dataname, MAX_NAME_LEN, "%s.%s", devsn, fixname[j]);
                    cJSON_AddStringToObject(dataItem, "name", dataname);
                    if (strstr(fixname[j], "Online"))
                        cJSON_AddNumberToObject(dataItem, "value", devonline?1:0);
                    else if (strstr(fixname[j], "alarmCnt"))
                        cJSON_AddNumberToObject(dataItem, "value", devalarmcnt);
                    else
                        cJSON_AddStringToObject(dataItem, "value", curstat->entryValue);
                    cJSON_AddNumberToObject(dataItem, "quality", 1);
                    cJSON_AddItemToArray(data, dataItem);
                }
            }

            snprintf(cmd, CMD_SHORT_LENGTH, "zrange %s_set 0 -1", devsn);
            redisReply *setReply = (redisReply *)redisCommand(var->context, cmd);
            if (NULL == reply || setReply->type != REDIS_REPLY_ARRAY)
            {
                dbg_syslog(LOG_ERR, "get dev data failed");
                freeReplyObject(setReply);
                continue;
            }
            for (int j = 0; j < setReply->elements; j++)
            {
                cJSON *dataItem = cJSON_CreateObject();
                char *dtok = strtok(setReply->element[j]->str, "|");
                char dataname[MAX_NAME_LEN] = {0};
                snprintf(dataname, MAX_NAME_LEN, "%s.%s", devsn, dtok);
                cJSON_AddStringToObject(dataItem, "name", dataname);
                
                set_entry_t *curset;
                HASH_FIND_STR(curdevtype->setnode, dtok, curset);
                if (curset == NULL)
                {
                    dbg_syslog(LOG_ERR, "data node %s not found", dtok);
                    continue;
                }
                dtok = strtok(0, "|");
                if (curset->type == 0)
                {
                    cJSON_AddNumberToObject(dataItem, "value", atoi(dtok));
                }
                else
                {
                    cJSON_AddNumberToObject(dataItem, "value", atof(dtok));
                }
                cJSON_AddNumberToObject(dataItem, "quality", 1);
                cJSON_AddItemToArray(data, dataItem);
            }
            freeReplyObject(setReply);
        }
        freeReplyObject(reply);
    }

    char *data_tmp = cJSON_PrintUnformatted(root);
    snprintf(topic, TOPIC_MAX_LEN, "%s/%s/%s/%s", PLATFORM, LCPOST, var->sn_str, FUNC_DATAREPORT);

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

void ems_smu_subdevmsg(ems_smu_var_t *var, cJSON *root)
{
    char cmd[CMD_SHORT_LENGTH] = {0};
    cJSON *msgs = cJSON_CreateArray();
    cJSON_AddItemToObject(root, "messages", msgs);

    snprintf(cmd, CMD_SHORT_LENGTH, "hgetall devs");
    redisReply *reply_devs = (redisReply *)redisCommand(var->context, cmd);
    if (NULL == reply_devs || (reply_devs->type != REDIS_REPLY_ARRAY) || (reply_devs->elements <= 0))
    {
        dbg_syslog(LOG_ERR, "getall devs failed");
        return;
    }
    for (int i = 0; i < reply_devs->elements; i+=2)
    {
        snprintf(cmd, CMD_MAX_LENGTH, "zrange %s 0 -1", reply_devs->element[i]->str);
        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", reply_devs->element[i]->str);
            continue;
        }
        cJSON *curdevdata = cJSON_CreateObject();
        cJSON_AddItemToArray(msgs, curdevdata);
        cJSON_AddStringToObject(curdevdata, "no", reply_devs->element[i]->str);
        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);
    }
    freeReplyObject(reply_devs);
}

int comm_status_poll(ems_smu_var_t *var)
{
    TIMER_CONFIRM(var->status_timer);
    time_t now = time(NULL);
    static time_t dev_report_time = 0;
    static time_t login_time = 0;

    if (var->state != STAT_CONNECTED && now >= login_time + 5)
    {
        login_time = now;
        gateway_login_msg(var);
    }
    if (var->state == STAT_CONNECTED && now >= dev_report_time + var->dev_interval)
    {
        dev_report_time = now;
        timer_dev_report(var);
    }
    return 0;
}

void ems_smu_pubstra(ems_smu_var_t *var, uint tm_score)
{
    char topic[TOPIC_MAX_LEN * 2] = {0};
    unsigned short segpower[MAX_TIME_SEG] = {0};
    unsigned char segmode[MAX_TIME_SEG] = {0};
    short segwggl[MAX_TIME_SEG] = {0};
    char cmd[CMD_MAX_LENGTH] = {0};
    char timestr[8] = {0};
    cJSON *rt_data = cJSON_CreateObject();
    cJSON *content = cJSON_CreateArray();
    cJSON_AddItemToObject(rt_data, "content", content);
    cJSON *dailystra = cJSON_CreateObject();
    cJSON_AddItemToArray(content, dailystra);
    cJSON *planinfo = cJSON_CreateArray();
    cJSON_AddItemToObject(dailystra, "PlanInfo", planinfo);

    snprintf(cmd, CMD_MAX_LENGTH, "zrangebyscore "REDIS_STRA_BASE" %u %u", tm_score, tm_score);
    redisReply *reply_stra = (redisReply *)redisCommand(var->context, cmd);
    if (NULL == reply_stra)
    {
        dbg_syslog(LOG_ERR, "execcmd %s failed", cmd);
    }
    else if (reply_stra->type == REDIS_REPLY_NIL)
    {
        dbg_syslog(LOG_ERR, "No strategy of specific date %u", tm_score);
    }
    else if (reply_stra->type == REDIS_REPLY_STRING)
    {
        parse_stra(reply_stra->str, segpower, segmode, segwggl);
    }
    else if (reply_stra->type == REDIS_REPLY_ARRAY && reply_stra->elements > 0)
    {
        parse_stra(reply_stra->element[0]->str, segpower, segmode, segwggl);
    }
    freeReplyObject(reply_stra);
    
    for (int i = 0; i < MAX_TIME_SEG; i++)
    {
        cJSON *timeseg = cJSON_CreateObject();
        cJSON_AddNumberToObject(timeseg, "option", segmode[i]);
        cJSON_AddNumberToObject(timeseg, "power", segpower[i]);
        cJSON_AddNumberToObject(timeseg, "Wggl", segwggl[i]);
        snprintf(timestr, 8, "P%02d%02d", i / 4, (((i % 4)<<4)-(i % 4)));
        cJSON_AddStringToObject(timeseg, "time", timestr);
        cJSON_AddItemToArray(planinfo, timeseg);
    }

    snprintf(cmd, CMD_MAX_LENGTH, "hgetall "REDIS_PROPARAM_BASE);
    redisReply *reply_pro = (redisReply *)redisCommand(var->context, cmd);
    if (NULL == reply_pro)
    {
        dbg_syslog(LOG_ERR, "execcmd %s failed", cmd);
    }
    else if (reply_pro->type == REDIS_REPLY_NIL)
    {
        dbg_syslog(LOG_ERR, "No protected param");
    }
    else
    {
        for(int k = 0; k < reply_pro->elements; k+=2)
        {
            tagmap_entry_t *curtag;
            double paramVal = atof(reply_pro->element[k+1]->str);  
            HASH_FIND_STR(var->tagmaps, reply_pro->element[k]->str, curtag);
            if (curtag != NULL)
            {
                cJSON_AddNumberToObject(dailystra, curtag->emsctrltag, paramVal);
            }
            else
            {
                cJSON_AddNumberToObject(dailystra, reply_pro->element[k]->str, paramVal);
            }
        }
    }
    freeReplyObject(reply_pro);
    cJSON_AddNumberToObject(dailystra, "repeat", 1);
    cJSON_AddNumberToObject(dailystra, "priority", 1);
    cJSON_AddNumberToObject(dailystra, "pvMax", 0);
    cJSON_AddNumberToObject(dailystra, "type", 1);
    cJSON_AddStringToObject(dailystra, "weekly", "");

    cJSON_AddStringToObject(rt_data, "identifier", "strategyTemplate");
    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);
}

static int ems_smu_subscribe_all(ems_smu_var_t *var)
{
    char topic[128];

    snprintf(topic, sizeof(topic), "%s/%s/%s/+", PLATFORM, LCPOSTRSP, var->sn_str);
    ipc_session_subscribe(var->session, topic);

    snprintf(topic, sizeof(topic), "%s/%s/%s/+", PLATFORM, GET, var->sn_str);
    ipc_session_subscribe(var->session, topic);

    snprintf(topic, sizeof(topic), "%s/%s/%s/+", PLATFORM, SET, var->sn_str);
    ipc_session_subscribe(var->session, topic);

    return 0;
}

void gateway_login_msg(ems_smu_var_t *var)
{
    char topic[TOPIC_MAX_LEN];

    cJSON *root = cJSON_CreateObject();
    if (!root)
    {
        dbg_syslog(LOG_ERR, "creaet json object error");
        return;
    }
    cJSON_AddStringToObject(root, "funcId", FUNC_LOGIN);
    cJSON_AddStringToObject(root, "lcSN", var->sn_str);
    cJSON_AddNumberToObject(root, "seq", var->seq++);
    cJSON_AddNumberToObject(root, "time", time(NULL));

    cJSON *logininfo = cJSON_CreateObject();
    cJSON_AddStringToObject(logininfo, "id", var->lcid);
    cJSON_AddStringToObject(logininfo, "station", var->lcsta);
    cJSON_AddNumberToObject(logininfo, "capacity", var->capacity);
    cJSON_AddNumberToObject(logininfo, "power", var->power);
    cJSON_AddStringToObject(logininfo, "time", var->loadtime);
    cJSON_AddStringToObject(logininfo, "local", var->loc);
    cJSON_AddStringToObject(logininfo, "station_name", var->staname);
    cJSON_AddStringToObject(logininfo, "vendor", var->vendor);
    cJSON_AddStringToObject(logininfo, "softVersion", var->softver);
    cJSON_AddStringToObject(logininfo, "model", var->model);
    cJSON_AddItemToObject(root, "info", logininfo);

    char *data_tmp = cJSON_PrintUnformatted(root);
    snprintf(topic, TOPIC_MAX_LEN, "%s/%s/%s/%s", PLATFORM, LCPOST, var->sn_str, FUNC_LOGIN);

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

// 建立与broker之间的MQTT连接
int ems_smu_mqtt_init(ems_smu_var_t *var)
{
    char clientId[128] = {0};

    snprintf(clientId, sizeof(clientId), "EMS_SMU_%s", var->sn_str);
    var->session = ipc_session_new(clientId, (void*)var, IPC_DEFAULT);
    if (var->session == NULL) return -1;
    ipc_session_set_callbacks(var->session, recv_host_msg, NULL);
    ems_smu_subscribe_all(var);
    ipc_session_start(var->session);
    return 0;
}
