#include <sys/select.h>
#include <sys/time.h>
#include <sys/types.h>
#include <unistd.h>
#include <sys/syscall.h>
#include <stdlib.h>
#include <stdio.h>
#include <string.h>
#include <sys/stat.h>
#include <sys/types.h>
#include <errno.h>
#include <poll.h>

#include "ems_redis.h"

ems_redis_var_t gvar = {0};
uint8_t racknuminpack[16] = {0};

//重连redis数据库
static int ems_redis_reconn(redisContext *rctx)
{
    dbg_syslog(LOG_ERR, "redisCommand ret null, reconnecting redis server");
    redisFree(rctx);
    rctx = redisConnect(REDIS_SERVER_IP, REDIS_SERVER_PORT);
    if (rctx->err)
    {
        redisFree(rctx);
        dbg_syslog(LOG_DEBUG, "%d connect redis server failure: %s\n", __LINE__, rctx->errstr);
        return false;
    }
    return true;
}

//更新设备在离线状态到数据库
static void ems_redis_writeonline(redisContext *rctx, unsigned char online, char *devsn)
{
    char cmd[CMD_MAX_LENGTH] = {0};
    snprintf(cmd, CMD_MAX_LENGTH, "zremrangebyscore %s 0 0", devsn);
    redisReply *reply = (redisReply *)redisCommand(rctx, cmd);
    freeReplyObject(reply);

    snprintf(cmd, CMD_MAX_LENGTH, "zadd %s 0 Online|%u", devsn, online);
    reply = (redisReply *)redisCommand(rctx, cmd);
    if (NULL == reply)
    {
        dbg_syslog(LOG_ERR, "set %s online %u failed", devsn, online);
    }
    freeReplyObject(reply);
}

//更新设备类型配置
static void parse_devtype(ems_redis_var_t *var, char *path)
{
    char dev_type[DEV_TYPE_LENGTH];    //设备类型
    char dev_type_alias[DEV_TYPE_LENGTH]; //设备类型别名，来自平台
    int common_addr;                    //设备公共地址
    int startaddr1;                     //遥测量起始地址
    int step1;                          //遥测量步长
    int startaddr2;                     //遥信量起始地址
    int step2;                          //遥信量步长
    int startModbus;                    //MobBus起始地址
    int stepModbus;                     //ModBus步长
    
    FILE *fp = fopen(path, "r");
    if (NULL == fp)
    {
        dbg_syslog(LOG_ERR, "fopen %s error\n", path);
        return;
    }
    while(fscanf(fp, "%s %s %d %d %d %d %d %d %d", dev_type, dev_type_alias, &common_addr, &startaddr1, &step1, &startaddr2, &step2, &startModbus, &stepModbus) != EOF)
    {
        devtype_entry_t *curdev;
        HASH_FIND_STR(var->devtypes, dev_type, curdev);
        if (NULL == curdev)
        {
            curdev = malloc(sizeof(devtype_entry_t));
            snprintf(curdev->dev_type, DEV_TYPE_LENGTH, "%s", dev_type);
            snprintf(curdev->dev_type_alias, DEV_TYPE_LENGTH, "%s", dev_type_alias);
            curdev->common_addr = common_addr;
            curdev->startaddr1 = startaddr1;
            curdev->step1 = step1;
            curdev->startaddr2 = startaddr2;
            curdev->step2 = step2;
            curdev->startModbus = startModbus;
            curdev->stepModbus = stepModbus;
            curdev->num = 0;
            curdev->tagscore = 1;
            curdev->tagsconf = NULL;
            HASH_ADD_STR(var->devtypes, dev_type, curdev);
            var->dtypenum++;
        }
    }
    fclose(fp);
}

//根据设备SN获取设备类型
static devtype_entry_t* get_devtype_by_sn(ems_redis_var_t *var, char *devsn)
{
    devtype_entry_t *pos;
    for(pos = var->devtypes; pos != NULL; pos = pos->hh.next)
    {
        if (strstr(devsn, pos->dev_type) || strstr(devsn, pos->dev_type_alias))
            return pos;
    }
    return NULL;
}

//根据设备类型名获取设备类型条目
static devtype_entry_t* get_devtype_by_dtype(ems_redis_var_t *var, char *dtype)
{
    devtype_entry_t *s;
    HASH_FIND_STR(var->devtypes, dtype, s);
    return s;
}

static void read_addr_tag(ems_redis_var_t *var, devtype_entry_t *devtype, char *path)
{
    //char cmd[CMD_MAX_LENGTH] = {0};
    int addr;
    char tag[TAG_MAX_LENGTH];
    int dtype;
    float k;

    FILE *fp = fopen(path, "r");
    if (NULL == fp)
    {
        dbg_syslog(LOG_ERR, "fopen %s error\n", path);
        return;
    }

    while(fscanf(fp, "%d %s %d %f", &addr, tag, &dtype, &k) != EOF)
    {    
        tagconf_entry_t *curtag;
        HASH_FIND_STR(devtype->tagsconf, tag, curtag);
        if (NULL == curtag)
        {
            curtag = malloc(sizeof(tagconf_entry_t));
            snprintf(curtag->tag, TAG_MAX_LENGTH, "%s", tag);
            curtag->addr = addr;
            curtag->dtype = dtype;
            curtag->k = k;
            memset(curtag->emmstag, 0, TAG_MAX_LENGTH);
            memset(curtag->smualarmtag, 0, TAG_MAX_LENGTH);
            HASH_ADD_STR(devtype->tagsconf, tag, curtag);
        }
    }
    fclose(fp);
    return;
}

static void read_emms2_tag(ems_redis_var_t *var, devtype_entry_t *devtype, char *path)
{
    char tag[TAG_MAX_LENGTH] = {0};
    char tag_emms2[TAG_MAX_LENGTH] = {0};
    FILE *fp = fopen(path, "r");
    if (NULL == fp)
    {
        dbg_syslog(LOG_ERR, "fopen %s error", path);
        return;
    }
    while(fscanf(fp, "%s %s", tag, tag_emms2) != EOF)
    {
        tagconf_entry_t *curtag;
        HASH_FIND_STR(devtype->tagsconf, tag, curtag);
        if (NULL == curtag)
        {
            dbg_syslog(LOG_ERR, "tag %s not in tagsconf of devtype %s", tag, devtype->dev_type);
            continue;
        }
        curtag->emmsscore = devtype->tagscore++;
        snprintf(curtag->emmstag, TAG_MAX_LENGTH, "%s", tag_emms2);
    }
    fclose(fp);
}

static void read_smualarm_tag(ems_redis_var_t *var, devtype_entry_t *devtype, char *path)
{
    char tag[TAG_MAX_LENGTH] = {0};
    char tag_smualarm[TAG_MAX_LENGTH] = {0};
    FILE *fp = fopen(path, "r");
    if (NULL == fp)
    {
        dbg_syslog(LOG_ERR, "fopen %s error", path);
        return;
    }
    while(fscanf(fp, "%s %s", tag, tag_smualarm) != EOF)
    {
        tagconf_entry_t *curtag;
        HASH_FIND_STR(devtype->tagsconf, tag, curtag);
        if (NULL == curtag)
        {
            dbg_syslog(LOG_ERR, "tag %s not in tagsconf of devtype %s", tag, devtype->dev_type);
            continue;
        }
        curtag->smualarmscore = devtype->alarmscore++;
        snprintf(curtag->smualarmtag, TAG_MAX_LENGTH, "%s", tag_smualarm);
    }
    fclose(fp);
}

static void read_alarm_tag(ems_redis_var_t *var, char *path)
{
    char tag[TAG_MAX_LENGTH] = {0};
    int alarmcode;
    int sevirity;
    int reasoncode;
    FILE *fp = fopen(path, "r");
    if (NULL == fp)
    {
        dbg_syslog(LOG_ERR, "fopen %s error", path);
        return;
    }
    while(fscanf(fp, "%s %d %d %d", tag, &alarmcode, &sevirity, &reasoncode) != EOF)
    {
        alarmtag_entry_t *curalarm;
        HASH_FIND_STR(var->alarmconf, tag, curalarm);
        if (NULL == curalarm)
        {
            curalarm = malloc(sizeof(alarmtag_entry_t));
            snprintf(curalarm->tag, TAG_MAX_LENGTH, "%s", tag);
            curalarm->alarmcode = alarmcode;
            curalarm->severity = sevirity;
            curalarm->reasoncode = reasoncode;
            HASH_ADD_STR(var->alarmconf, tag, curalarm);
        }
    }
    fclose(fp);
}

static tagconf_entry_t* get_tag_index(devtype_entry_t *devtype, char *str)
{
    tagconf_entry_t *s;
    HASH_FIND_STR(devtype->tagsconf, str, s);
    return s;
}

static int parse_nodecfg(ems_redis_var_t *var, char *nodecfgpath)
{
    char buf[JSON_MAX_LENGTH] = {0};
    char devsn[DEV_SN_LEN];
    char cmd[CMD_MAX_LENGTH];
    //int devs_index = 0;
    int common_addr;
    int dev_number; 

    FILE *fp = fopen(nodecfgpath, "r");
    if (NULL == fp)
    {
        dbg_syslog(LOG_ERR, "No %s file", nodecfgpath);
        return -1;
    }
    fread(buf, 1, JSON_MAX_LENGTH, fp);

    cJSON *nodes = cJSON_Parse(buf);
    cJSON *nodecfg = cJSON_GetObjectItem(nodes, "nodes_cfg");
    if (NULL == nodecfg)
    {
        dbg_syslog(LOG_ERR, "No nodes_cfg exist, exit");
        return -1;
    }

    int node_cnt = cJSON_GetArraySize(nodecfg);
    if (node_cnt <= 0)
    {
        dbg_syslog(LOG_ERR,"Empty dev array");
        return -1;
    }
    
    for (int  i = 0; i < node_cnt; i++)
    {
        cJSON *node = cJSON_GetArrayItem(nodecfg, i);
        
        GET_JSON_VALUE_STRING(node, "sn", devsn);

        devtype_entry_t *devtype = get_devtype_by_sn(var, devsn);
        if (NULL == devtype)
        {
            dbg_syslog(LOG_ERR, "invalid devtype %s in json format", devsn);
            continue;
        }

        if (strstr(devtype->dev_type, "PCS") || strstr(devtype->dev_type, "DUI"))
        {
            common_addr = devtype->common_addr++;
            dev_number = 0;
        }
        else
        {
            common_addr = devtype->common_addr;
            dev_number = devtype->num++;
        }
        
        snprintf(var->dev_table.devtab[var->dev_table.devcount].devsn, DEV_SN_LEN, "%s", devsn);
        snprintf(var->dev_table.devtab[var->dev_table.devcount].devtype, DEV_TYPE_LENGTH, "%s", devtype->dev_type);
        var->dev_table.devtab[var->dev_table.devcount].dev_no = dev_number;
        var->dev_table.devtab[var->dev_table.devcount].common_addr = common_addr;
        var->dev_table.devcount++;
        snprintf(cmd, CMD_MAX_LENGTH, "hmset devs %s %s|%d|%d", devsn, devtype->dev_type, dev_number, common_addr);
        redisReply *reply = (redisReply *)redisCommand(var->context, cmd);
        if (NULL == reply)
        {
            dbg_syslog(LOG_DEBUG, "%d execcmd %s error.\n", __LINE__, cmd);
        }
        freeReplyObject(reply);
    }
    return 0;
}

static int find_devindex(ems_redis_var_t *var, char *devsn)
{
    for(int i = 0; i < var->dev_table.devcount; i++)
    {
        if ((strlen(devsn) == strlen(var->dev_table.devtab[i].devsn)) && (0 == strncmp(devsn, var->dev_table.devtab[i].devsn, strlen(devsn))))
            return i;
    }
    return -1;
}

#ifdef DB_IEC104
static void ems_redis_mqtt_spontaneous_report(ems_redis_var_t *var, cJSON *spon_yx, cJSON *spon_yc, char *devsn, int cmaddr, char *devtype)
{
    char topic[TOPIC_MAX_LEN] = {0};
    snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/lnxall/device/changeReport", var->sn_str);

    cJSON *rt_data = cJSON_CreateObject();
    cJSON_AddStringToObject(rt_data, "devsn", devsn);
    cJSON_AddStringToObject(rt_data, "devtype", devtype);
    cJSON_AddStringToObject(rt_data, "identifier", "changeReport");
    cJSON_AddNumberToObject(rt_data, "mi", var->mi++);
    cJSON_AddNumberToObject(rt_data, "time", time(NULL));
    cJSON_AddNumberToObject(rt_data, "ca", cmaddr);
    cJSON_AddItemToObject(rt_data, "yx", spon_yx);
    cJSON_AddItemToObject(rt_data, "yc", spon_yc);

    char *data_tmp = cJSON_Print(rt_data);
    dbg_syslog(LOG_DEBUG, "publish %s payload %s", topic, data_tmp);
    ipc_session_publish(var->session, topic, data_tmp, strlen(data_tmp));
    free(data_tmp);
    cJSON_Delete(rt_data);
}
#endif
static void ems_redis_mqtt_alarm_report(ems_redis_var_t *var, cJSON *msg)
{
    char topic[TOPIC_MAX_LEN] = {0};
    snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/lnxall/device/alarmReport", var->sn_str);

    cJSON *rt_data = cJSON_CreateObject();
    cJSON_AddStringToObject(rt_data, "identifier", "alarmReport");
    cJSON_AddNumberToObject(rt_data, "mi", var->mi++);
    cJSON_AddNumberToObject(rt_data, "time", time(NULL));
    cJSON_AddItemToObject(rt_data, "alarmmsg", msg);

    char *data_tmp = cJSON_Print(rt_data);
    dbg_syslog(LOG_DEBUG, "publish %s payload %s", topic, data_tmp);
    ipc_session_publish(var->session, topic, data_tmp, strlen(data_tmp));
    free(data_tmp);
    cJSON_Delete(rt_data);
}

static int ems_redis_getglobal_rackno(int devno, int num)
{
    int globalno = 0;
    for(int i = 0; i + 18 < devno; i++)
        globalno += racknuminpack[i];
    return globalno + num;
}

static int ems_redis_mqtt_handle_recv_msg(void *obj, ipc_msg_t *mqtt_msg)
{
    char cmd[CMD_MAX_LENGTH] = {0};
#ifdef DB_MODBUS
    char cmd_modbus[CMD_MAX_LENGTH]  = {0};
#endif
    char devsn[DEV_SN_LEN];
    ems_redis_var_t *var = (ems_redis_var_t *)obj;
    devtype_entry_t *cur_deventry;
    devtype_entry_t *dx_deventry;
    time_t now = time(NULL);
    dbg_syslog(LOG_DEBUG, "MQTT client: received MQTT topic: %s payload length: %d payload %s", 
        mqtt_msg->topic, mqtt_msg->payloadLen, mqtt_msg->payload);

    if (strstr(mqtt_msg->topic, TOPIC_IPCALARM)) //EMS告警信息
    {
        int alarmCode=0;
        int reasonCode=-1;
        int severity=0;
        time_t raisetime;
        char entityInst[DEV_SN_LEN] = {0};
        char addInfo[ADD_INFO_LENGTH] = {0};
        struct tm* tm_detail;

        cJSON *root = cJSON_Parse(mqtt_msg->payload);
        if (!root)
        {
            dbg_syslog(LOG_DEBUG, "parse payload failed");
            return -1;
        }

        GET_JSON_VALUE_INT(root, "alarmCode", alarmCode);
        GET_JSON_VALUE_INT(root, "reasonCode", reasonCode);
        GET_JSON_VALUE_INT(root, "severity", severity);
        GET_JSON_VALUE_INT(root, "raiseTime", raisetime);
        GET_JSON_VALUE_STRING(root, "entityInstance", entityInst);
        if (cJSON_HasObjectItem(root, "additionalParam"))
            GET_JSON_VALUE_STRING(root, "additionalParam", addInfo);
        
        tm_detail = localtime(&raisetime);
        if (-1 == reasonCode)
        {
            snprintf(cmd, CMD_MAX_LENGTH, "hdel "REDIS_REAL_ALARM_TABLE" %d|%d|%s", alarmCode, severity, entityInst);  //告警清除，从实时告警列表删除
        } 
        else
        {
            snprintf(cmd, CMD_MAX_LENGTH, "hset "REDIS_REAL_ALARM_TABLE" %d|%d|%s %d|%lu", alarmCode, severity, entityInst, reasonCode, raisetime); //告警发生，存入实时告警列表
            if (addInfo[0] != 0)
            {
                strcat(cmd, "|");
                strncat(cmd, addInfo, ADD_INFO_LENGTH);
            }
        }
        redisReply *reply_alarm = (redisReply *)redisCommand(var->context, cmd);
        if (NULL == reply_alarm)
        {
            dbg_syslog(LOG_DEBUG, "execcmd %s error", cmd);
        }
        freeReplyObject(reply_alarm);


        snprintf(cmd, CMD_MAX_LENGTH, "zadd alarm-%04d%02d%02d %lu %d|%d|%d|%lu|%s", tm_detail->tm_year+1900, tm_detail->tm_mon+1, tm_detail->tm_mday,
                raisetime, alarmCode, reasonCode, severity, raisetime, entityInst);  //告警记录存入历史告警列表
        if (addInfo[0] != 0)
        {
            strcat(cmd, "|");
            strncat(cmd, addInfo, ADD_INFO_LENGTH);
        }
        reply_alarm = (redisReply *)redisCommand(var->context, cmd);
        if (NULL == reply_alarm)
        {
            dbg_syslog(LOG_DEBUG, "execcmd %s error", cmd);
        }
        freeReplyObject(reply_alarm);

        cJSON *alarmmsg = cJSON_CreateObject();
        cJSON_AddNumberToObject(alarmmsg, "alarmCode", alarmCode);
        cJSON_AddNumberToObject(alarmmsg, "reasonCode", reasonCode);
        cJSON_AddNumberToObject(alarmmsg, "severity", severity);
        cJSON_AddNumberToObject(alarmmsg, "raiseTime", raisetime);
        cJSON_AddStringToObject(alarmmsg, "entityInstance", entityInst);
        if (addInfo[0] != 0)
        {
            cJSON_AddStringToObject(alarmmsg, "additionalParam", addInfo);
        }
        ems_redis_mqtt_alarm_report(var, alarmmsg);
    }
    else if (strstr(mqtt_msg->topic, TOPIC_SERVICE))
    {
        //int ret;
        //char *sn = NULL;
        char devtype[DEV_TYPE_LENGTH] = {0};
        int devno = 0;
        int cmaddr = 0; //common address，设备公共地址
        int dx_ca = 0; //对于簇数据，计录电芯的北向公共地址
        int node_cnt;
        int i;
        //int addr;  //点位地址
        int addr_modbus; //modbus 点位地址
        int dtype; //点位类型
        float k; //点位系数k
        //int step1,step2;
        //int devs_index;
        time_t datatime; //数据报文中的时间
        dx_deventry = get_devtype_by_dtype(var, "DX");
        char zsetname[DEV_SN_LEN] = {0};

        cJSON *root = cJSON_Parse(mqtt_msg->payload);
        if (!root)
        {
            dbg_syslog(LOG_DEBUG, "parse payload failed");
            return -1;
        }
        GET_JSON_VALUE_STRING(root, "sn", devsn);

        GET_JSON_VALUE_INT(root, "time", datatime);

        struct tm* tm_detail;
        tm_detail = localtime(&datatime); 

        int curdevind = find_devindex(var, devsn);
        if (curdevind == -1)
        {
            dbg_syslog(LOG_ERR, "device %s not found in devs table", devsn);
            cJSON_Delete(root);
            return -1;
        }
        snprintf(devtype, DEV_TYPE_LENGTH, "%s", var->dev_table.devtab[curdevind].devtype);
        devno = var->dev_table.devtab[curdevind].dev_no;
        cmaddr = var->dev_table.devtab[curdevind].common_addr;
        var->dev_table.devtab[curdevind].lastmsgtime = now;
        if (var->dev_table.devtab[curdevind].online == 0)
        {
            var->dev_table.devtab[curdevind].online = 1;
            ems_redis_writeonline(var->context, 1, devsn);
        }
        dbg_syslog(LOG_ERR, "devtype %s %d %d", devtype, devno, cmaddr);

        cur_deventry = get_devtype_by_dtype(var, devtype);
        if (NULL == cur_deventry)
        {
            dbg_syslog(LOG_ERR, "invalid devtype %s in redis", devtype);
            cJSON_Delete(root);
            return -1;
        }
        
        cJSON *tags = cJSON_GetObjectItem(root, "tags");
        if (!tags)
        {
            dbg_syslog(LOG_ERR, "%s get tags failed.",devsn);
            cJSON_Delete(root);
            return -1;
        }

#ifdef DB_IEC104
        cJSON *spon_yx = cJSON_CreateObject(); //变位遥信
        cJSON *spon_yc = cJSON_CreateObject(); //变位遥测
        cJSON *spon_yx_dx = cJSON_CreateObject();
        cJSON *spon_yc_dx = cJSON_CreateObject();
        bool spon_rep = false;
        bool spon_rep_dx = false;
        step1 = cur_deventry->step1;
        step2 = cur_deventry->step2;
#endif
        node_cnt = cJSON_GetArraySize(tags);
        dbg_syslog(LOG_DEBUG, "devtype %s devno %d common address %d node_cnt %d", devtype, devno, cmaddr, node_cnt);
       
        for (i = 0; i < node_cnt; i++)
        {
            int dxflag = 0; //电芯数据在簇中，检查此tag是否是电芯数据
            char db_data[16] = {0}; //将数据转换成存入redis数据库的字符串
            tagconf_entry_t* tagind;
            cJSON *node = cJSON_GetArrayItem(tags, i);
            if (node->type != cJSON_Number)
                continue;
#ifdef DB_IEC104
            tagind = get_tag_index(cur_deventry, node->string);
            if (NULL == tagind)
            {
                if (strstr(devtype, "CU") && (dx_deventry != NULL))
                {
                    tagind = get_tag_index(dx_deventry, node->string);
                    if (NULL == tagind)
                    {
                        dbg_syslog(LOG_DEBUG, "TAG %s not found in dx and %s", node->string, devtype);
                        continue;
                    }
                    else
                    {
                        dxflag = 1;
                    }
                }
                else
                {
                    dbg_syslog(LOG_DEBUG, "TAG %s not found in %s", node->string, devtype);
                    continue;
                }
            }
            addr = tagind->addr;
            k = tagind->k;
            dtype = tagind->dtype;

            if (1 == dxflag)
            {
                if (dx_deventry != NULL)
                {
                    int global_rackno = ems_redis_getglobal_rackno(cmaddr ,devno);
                    dx_ca = 49 + global_rackno / 8;
                    addr += dx_deventry->step1 * (global_rackno % 8);
                }
            }
            else
            {
                if (addr >= cur_deventry->startaddr1 && addr <= cur_deventry->startaddr1 + step1)
                    addr += step1 * devno;
                else
                    addr += step2 * devno;
            }
            
            if (dtype < REDIS_DATA_F32)
                snprintf(db_data, 16, "%d", (int)(node->valuedouble/k));
            else
                snprintf(db_data, 16, "%.4lf", node->valuedouble/k);
#endif

#ifdef DB_MODBUS
            tagind = get_tag_index(cur_deventry, node->string);
            if (tagind == NULL && dx_deventry != NULL)
            {
                if (strstr(devtype, "CU") && dx_deventry != NULL)
                {
                    tagind = get_tag_index(dx_deventry, node->string);
                    if (NULL == tagind)
                    {
                        dbg_syslog(LOG_DEBUG, "TAG %s not found in dx and %s", node->string, devtype);
                        continue;
                    }
                    else
                    {
                        dxflag = 1;
                    }
                }
                else
                {
                    dbg_syslog(LOG_DEBUG, "TAG %s not found in %s", node->string, devtype);
                    continue;
                }
            }
            
            addr_modbus = tagind->addr;
            dtype = tagind->dtype;
            k = tagind->k;
            if (1 == dxflag)
            {
                if (dx_deventry != NULL)
                {
                    int global_rackno = ems_redis_getglobal_rackno(cmaddr ,devno);
                    dx_ca = 49 + global_rackno / 8;
                    addr_modbus += dx_deventry->stepModbus * (global_rackno % 8);
                }
            }
            else
            {
                addr_modbus += cur_deventry->stepModbus * devno;
            }

            if (dtype < REDIS_DATA_F32)
                snprintf(db_data, 16, "%d", (int)(node->valuedouble/k));
            else
                snprintf(db_data, 16, "%.4lf", node->valuedouble/k);
#endif /* DB_MODBUS */

            //检查该tag的值是否发生变化
            bool value_changed = false;
            if (dxflag)
            {
#ifdef DB_IEC104
                snprintf(cmd, CMD_MAX_LENGTH, "zrangebyscore %d %d %d", dx_ca, addr, addr);
#else
                snprintf(cmd, CMD_MAX_LENGTH, "zrangebyscore %d_modbus %d %d", dx_ca, addr_modbus, addr_modbus);
#endif
            }
            else
            {
#ifdef DB_IEC104
                snprintf(cmd, CMD_MAX_LENGTH, "zrangebyscore %d %d %d", cmaddr, addr, addr);
#else
                snprintf(cmd, CMD_MAX_LENGTH, "zrangebyscore %d_modbus %d %d", cmaddr, addr_modbus, addr_modbus);
#endif
            }
            redisReply *reply_val = (redisReply *)redisCommand(var->context, cmd);
            if (NULL == reply_val)
            {
                dbg_syslog(LOG_ERR, "%d execcmd %s error.\n", __LINE__, cmd);
                ems_redis_reconn(var->context);
                continue;
            }
            if ((reply_val->type == REDIS_REPLY_ARRAY) && (reply_val->elements > 0))
            {
                //char addr_str[8] = {0};
                if ((strncmp(db_data, reply_val->element[0]->str, strlen(db_data)) != 0) || (reply_val->element[0]->str[strlen(db_data)] != '|'))
                {
                    value_changed = true;
#ifdef DB_IEC104
                    snprintf(addr_str, 8, "%d", addr);
                    if (dxflag == 1)
                    {
                        spon_rep_dx = true;
                        if (dtype < REDIS_DATA_U32)
#ifdef FEATURE_MQTT
                            cJSON_AddNumberToObject(spon_yx_dx, node->string, node->valueint);
#else
                            cJSON_AddNumberToObject(spon_yx_dx, addr_str, node->valueint);
#endif
                        else
#ifdef FEATURE_MQTT
                            cJSON_AddStringToObject(spon_yc_dx, node->string, db_data);
#else
                            cJSON_AddStringToObject(spon_yc_dx, addr_str, db_data);
#endif
                    }
                    else 
                    {
                        spon_rep = true;
                        if (dtype < REDIS_DATA_U32)
#ifdef FEATURE_MQTT
                            cJSON_AddNumberToObject(spon_yx, node->string, node->valueint);
#else
                            cJSON_AddNumberToObject(spon_yx, addr_str, node->valueint);
#endif
                        else
#ifdef FEATURE_MQTT
                            cJSON_AddStringToObject(spon_yc, node->string, db_data);
#else
                            cJSON_AddStringToObject(spon_yc, addr_str, db_data);
#endif
                    }
#endif
                }
            }
            if (((REDIS_REPLY_ARRAY == reply_val->type) && (reply_val->elements == 0))|| value_changed)
            {
#ifdef DB_IEC104
                if (dxflag)
                {
                    snprintf(cmd, CMD_MAX_LENGTH, "zremrangebyscore %d %u %u", dx_ca, addr, addr)
                }
                else 
                {
                    snprintf(cmd, CMD_MAX_LENGTH, "zremrangebyscore %d %u %u", cmaddr, addr, addr)
                }
#endif
#ifdef DB_MODBUS
                if (dxflag)
                {
                    snprintf(cmd, CMD_MAX_LENGTH, "zremrangebyscore %d_modbus %u %u", dx_ca, addr_modbus, addr_modbus);
                }
                else 
                {
                    snprintf(cmd, CMD_MAX_LENGTH, "zremrangebyscore %d_modbus %u %u", cmaddr, addr_modbus, addr_modbus);
                }
#endif 
                redisReply *reply_addr = (redisReply *)redisCommand(var->context, cmd);
                if (reply_addr == NULL)
                {
                    dbg_syslog(LOG_INFO, "execcmd %s error", cmd);
                }
                freeReplyObject(reply_addr);

#ifdef DB_IEC104
                if (dxflag)
                {
                    snprintf(cmd, CMD_MAX_LENGTH, "zadd %d %d %s|%d|%s|%lu", 
                        dx_ca, addr, db_data, dtype, node->string, now);
                }
                else
                {
                    snprintf(cmd, CMD_MAX_LENGTH, "zadd %d %d %s|%d|%s|%lu", 
                        cmaddr, addr, db_data, dtype, node->string, now);
                }
#endif
#ifdef DB_MODBUS
                if (dxflag)
                {
                    snprintf(cmd_modbus, CMD_MAX_LENGTH, "zadd %d_modbus %d %s|%d|%s|%lu", 
                        dx_ca, addr_modbus, db_data, dtype, node->string, now);
                }
                else
                {
                    snprintf(cmd_modbus, CMD_MAX_LENGTH, "zadd %d_modbus %d %s|%d|%s|%lu", 
                        cmaddr, addr_modbus, db_data, dtype, node->string, now);
                }
#endif
                
#ifdef DB_IEC104
                dbg_syslog(LOG_DEBUG, "cmd: %s", cmd);
                redisReply *reply4 = (redisReply *)redisCommand(var->context, cmd);
                if (NULL == reply4)
                {
                    dbg_syslog(LOG_DEBUG, "%d execcmd %s error.\n", __LINE__, cmd);
                }
                freeReplyObject(reply4);
                reply4 = NULL;
#endif

#ifdef DB_MODBUS
                dbg_syslog(LOG_DEBUG, "cmd: %s", cmd_modbus);
                redisReply *reply5 = (redisReply *)redisCommand(var->context, cmd_modbus);
                if (NULL == reply5)
                {
                    dbg_syslog(LOG_DEBUG, "%d execcmd %s error.\n", __LINE__, cmd_modbus);
                }
                freeReplyObject(reply5);
#endif

                if (tagind->smualarmtag[0] != 0) //此tag有对应的smu alarm tag，为设备告警数据且需存入设备告警数据表
                {                    
                    if (dxflag)
                    {
                        snprintf(zsetname, DEV_SN_LEN, "%s_CELL_alarm", devsn);
                    }
                    else
                    {
                        snprintf(zsetname, DEV_SN_LEN, "%s_alarm", devsn);
                    }
                    
                    //redis有序集同一个分数可以存多个条目，因此先删掉此分数下的旧数据再添加新数据
                    snprintf(cmd, CMD_MAX_LENGTH, "zremrangebyscore %s %u %u", zsetname, tagind->smualarmscore, tagind->smualarmscore);
                    redisReply *reply_emmsv2 = (redisReply *)redisCommand(var->context, cmd);
                    if (NULL == reply_emmsv2)
                    {
                        dbg_syslog(LOG_ERR, "execcmd %s error", cmd);
                    }
                    freeReplyObject(reply_emmsv2);

                    snprintf(cmd, CMD_MAX_LENGTH, "zadd %s %u %s|%d", zsetname, tagind->smualarmscore, tagind->smualarmtag, node->valueint);
                    reply_emmsv2 = (redisReply *)redisCommand(var->context, cmd);
                    if (NULL == reply_emmsv2)
                    {
                        dbg_syslog(LOG_ERR, "execcmd %s error", cmd);
                    }
                    freeReplyObject(reply_emmsv2);
                }

                if (tagind->emmstag[0] != 0) //此tag有对应的emmsv2 tag，为设备数据且需存入实时数据表
                {                    
                    if (dxflag)
                    {
                        snprintf(zsetname, DEV_SN_LEN, "%s_CELL", devsn);
                    }
                    else
                    {
                        snprintf(zsetname, DEV_SN_LEN, "%s", devsn);
                    }
                    
                    //redis有序集同一个分数可以存多个条目，因此先删掉此分数下的旧数据再添加新数据
                    snprintf(cmd, CMD_MAX_LENGTH, "zremrangebyscore %s %u %u", zsetname, tagind->emmsscore, tagind->emmsscore);
                    redisReply *reply_emmsv2 = (redisReply *)redisCommand(var->context, cmd);
                    if (NULL == reply_emmsv2)
                    {
                        dbg_syslog(LOG_ERR, "execcmd %s error", cmd);
                    }
                    freeReplyObject(reply_emmsv2);

                    snprintf(cmd, CMD_MAX_LENGTH, "zadd %s %u %s|%.3lf", zsetname, tagind->emmsscore, tagind->emmstag, node->valuedouble);
                    reply_emmsv2 = (redisReply *)redisCommand(var->context, cmd);
                    if (NULL == reply_emmsv2)
                    {
                        dbg_syslog(LOG_ERR, "execcmd %s error", cmd);
                    }
                    freeReplyObject(reply_emmsv2);
                }
                else //查找此tag是否在全局告警tag表中
                {
                    alarmtag_entry_t *curtag;
                    HASH_FIND_STR(var->alarmconf, node->string, curtag);
                    if (curtag != NULL) //此tag为告警数据且发生突变，需更新实时告警表并将告警发生或清除的记录存入历史数据表
                    {
                        int reasCode = curtag->reasoncode;
                        if (node->valueint <= 0)
                        {
                            reasCode = -1;
                        }
                        if (-1 == reasCode)
                        {
                            snprintf(cmd, CMD_MAX_LENGTH, "hdel "REDIS_REAL_ALARM_TABLE" %d|%d|%s", curtag->alarmcode, curtag->severity, devsn);  //告警清除，从实时告警列表删除
                        }
                        else
                        {
                            snprintf(cmd, CMD_MAX_LENGTH, "hset "REDIS_REAL_ALARM_TABLE" %d|%d|%s %d|%lu", curtag->alarmcode, curtag->severity, devsn, reasCode, datatime); //告警发生，存入实时告警列表
                        }
                        redisReply *reply_alarm = (redisReply *)redisCommand(var->context, cmd);
                        if (NULL == reply_alarm)
                        {
                            dbg_syslog(LOG_DEBUG, "execcmd %s error", cmd);
                        }
                        freeReplyObject(reply_alarm);

                        snprintf(cmd, CMD_MAX_LENGTH, "zadd alarm-%04d%02d%02d %lu %d|%d|%d|%lu|%s", tm_detail->tm_year+1900, tm_detail->tm_mon+1, tm_detail->tm_mday,
                                datatime, curtag->alarmcode, reasCode, curtag->severity, datatime, devsn);  //告警记录存入历史告警列表
                        reply_alarm = (redisReply *)redisCommand(var->context, cmd);
                        if (NULL == reply_alarm)
                        {
                            dbg_syslog(LOG_DEBUG, "execcmd %s error", cmd);
                        }
                        freeReplyObject(reply_alarm);

                        cJSON *alarmmsg = cJSON_CreateObject();
                        cJSON_AddNumberToObject(alarmmsg, "alarmCode", curtag->alarmcode);
                        cJSON_AddNumberToObject(alarmmsg, "reasonCode", reasCode);
                        cJSON_AddNumberToObject(alarmmsg, "severity", curtag->severity);
                        cJSON_AddNumberToObject(alarmmsg, "raiseTime", datatime);
                        cJSON_AddStringToObject(alarmmsg, "entityInstance", devsn);
                        ems_redis_mqtt_alarm_report(var, alarmmsg);
                    }
                }

            }
            freeReplyObject(reply_val);
            reply_val = NULL;
        }
        cJSON_Delete(root);

#ifdef DB_IEC104
        if (spon_rep)
        {
            ems_redis_mqtt_spontaneous_report(var, spon_yx, spon_yc, devsn, cmaddr, devtype);
        }
        else 
        {
            cJSON_Delete(spon_yx);
            cJSON_Delete(spon_yc);
        }
        if (spon_rep_dx)
        {
            ems_redis_mqtt_spontaneous_report(var, spon_yx_dx, spon_yc_dx, devsn, dx_ca, devtype);
        }
        else 
        {
            cJSON_Delete(spon_yx_dx);
            cJSON_Delete(spon_yc_dx);
        }
#endif
    }
    else if (strstr(mqtt_msg->topic, TOPIC_SETEMS))
    {
        cJSON *root = cJSON_Parse(mqtt_msg->payload);
        if (!root)
        {
            dbg_syslog(LOG_ERR, "parse setems failed");
            return -1;
        }

    }
    dbg_syslog(LOG_DEBUG, "handle recv msg finish.");
    return 0;
}

static int ds_flush_timer(ems_redis_var_t *var)
{
    TIMER_CONFIRM(var->ems_status_timer);

    static time_t last_snapshot_time = 0;
    struct tm* t;
    time_t now;
    time(&now);
    t = localtime(&now);
    char cmd[DEVDATA_CMD_LENGTH] = {0};
    
    //检查设备在线状态
    for (int devind = 0; devind < var->dev_table.devcount; devind++)
    {
        if ((now - var->dev_table.devtab[devind].lastmsgtime > var->devstatus_interval) && (var->dev_table.devtab[devind].online == 1)) //设备离线
        {
            var->dev_table.devtab[devind].online = 0;
            ems_redis_writeonline(var->snapcontext, 0, var->dev_table.devtab[devind].devsn);
        }
    }

    //设备历史数据保存
    if (var->snapshot_enable && (now >= last_snapshot_time + var->snapshot_interval))
    {
        last_snapshot_time = now;
        var->snapshot_score = t->tm_hour * 10000 + t->tm_min * 100 + t->tm_sec;
        
        //为避免过多占用内存，每日凌晨删除电表具体点位的两天前的数据
        if (t->tm_hour == 0 && t->tm_min == 0)
        {
            struct tm* t_del;
            time_t now_del = now - (2 * 24 * 3600); //根据当前时间计算要删除的时间的日期
            t_del = localtime(&now_del);
            snprintf(cmd, CMD_MAX_LENGTH, "keys *meter*.*-%04d-%02d-%02d", t_del->tm_year+1900, t_del->tm_mon+1, t_del->tm_mday);
            redisReply *reply_del = (redisReply *)redisCommand(var->snapcontext, cmd);
            if ((NULL == reply_del) || (reply_del->type != REDIS_REPLY_ARRAY) || (reply_del->elements <= 0))
            {
                dbg_syslog(LOG_ERR, "acquire meter history data failed");
            }
            else
            {
                for(int i = 0; i < reply_del->elements; i++)
                {
                    snprintf(cmd, CMD_MAX_LENGTH, "del %s", reply_del->element[i]->str);
                    redisReply *reply_kdel = (redisReply *)redisCommand(var->snapcontext, cmd);
                    freeReplyObject(reply_kdel);
                }
            }
            freeReplyObject(reply_del);
        }

        //遍历所有设备，将设备数据添加到设备历史数据表
        for (int i = 0; i < var->dev_table.devcount; i++)
        {
            dbg_syslog(LOG_ERR, "index %d dev %s", i, var->dev_table.devtab[i].devsn);
            char fulldata[DEVDATA_MAX_LENGTH] = {0};
            snprintf(cmd, DEVDATA_CMD_LENGTH, "zrange %s 0 -1", var->dev_table.devtab[i].devsn);
            redisReply *reply = (redisReply *)redisCommand(var->snapcontext, cmd);
            
            if (NULL == reply)
            {
                dbg_syslog(LOG_ERR, "execcmd %s error", cmd);
                continue;
            }   
            dbg_syslog(LOG_ERR, "reply->type %u", reply->type);
            if (reply->elements <= 0)
            {
                freeReplyObject(reply);
                continue;
            }
            dbg_syslog(LOG_ERR, "reply->elements %lu", reply->elements);
            for (int eleno = 0; eleno < reply->elements; eleno++)
            {
                char *tok;
                tok = strtok(reply->element[eleno]->str, "|");
                tok = strtok(0, "|");
                strcat(fulldata, tok);
                if (eleno + 1 < reply->elements)
                    strcat(fulldata, "|");
            }
            freeReplyObject(reply);
            snprintf(cmd, DEVDATA_CMD_LENGTH, "zadd %s-%04d%02d%02d %u %lu|%s", var->dev_table.devtab[i].devsn,
                    t->tm_year + 1900, t->tm_mon + 1, t->tm_mday, var->snapshot_score, now, fulldata);
            dbg_syslog(LOG_ERR, "cmd %s", cmd);
            reply = (redisReply *)redisCommand(var->snapcontext, cmd);
            if (NULL == reply)
            {
                dbg_syslog(LOG_ERR, "execcmd zadd %s error", var->dev_table.devtab[i].devsn);
                continue;
            }
            freeReplyObject(reply);

            if (strstr(var->dev_table.devtab[i].devsn, "Rack")) //遍历到簇，处理电芯数据
            {
                memset(fulldata, 0, DEVDATA_MAX_LENGTH);
                snprintf(cmd, DEVDATA_CMD_LENGTH, "zrange %s_CELL 0 -1", var->dev_table.devtab[i].devsn);
                redisReply *reply = (redisReply *)redisCommand(var->snapcontext, cmd);
                if (NULL == reply)
                {
                    dbg_syslog(LOG_ERR, "execcmd %s error", cmd);
                    continue;
                }   
                if (reply->elements <= 0)
                {
                    freeReplyObject(reply);
                    continue;
                }
                for (int eleno = 0; eleno < reply->elements; eleno++)
                {
                    char *tok;
                    tok = strtok(reply->element[eleno]->str, "|");
                    tok = strtok(0, "|");
                    strcat(fulldata, tok);
                    if (eleno + 1 < reply->elements)
                        strcat(fulldata, "|");
                }
                freeReplyObject(reply);
                snprintf(cmd, DEVDATA_CMD_LENGTH, "zadd %s_CELL-%04d-%02d-%02d %u %lu|%s", var->dev_table.devtab[i].devsn,
                    t->tm_year + 1900, t->tm_mon + 1, t->tm_mday, var->snapshot_score, now, fulldata);
                reply = (redisReply *)redisCommand(var->snapcontext, cmd);
                if (NULL == reply)
                {
                    dbg_syslog(LOG_ERR, "execcmd zadd %s error", var->dev_table.devtab[i].devsn);
                    continue;
                }
                freeReplyObject(reply);
            }
            else if (strstr(var->dev_table.devtab[i].devsn, "meter"))   //电表历史数据
            {
                snprintf(cmd, CMD_MAX_LENGTH, "zrange %s 0 -1", var->dev_table.devtab[i].devsn);
                reply = (redisReply *)redisCommand(var->snapcontext, cmd);
                if (NULL == reply)
                {
                    dbg_syslog(LOG_ERR, "execcmd %s error", cmd);
                    continue;
                }
                if (reply->elements <= 0)
                {
                    freeReplyObject(reply);
                    continue;
                }
                for (int eleno = 0; eleno < reply->elements; eleno++)
                {
                    char *tok = strtok(reply->element[eleno]->str, "|");
                    char *tok2 = strtok(0, "|");
                    snprintf(cmd, CMD_MAX_LENGTH, "zadd %s.%s-%04d-%02d-%02d %u %s|%lu", var->dev_table.devtab[i].devsn, tok,
                            t->tm_year+1900, t->tm_mon+1, t->tm_mday, var->snapshot_score, tok2, now);
                    redisReply *devreply = (redisReply *)redisCommand(var->snapcontext, cmd);
                    freeReplyObject(devreply);
                }
                freeReplyObject(reply);
            }
        }
    }

    
    return 0;
}

static int ems_redis_mqtt_client_init(ems_redis_var_t *var)
{
    char clientId[MAX_CLIENT_ID_LEN] = {0};
    char topic[TOPIC_MAX_LEN] = {0};
    snprintf(clientId, MAX_CLIENT_ID_LEN, "INIT_EMS_REDIS");
    var->session = ipc_session_new(clientId, (void*)var, IPC_DEFAULT);
    if (var->session == NULL)
        return -1;

    ipc_session_set_callbacks(var->session, ems_redis_mqtt_handle_recv_msg, NULL);
    snprintf(topic, TOPIC_MAX_LEN, "ipc/+/+/device/+/data_filtered/service/+"); //订阅此topic下的所有设备采集的数据
    ipc_session_subscribe(var->session, topic);
    ipc_session_start(var->session);
    return 0;
}

static void ems_redis_loop(ems_redis_var_t *var)
{
    int ret = -1;
    struct pollfd pfd;
    while(1)
    {
        pfd.fd = var->ems_status_timer;
        pfd.events = POLLIN;
        pfd.revents = 0;
        ret = poll(&pfd, 0x1, 2000);
        if (ret < 0)
        {
            int error = errno;	
            dbg_syslog(LOG_WARNING, "ems_status_timer %d errno %d\n", var->ems_status_timer, error);

            if (error == EINTR) {
                continue;
            }
            else {
                break;
            }
        }
        else if(ret == 0)
        {
            /* do nothing */
        }
        else if(ret > 0)
        {
            if (var->ems_status_timer > 0 && pfd.revents)
            {
                ds_flush_timer(var);
            }
        }
    }
} 

//加载ems_redis_cfg.json配置文件
static int ems_redis_loadconfig(ems_redis_var_t *var, char *cfg_file)
{
    char *data = NULL;
    cJSON *root = NULL;
    int ret = -1;
    int lens;

    if (cfg_file == NULL)
    {
        dbg_syslog(LOG_ERR, "cfg_file:%s not exist, use defaut MQTT server", cfg_file);
        goto out;
    }

    data = read_file_data_and_lens(cfg_file, &lens);
    if (data == NULL)
    {
        dbg_syslog(LOG_ERR, "read cfg file %s error, load default", cfg_file);
        goto out;
    }

    root = cJSON_Parse(data);
    if (!root)
    {
        dbg_syslog(LOG_ERR, "parse cfg file %s error, load default", cfg_file);
        goto out;
    }

    GET_JSON_VALUE_INT(root, "snapshot_interval", var->snapshot_interval);
    GET_JSON_VALUE_INT(root, "devstatus_interval", var->devstatus_interval);
    GET_JSON_VALUE_INT(root, "snapshot_enable", var->snapshot_enable);

    cJSON_Delete(root);
    ret = 0;

out:
    if (data)
    {
        free(data);
    }
    if (var->snapshot_interval == 0)
    {
        var->snapshot_interval = 60; // default 1 min
    }
    if (var->devstatus_interval == 0)
    {
        var->devstatus_interval = 10; // default 10 secs
    }
    return ret;
}

static void ems_redis_init(ems_redis_var_t *var)
{
    char addr_tag_path[ADDR_TAG_PATH_LEN] = {0};
    //char devtype_modbus[DEV_TYPE_MODBUS_LENGTH] = {0};
    //char cmd[CMD_MAX_LENGTH] = {0};

    get_board_sn(var->sn_str);
    dbg_syslog(LOG_INFO, "board SN:%s", var->sn_str);

    ems_redis_loadconfig(var, EMS_REDIS_CFG_PATH);
    dbg_syslog(LOG_INFO, "var->snapshot_interval: %d", var->snapshot_interval);

    parse_devtype(var, DEV_TYPE_PATH);
    dbg_syslog(LOG_INFO, "typenum %d", var->dtypenum);

    var->context = redisConnect(REDIS_SERVER_IP, REDIS_SERVER_PORT);
    if (var->context->err)
    {
        dbg_syslog(LOG_ERR, "connect redis server failure: %s\n", var->context->errstr);
        redisFree(var->context);
        return;
    }
    dbg_syslog(LOG_DEBUG, "connect redis server success.");

    var->snapcontext = redisConnect(REDIS_SERVER_IP, REDIS_SERVER_PORT);
    if (var->snapcontext->err)
    {
        dbg_syslog(LOG_ERR, "connect redis server failure: %s\n", var->snapcontext->errstr);
        redisFree(var->snapcontext);
        return;
    }
    dbg_syslog(LOG_DEBUG, "snapcontext connect redis server success.");
    
    var->ems_status_timer = my_timer_create();
    if (var->ems_status_timer > 0)
    {
        my_timer_set(var->ems_status_timer, 1, 2000);
    }

    devtype_entry_t *pos;
    for(pos = var->devtypes; pos != NULL; pos = pos->hh.next)
    {
        dbg_syslog(LOG_DEBUG, "devtype: %s %s", pos->dev_type, pos->dev_type_alias);

#ifdef DB_IEC104
        snprintf(addr_tag_path, ADDR_TAG_PATH_LEN, DEV_TAG"%s.tag", pos->dev_type);
        read_addr_tag(var, pos, addr_tag_path);
#endif

#ifdef DB_MODBUS
        snprintf(addr_tag_path, ADDR_TAG_PATH_LEN, DEV_TAG_MODBUS"%s_MODBUS.tag", pos->dev_type);
        read_addr_tag(var, pos, addr_tag_path);
#endif

        snprintf(addr_tag_path, ADDR_TAG_PATH_LEN, DEV_TAG_EMMS2"%s_EMMS2.tag", pos->dev_type);
        read_emms2_tag(var, pos, addr_tag_path);

        snprintf(addr_tag_path, ADDR_TAG_PATH_LEN, DEV_TAG_SMUALARM"%s_SMUALARM.tag", pos->dev_type);
        read_smualarm_tag(var, pos, addr_tag_path);
    }

    read_alarm_tag(var, DEV_TAG_ALARM);
    parse_nodecfg(var, NODE_CFG_PATH);
    ems_redis_mqtt_client_init(var);
}

int main(int argc, char *argv[])
{
    ems_redis_var_t *var = &gvar;
    memset(var, 0, sizeof(ems_redis_var_t));
    lnxall_loglevel_set(LNXALL_LOGERR, 1);

    ems_redis_init(var);
    ems_redis_loop(var);

    redisFree(var->context);
    dbg_syslog(LOG_DEBUG, "connection with redis server finished.\n");

    return 0;
}