#include <stdlib.h>
#include <stdbool.h>
#include <stdio.h>
#include <string.h>
#include <signal.h>

#include "hal_thread.h"
#include "hal_time.h"
#include "iec104_server_new.h"

static bool running;
#define IEC104_PERIODIC_INTERVAL 30

#define FEATURE_PERIOD_REPORT //allow periodical report
#undef CP56_SPON //spontaneous report with CP56time
#undef CP56_IRG //interrogation with CP56time

#undef KELIN_IEC104
#undef BONA_IEC104
#undef XINYI_IEC104

#if defined(KELIN_IEC104)
uint8_t asduca[] = {18, 49, 51, 54, 55, 60, 61, 62};
#define FEATURE_PERIOD_REPORT
#define CP56_SPON
#define CP56_IRG
#elif defined(BONA_IEC104)
uint8_t asduca[] = {1};
#elif defined(XINYI_IEC104)
uint8_t asduca[] = {2, 18,  49, 50, 51, 52, 53, 56, 57, 59};
#else 
uint8_t asduca[] = {1, 2};
#endif

#ifdef CP56_SPON
#define MAXASDUIOS 14
#else
#define MAXASDUIOS 25
#endif

int asdus = sizeof(asduca);

static void sigint_handler(int sino)
{
    running = false;
    fprintf(stderr, "received signal: %d\n", sino);
    fflush(stderr);
}

static const char * iec104_systime(char * tbuf, size_t tlen)
{
    int ret;
    time_t nowt;
    struct timespec spec;
    struct tm tim, * ptm;

    spec.tv_sec = 0;
    spec.tv_nsec = 0;
    ret = clock_gettime(CLOCK_REALTIME, &spec);
    if (ret == -1) {
        tbuf[0] = '\0';
        return tbuf;
    }

    nowt = spec.tv_sec;
    ptm = localtime_r(&nowt, &tim);

    if (ptm == NULL) {
        ret = snprintf(tbuf, tlen, "%d.%06d",
            (int) spec.tv_sec, (int) (spec.tv_nsec / 1000));
    } else {
        ret = snprintf(tbuf, tlen, "%d-%02d-%02d %02d:%02d:%02d.%06d",
            ptm->tm_year + 1900, ptm->tm_mon + 1, ptm->tm_mday,
            ptm->tm_hour, ptm->tm_min, ptm->tm_sec,
            (int) (spec.tv_nsec / 1000));
    }
    if (ret <= 0) {
        tbuf[0] = '\0';
        return tbuf;
    } else if (ret >= (int) tlen)
        ret = (int) (tlen - 1);
    tbuf[ret] = '\0';
    return tbuf;
}

static void printCP56Time2a(CP56Time2a time)
{
    char nowt[64];
    dbg_syslog(LOG_DEBUG, "[%s]: %02i:%02i:%02i %02i/%02i/%04i\n",
        iec104_systime(nowt, sizeof(nowt)),
        CP56Time2a_getHour(time),
        CP56Time2a_getMinute(time),
        CP56Time2a_getSecond(time),
        CP56Time2a_getDayOfMonth(time),
        CP56Time2a_getMonth(time),
        CP56Time2a_getYear(time) + 2000);
}

static bool clockSyncHandler(void* parameter,
    IMasterConnection connection, CS101_ASDU asdu, CP56Time2a newTime)
{
    char nowt[64];
    dbg_syslog(LOG_DEBUG, "[%s]: Process time sync command with time ",
        iec104_systime(nowt, sizeof(nowt)));
    printCP56Time2a(newTime);
    dbg_syslog(LOG_DEBUG, "\n");

    /* uint64_t newSystemTimeInMs = CP56Time2a_toMsTimestamp(newTime); */

    /* Set time for ACT_CON message */
    CP56Time2a_setFromMsTimestamp(newTime, Hal_getTimeInMs());

    /* update system time here */

    return true;
}

static bool interrogationHandler(void* parameter,
    IMasterConnection connection, CS101_ASDU asdu, uint8_t qoi)
{
    char nowt[64];
    int oaddr;
    int caddr; //common address
    struct iec104_var * var;
    CS101_AppLayerParameters alParams;
    char cmd[CMD_MAX_LENGTH] = {0};
    int i = 0;
    int val_addr, val_type;
    int cur_group = 0;
    int added = 0;
    ushort addrlist[1024 * 10] = {0};
    int addrlen = 0;
    int lasttype = 1;
    int singlegroup = 25;
#ifdef CP56_IRG
    struct sCP56Time2a cp56;
    singlegroup = 15;
    uint64_t dtime;
#endif

    alParams = NULL;
    var = (struct iec104_var *) parameter;

    oaddr = CS101_ASDU_getOA(asdu);
    caddr = CS101_ASDU_getCA(asdu);
    dbg_syslog(LOG_DEBUG, "[%s]: Received interrogation for group %i, OA: %d, CA: %d\n",
        iec104_systime(nowt, sizeof(nowt)), qoi, oaddr, caddr);

    if (qoi != 20)  // only handle station interrogation 
        goto exit;

    alParams = IMasterConnection_getApplicationLayerParameters(connection);
    if (alParams == NULL)
        goto exit;

    IMasterConnection_sendACT_CON(connection, asdu, false);
    CS101_ASDU newAsdu = CS101_ASDU_create(alParams, false, CS101_COT_INTERROGATED_BY_STATION, oaddr, caddr, false, false);
    InformationObject io = NULL;

    snprintf(cmd, CMD_MAX_LENGTH, "zrange %d 0 -1 withscores", caddr);
    redisReply* reply = (redisReply* )redisCommand(var->context, cmd);
    if (NULL == reply || (REDIS_REPLY_NIL == reply->type))
    {
        dbg_syslog(LOG_DEBUG, "%d execcmd %s error\n", __LINE__, cmd);
        freeReplyObject(reply);
        goto exit;
    }
    for (; i < reply->elements; i+=2)
    {
        char *tok = strtok(reply->element[i]->str, "|");
        val_type = atoi(strtok(0, "|"));
        val_addr = atoi(reply->element[i+1]->str);
#ifdef CP56_IRG
        char *dtimestr = strtok(0, "|");
        dtimestr = strtok(0, "|");
        dtime = strtoull(dtimestr, NULL, 10);
#endif
        dbg_syslog(LOG_DEBUG, "Group: %d cur_group %d val_addr %d val_type %d lasttype %d added %d\n", val_addr >> 12, cur_group, val_addr, val_type, lasttype, added);
        if (((val_addr >> 12) != cur_group) || (val_type != lasttype) || (added > singlegroup)) //create new group
        {
            if (added > 0)
            {
                IMasterConnection_sendASDU(connection, newAsdu);
                CS101_ASDU_destroy(newAsdu);
                newAsdu = CS101_ASDU_create(alParams, false, CS101_COT_INTERROGATED_BY_STATION, oaddr, caddr, false, false);
                InformationObject_destroy(io);
                io = NULL;
                added = 0;
            }
            cur_group = val_addr >> 12;
            dbg_syslog(LOG_DEBUG, "create new group\n");
        }
        if (val_type == REDIS_DATA_U8)
        {
            bool val = (bool)atoi(tok);
            if (NULL == io)
            {
#ifdef CP56_IRG
                CP56Time2a_createFromMsTimestamp(&cp56, dtime * 1000);
                printCP56Time2a(&cp56);
                io = (InformationObject)SinglePointWithCP56Time2a_create(NULL, val_addr, val, IEC60870_QUALITY_GOOD, &cp56);
                CS101_ASDU_addInformationObject(newAsdu, io);
#else
                io = (InformationObject)SinglePointInformation_create(NULL, val_addr, val, IEC60870_QUALITY_GOOD);
                CS101_ASDU_addInformationObject(newAsdu, io);
#endif
            }
            else
            {
#ifdef CP56_IRG
                CP56Time2a_createFromMsTimestamp(&cp56, dtime * 1000);
                CS101_ASDU_addInformationObject(newAsdu, (InformationObject)
                    SinglePointWithCP56Time2a_create((SinglePointInformation)io, val_addr, val, IEC60870_QUALITY_GOOD, &cp56));
#else
                CS101_ASDU_addInformationObject(newAsdu, (InformationObject)
                    SinglePointInformation_create((SinglePointInformation)io, val_addr, val, IEC60870_QUALITY_GOOD));
#endif
            }
        }
        else if (val_type == REDIS_DATA_U16 || val_type == REDIS_DATA_I16 || val_type == REDIS_DATA_U32 || val_type == REDIS_DATA_I32)
        {
            int val = atoi(tok);
            if (NULL == io)
            {
#ifdef CP56_IRG 
                CP56Time2a_createFromMsTimestamp(&cp56, dtime * 1000);
                printCP56Time2a(&cp56);
                io = (InformationObject)MeasuredValueScaledWithCP56Time2a_create(NULL, val_addr, val, IEC60870_QUALITY_GOOD, &cp56);
                CS101_ASDU_addInformationObject(newAsdu, io);
#else
                io = (InformationObject)MeasuredValueScaled_create(NULL, val_addr, val, IEC60870_QUALITY_GOOD);
                CS101_ASDU_addInformationObject(newAsdu, io);
#endif
            }
            else
            {
#ifdef CP56_IRG
                CP56Time2a_createFromMsTimestamp(&cp56, dtime * 1000);
                CS101_ASDU_addInformationObject(newAsdu, (InformationObject)
                    MeasuredValueScaledWithCP56Time2a_create((MeasuredValueScaled)io, val_addr, val, IEC60870_QUALITY_GOOD, &cp56));
#else
                CS101_ASDU_addInformationObject(newAsdu, (InformationObject)
                    MeasuredValueScaled_create((MeasuredValueScaled)io, val_addr, val, IEC60870_QUALITY_GOOD));
#endif
            }
        }
        else if (val_type == REDIS_DATA_F32)
        {
            float val = (float)atof(tok);
            if (NULL == io)
            {
#ifdef CP56_IRG
                CP56Time2a_createFromMsTimestamp(&cp56, dtime * 1000);
                printCP56Time2a(&cp56);
                io = (InformationObject)MeasuredValueShortWithCP56Time2a_create(NULL, val_addr, val, IEC60870_QUALITY_GOOD, &cp56);
                CS101_ASDU_addInformationObject(newAsdu, io);
#else
                io = (InformationObject)MeasuredValueShort_create(NULL, val_addr, val, IEC60870_QUALITY_GOOD);
                CS101_ASDU_addInformationObject(newAsdu, io);
#endif
            }
            else
            {
#ifdef CP56_IRG
                CP56Time2a_createFromMsTimestamp(&cp56, dtime * 1000);
                CS101_ASDU_addInformationObject(newAsdu, (InformationObject)
                    MeasuredValueShortWithCP56Time2a_create((MeasuredValueShort)io, val_addr, val, IEC60870_QUALITY_GOOD, &cp56));
#else
                CS101_ASDU_addInformationObject(newAsdu, (InformationObject)
                    MeasuredValueShort_create((MeasuredValueShort)io, val_addr, val, IEC60870_QUALITY_GOOD));
#endif
            }
        }
        added++;
        lasttype = val_type;
    }
    freeReplyObject(reply);
    
    if (added > 0)
        IMasterConnection_sendASDU(connection, newAsdu);
    CS101_ASDU_destroy(newAsdu);
    if (io != NULL)
    {
        InformationObject_destroy(io);
    }

    IMasterConnection_sendACT_TERM(connection, asdu);
    return true;

exit:
    IMasterConnection_sendACT_CON(connection, asdu, true);
    return true;
}

static bool counterinterrogationhandler(void * parameter,
    IMasterConnection connection, CS101_ASDU asdu, QualifierOfCIC qcc)
{
    char nowt[64];
    int oaddr;
    int caddr; //common address
    struct iec104_var * var;
    CS101_AppLayerParameters alParams;
    char cmd[CMD_MAX_LENGTH] = {0};
    int val_addr, val_type;
    int cur_group = 0;
    int added = 0;
    int addrlist[2048] = {0};
    int addrlen = 0;
    int lasttype = 1;

    var = (struct iec104_var *) parameter;
    oaddr = CS101_ASDU_getOA(asdu);
    caddr = CS101_ASDU_getCA(asdu);
    var->lastoaddr = oaddr;
    var->lastcaddr = caddr;

    alParams = IMasterConnection_getApplicationLayerParameters(connection);
    if (alParams == NULL)
        goto exit;
 
    dbg_syslog(LOG_DEBUG, "[%s]: Received interrogation for group , OA: %d, CA: %d\n",
        iec104_systime(nowt, sizeof(nowt)), oaddr, caddr);

    IMasterConnection_sendACT_CON(connection, asdu, false);

    //int cmaddr = 2; //redis addr
    int startaddr = 25601;
    
    CS101_ASDU newAsdu = CS101_ASDU_create(alParams, false, CS101_COT_INTERROGATED_BY_STATION, oaddr, caddr, false, false);
    InformationObject io = NULL;
    
    dbg_syslog(LOG_DEBUG, "addrlen %d addrlist[0] %d\n", addrlen, addrlist[0]);
    while(true)
    {
        val_addr = startaddr++;
        snprintf(cmd, CMD_MAX_LENGTH, "zrangebyscore %d %d %d", caddr, val_addr, val_addr);
        redisReply* reply = (redisReply* )redisCommand(var->context, cmd);
        if(NULL == reply)
        {
            dbg_syslog(LOG_DEBUG, "%d execcmd %s error\n", __LINE__, cmd);
            break;
        }
        if (reply->elements < 1 || reply->element[0]->str == NULL)
        {
            dbg_syslog(LOG_DEBUG, "%d execcmd %s error\n", __LINE__, cmd);
            freeReplyObject(reply);
            break;
        }
        else
        {
            char *tok = strtok(reply->element[0]->str, "|");
            val_type = atoi(strtok(0, "|"));
            dbg_syslog(LOG_DEBUG, "Group: %d cur_group %d val_addr %d val_type %d lasttype %d added %d\n", val_addr >> 12, cur_group, val_addr, val_type, lasttype, added);
            if (((val_addr >> 12) != cur_group) || (val_type != lasttype) || (added > 40)) //create new group
            {
                if (added > 0)
                {
                    IMasterConnection_sendASDU(connection, newAsdu);
                    CS101_ASDU_destroy(newAsdu);
                    newAsdu = CS101_ASDU_create(alParams, false, CS101_COT_INTERROGATED_BY_STATION, oaddr, caddr, false, false);
                    InformationObject_destroy(io);
                    io = NULL;
                    added = 0;
                }
                cur_group = val_addr >> 12;
                dbg_syslog(LOG_DEBUG, "create new group\n");
            }
            int val = atoi(tok);
            BinaryCounterReading bcr1 = BinaryCounterReading_create(NULL, val, 0, false, false, false);
            if (NULL == io)
            {
                io = (InformationObject)IntegratedTotals_create(NULL, val_addr, bcr1);
                CS101_ASDU_addInformationObject(newAsdu, io);
            }
            else
            {
                CS101_ASDU_addInformationObject(newAsdu, (InformationObject)
                    IntegratedTotals_create((IntegratedTotals)io, val_addr, bcr1));
            }
            BinaryCounterReading_destroy(bcr1);   
            added++;
            lasttype = val_type;
            dbg_syslog(LOG_DEBUG, "add io value %d\n", val);
        }
        freeReplyObject(reply);
    }
    if (added > 0)
        IMasterConnection_sendASDU(connection, newAsdu);
    CS101_ASDU_destroy(newAsdu);
    if (io != NULL)
    {
        InformationObject_destroy(io);
    }

    IMasterConnection_sendACT_TERM(connection, asdu);

    return true;

exit:
    IMasterConnection_sendACT_CON(connection, asdu, true);
    return true;
}

static bool asduHandler(void* parameter,
    IMasterConnection connection, CS101_ASDU asdu)
{
    int topublish = 0;
    int caddr; //common address
    char nowt[64];
    char addr[8];
    char value[32];
    char topic[TOPIC_MAX_LEN] = {0};
    struct iec104_var * var;
    IEC60870_5_TypeID typeid;
    typeid = CS101_ASDU_getTypeID(asdu);
    var = (struct iec104_var *) parameter;

    caddr = CS101_ASDU_getCA(asdu);
    iec104_systime(nowt, sizeof(nowt)),
    fprintf(stdout, "[%s]: Received an message, type: %d (%#x)\n",
        nowt, (int) typeid, (unsigned int) typeid);
    fflush(stdout);
    snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/lnxall/device/setems", var->sn_str);
    cJSON *rt_data = cJSON_CreateObject();
    cJSON *pts = cJSON_CreateObject();
    cJSON_AddStringToObject(rt_data, "proto", "iec104");
    cJSON_AddNumberToObject(rt_data, "caddr", caddr);
    cJSON_AddItemToObject(rt_data, "pts", pts);
    fprintf(stdout, "topic : %s %u\n", topic, (typeid == C_SC_NA_1)?1:0);
    fflush(stdout);
    if (typeid == C_SE_NA_1)
    {
        if  (CS101_ASDU_getCOT(asdu) == CS101_COT_ACTIVATION) {
            InformationObject io = CS101_ASDU_getElement(asdu, 0);

            if (io) {
                CS101_ASDU_setCOT(asdu, CS101_COT_ACTIVATION_CON);            
                snprintf(addr, 8, "%d", InformationObject_getObjectAddress(io));
                snprintf(value, 32, "%lf", SetpointCommandNormalized_getValue(io));
                cJSON_AddStringToObject(pts, addr, value);
                topublish = 1;
                InformationObject_destroy(io);
            }
            else {
                dbg_syslog(LOG_DEBUG, "ERROR: message has no valid information object\n");
                return true;
            }
        }
        else
            CS101_ASDU_setCOT(asdu, CS101_COT_UNKNOWN_COT);

        IMasterConnection_sendASDU(connection, asdu);
    }
    else if (typeid == C_SE_NB_1)
    {
        if  (CS101_ASDU_getCOT(asdu) == CS101_COT_ACTIVATION) {
            InformationObject io = CS101_ASDU_getElement(asdu, 0);

            if (io) {
                CS101_ASDU_setCOT(asdu, CS101_COT_ACTIVATION_CON);            
                snprintf(addr, 8, "%d", InformationObject_getObjectAddress(io));
                snprintf(value, 32, "%d", SetpointCommandScaled_getValue(io));
                cJSON_AddStringToObject(pts, addr, value);
                topublish = 1;
                InformationObject_destroy(io);
            }
            else {
                dbg_syslog(LOG_DEBUG, "ERROR: message has no valid information object\n");
                return true;
            }
        }
        else
            CS101_ASDU_setCOT(asdu, CS101_COT_UNKNOWN_COT);

        IMasterConnection_sendASDU(connection, asdu);
    }
    else if (typeid == C_SE_NC_1)
    {
        if  (CS101_ASDU_getCOT(asdu) == CS101_COT_ACTIVATION) {
            InformationObject io = CS101_ASDU_getElement(asdu, 0);

            if (io) {
                CS101_ASDU_setCOT(asdu, CS101_COT_ACTIVATION_CON);            
                snprintf(addr, 8, "%d", InformationObject_getObjectAddress(io));
#ifdef BONA_IEC104
                snprintf(value, 32, "%d", (int)SetpointCommandShort_getValue(io));
#else
                snprintf(value, 32, "%lf", SetpointCommandShort_getValue(io));
#endif
                cJSON_AddStringToObject(pts, addr, value);
                topublish = 1;
                InformationObject_destroy(io);
            }
            else {
                dbg_syslog(LOG_DEBUG, "ERROR: message has no valid information object\n");
                return true;
            }
        }
        else
            CS101_ASDU_setCOT(asdu, CS101_COT_UNKNOWN_COT);

        IMasterConnection_sendASDU(connection, asdu);
    }
    else if (typeid == C_SC_NA_1)
    {
        if  (CS101_ASDU_getCOT(asdu) == CS101_COT_ACTIVATION) {
            InformationObject io = CS101_ASDU_getElement(asdu, 0);

            if (io) {
                CS101_ASDU_setCOT(asdu, CS101_COT_ACTIVATION_CON);            
                snprintf(addr, 8, "%d", InformationObject_getObjectAddress(io));
#ifdef KELIN_IEC104
                snprintf(value, 32, "%d", SingleCommand_getState(io) ? 1 : 2);
#else
                snprintf(value, 32, "%d", SingleCommand_getState(io) ? 1 : 0);
#endif
                cJSON_AddStringToObject(pts, addr, value);
                topublish = 1;
                InformationObject_destroy(io);
                fprintf(stdout, "topublish: %u addr %s value %s", topublish, addr, value);
                fflush(stdout);
            }
            else {
                dbg_syslog(LOG_DEBUG, "ERROR: message has no valid information object\n");
                return true;
            }
        }
        else
            CS101_ASDU_setCOT(asdu, CS101_COT_UNKNOWN_COT);

        IMasterConnection_sendASDU(connection, asdu);
    }
    else if (typeid == C_DC_NA_1)
    {
        if (CS101_ASDU_getCOT(asdu) == CS101_COT_ACTIVATION) {
            InformationObject io = CS101_ASDU_getElement(asdu, 0);
            if (io) {
                int dcstate = DoubleCommand_getState(io);
                if ((dcstate > 2) || (dcstate < 1))
                {
                    dbg_syslog(LOG_ERR, "invalid dcstate %d", dcstate);
                }
                else 
                {
                    CS101_ASDU_setCOT(asdu, CS101_COT_ACTIVATION_CON);
                    snprintf(addr, 8, "%d", InformationObject_getObjectAddress(io));
                    snprintf(value, 4, "%d", dcstate - 1);
                    cJSON_AddStringToObject(pts, addr, value);
                    topublish = 1;
                    fprintf(stdout, "topublish: %u addr %s value %s\n", topublish, addr, value);
                    fflush(stdout);
                }
                InformationObject_destroy(io);
            }
            else {
                dbg_syslog(LOG_DEBUG, "ERROR: message has no valid information object\n");
                return true;
            }
        }
        else
            CS101_ASDU_setCOT(asdu, CS101_COT_UNKNOWN_COT);
        
        IMasterConnection_sendASDU(connection, asdu);
    }
    else
    {
        //do nothing.
    }
    
    if (topublish)
    {
        char *data_tmp = cJSON_Print(rt_data);
        fprintf(stdout, "data_tmp : %s %u\n", data_tmp, strlen(data_tmp));
        fflush(stdout);
        ipc_session_publish(var->session, topic, data_tmp, strlen(data_tmp));
        free(data_tmp);
    }

    fflush(stdout);
    cJSON_Delete(rt_data);
    return true;
}

static bool connectionRequestHandler(void* parameter, const char* ipAddress)
{
    char nowt[64];
    fprintf(stdout, "[%s]: New connection request from %s\n",
        iec104_systime(nowt, sizeof(nowt)), ipAddress);
    fflush(stdout);
    return true;
}

static void connectionEventHandler(void* parameter,
    IMasterConnection con, CS104_PeerConnectionEvent event)
{
    char nowt[64];
    iec104_systime(nowt, sizeof(nowt));
    if (event == CS104_CON_EVENT_CONNECTION_OPENED) {
        dbg_syslog(LOG_DEBUG, "[%s]: Connection opened (%p)\n", nowt, con);
    }
    else if (event == CS104_CON_EVENT_CONNECTION_CLOSED) {
        dbg_syslog(LOG_DEBUG, "[%s]: Connection closed (%p)\n", nowt, con);
    }
    else if (event == CS104_CON_EVENT_ACTIVATED) {
        dbg_syslog(LOG_DEBUG, "[%s]: Connection activated (%p)\n", nowt, con);
    }
    else if (event == CS104_CON_EVENT_DEACTIVATED) {
        dbg_syslog(LOG_DEBUG, "[%s]: Connection deactivated (%p)\n", nowt, con);
    } else {
        dbg_syslog(LOG_DEBUG, "[%s]: unknown event: %d (%#x)\n", nowt, (int) event, (unsigned int) event);
    }
    fflush(stdout);
}

static void periodReport(int caddr, struct iec104_var *var, CS104_Slave slave)
{
    char cmd[CMD_MAX_LENGTH] = {0};
    int i = 0;
    int val_addr, val_type;
    int cur_group = 0;
    int added = 0;
    ushort addrlist[2048 * 8] = {0};
    int addrlen = 0;
    int lasttype = 1;
    int singlegroup = 25;
    CS101_AppLayerParameters alParams = CS104_Slave_getAppLayerParameters(var->slave);

    snprintf(cmd, CMD_MAX_LENGTH, "zrange %d 0 -1 withscores", caddr);
    redisReply* reply = (redisReply* )redisCommand(var->context, cmd);
    if (NULL == reply || (REDIS_REPLY_NIL == reply->type))
    {
        dbg_syslog(LOG_DEBUG, "%d execcmd %s error\n", __LINE__, cmd);
        freeReplyObject(reply);
        return false;
    }

    CS101_ASDU newAsdu = CS101_ASDU_create(alParams, false, CS101_COT_PERIODIC, var->lastoaddr, caddr, false, false);
    InformationObject io = NULL;
    for (; i < reply->elements; i+=2)
    {
        char *tok = strtok(reply->element[i]->str, "|");
        val_type = atoi(strtok(0, "|"));
        val_addr = atoi(reply->element[i+1]->str);
        if (val_addr > 25600) //counter value
            continue;
#ifdef CP56_IRG
        char *dtimestr = strtok(0, "|");
        dtimestr = strtok(0, "|");
        dtime = strtoull(dtimestr, NULL, 10);
#endif
        dbg_syslog(LOG_DEBUG, "Group: %d cur_group %d val_addr %d val_type %d lasttype %d added %d\n", val_addr >> 12, cur_group, val_addr, val_type, lasttype, added);
        if (((val_addr >> 12) != cur_group) || (val_type != lasttype) || (added > singlegroup)) //create new group
        {
            if (added > 0)
            {
                CS104_Slave_enqueueASDU(var->slave, newAsdu);
                CS101_ASDU_destroy(newAsdu);
                newAsdu = CS101_ASDU_create(alParams, false, CS101_COT_PERIODIC, var->lastoaddr, caddr, false, false);
                InformationObject_destroy(io);
                io = NULL;
                added = 0;
            }
            cur_group = val_addr >> 12;
            dbg_syslog(LOG_DEBUG, "create new group\n");
        }
        if (val_type == REDIS_DATA_U8)
        {
            bool val = (bool)atoi(tok);
            if (NULL == io)
            {
#ifdef CP56_IRG
                CP56Time2a_createFromMsTimestamp(&cp56, dtime * 1000);
                printCP56Time2a(&cp56);
                io = (InformationObject)SinglePointWithCP56Time2a_create(NULL, val_addr, val, IEC60870_QUALITY_GOOD, &cp56);
                CS101_ASDU_addInformationObject(newAsdu, io);
#else
                io = (InformationObject)SinglePointInformation_create(NULL, val_addr, val, IEC60870_QUALITY_GOOD);
                CS101_ASDU_addInformationObject(newAsdu, io);
#endif
            }
            else
            {
#ifdef CP56_IRG
                CP56Time2a_createFromMsTimestamp(&cp56, dtime * 1000);
                CS101_ASDU_addInformationObject(newAsdu, (InformationObject)
                    SinglePointWithCP56Time2a_create((SinglePointInformation)io, val_addr, val, IEC60870_QUALITY_GOOD, &cp56));
#else
                CS101_ASDU_addInformationObject(newAsdu, (InformationObject)
                    SinglePointInformation_create((SinglePointInformation)io, val_addr, val, IEC60870_QUALITY_GOOD));
#endif
            }
        }
        else if (val_type == REDIS_DATA_U16 || val_type == REDIS_DATA_I16 || val_type == REDIS_DATA_U32 || val_type == REDIS_DATA_I32)
        {
            int val = atoi(tok);
            if (NULL == io)
            {
#ifdef CP56_IRG 
                CP56Time2a_createFromMsTimestamp(&cp56, dtime * 1000);
                printCP56Time2a(&cp56);
                io = (InformationObject)MeasuredValueScaledWithCP56Time2a_create(NULL, val_addr, val, IEC60870_QUALITY_GOOD, &cp56);
                CS101_ASDU_addInformationObject(newAsdu, io);
#else
                io = (InformationObject)MeasuredValueScaled_create(NULL, val_addr, val, IEC60870_QUALITY_GOOD);
                CS101_ASDU_addInformationObject(newAsdu, io);
#endif
            }
            else
            {
#ifdef CP56_IRG
                CP56Time2a_createFromMsTimestamp(&cp56, dtime * 1000);
                CS101_ASDU_addInformationObject(newAsdu, (InformationObject)
                    MeasuredValueScaledWithCP56Time2a_create((MeasuredValueScaled)io, val_addr, val, IEC60870_QUALITY_GOOD, &cp56));
#else
                CS101_ASDU_addInformationObject(newAsdu, (InformationObject)
                    MeasuredValueScaled_create((MeasuredValueScaled)io, val_addr, val, IEC60870_QUALITY_GOOD));
#endif
            }
        }
        else if (val_type == REDIS_DATA_F32)
        {
            float val = (float)atof(tok);
            if (NULL == io)
            {
#ifdef CP56_IRG
                CP56Time2a_createFromMsTimestamp(&cp56, dtime * 1000);
                printCP56Time2a(&cp56);
                io = (InformationObject)MeasuredValueShortWithCP56Time2a_create(NULL, val_addr, val, IEC60870_QUALITY_GOOD, &cp56);
                CS101_ASDU_addInformationObject(newAsdu, io);
#else
                io = (InformationObject)MeasuredValueShort_create(NULL, val_addr, val, IEC60870_QUALITY_GOOD);
                CS101_ASDU_addInformationObject(newAsdu, io);
#endif
            }
            else
            {
#ifdef CP56_IRG
                CP56Time2a_createFromMsTimestamp(&cp56, dtime * 1000);
                CS101_ASDU_addInformationObject(newAsdu, (InformationObject)
                    MeasuredValueShortWithCP56Time2a_create((MeasuredValueShort)io, val_addr, val, IEC60870_QUALITY_GOOD, &cp56));
#else
                CS101_ASDU_addInformationObject(newAsdu, (InformationObject)
                    MeasuredValueShort_create((MeasuredValueShort)io, val_addr, val, IEC60870_QUALITY_GOOD));
#endif
            }
        }
        added++;
        lasttype = val_type;
    }
    
    if (added > 0)
        CS104_Slave_enqueueASDU(var->slave, newAsdu);
    CS101_ASDU_destroy(newAsdu);
    if (io != NULL)
    {
        InformationObject_destroy(io);
    }
    freeReplyObject(reply);
    return;
}

static int iec104_server_new_mqtt_handle_recv_msg(void *obj, ipc_msg_t *mqtt_msg)
{
    struct iec104_var *var = (struct iec104_var *)obj;
    CS101_ASDU newAsdu = NULL;
    int caddr = 0;
    InformationObject io = NULL;
    CS101_AppLayerParameters alParams = CS104_Slave_getAppLayerParameters(var->slave);
    int val_addr;
    int val_int;
    float val_float;
#ifdef CP56_SPON
    struct sCP56Time2a cp56;
#endif
    
    cJSON *root = cJSON_Parse(mqtt_msg->payload);
    if (!root)
    {
        dbg_syslog(LOG_ERR, "parse setems failed.");
        return -1;
    }

    GET_JSON_VALUE_INT(root, "ca", caddr);
    cJSON *yx = cJSON_GetObjectItem(root, "yx");
    if (!yx)
    {
        dbg_syslog(LOG_ERR, "get yx array failed.");
    }
    else
    {
        int yxcount = cJSON_GetArraySize(yx);
        int addedios = 0;
        if (yxcount > 0){
            newAsdu = CS101_ASDU_create(alParams, false,
                    CS101_COT_SPONTANEOUS, var->lastoaddr, caddr > 0 ? caddr : (var->lastcaddr), false, false);
            for (int i = 0; i < yxcount; i++)
            {
                cJSON *node = cJSON_GetArrayItem(yx, i);
                val_addr = atoi(node->string);
                val_int =  node->valueint;
                if (NULL == io)
                {
#ifdef CP56_SPON
                    CP56Time2a_createFromMsTimestamp(&cp56, Hal_getTimeInMs());
                    io = (InformationObject)SinglePointWithCP56Time2a_create(NULL, val_addr, val_int, IEC60870_QUALITY_GOOD, &cp56);
                    CS101_ASDU_addInformationObject(newAsdu, io);
#else
                    io = (InformationObject)SinglePointInformation_create(NULL, val_addr, val_int, IEC60870_QUALITY_GOOD);
                    CS101_ASDU_addInformationObject(newAsdu, io);
#endif
                }
                else
                {
#ifdef CP56_SPON
                    CP56Time2a_createFromMsTimestamp(&cp56, Hal_getTimeInMs());
                    CS101_ASDU_addInformationObject(newAsdu, (InformationObject)
                        SinglePointWithCP56Time2a_create((SinglePointInformation)io, val_addr, val_int, IEC60870_QUALITY_GOOD, &cp56));
#else
                    CS101_ASDU_addInformationObject(newAsdu, (InformationObject)
                        SinglePointInformation_create((SinglePointInformation)io, val_addr, val_int, IEC60870_QUALITY_GOOD));
#endif
                }
                addedios++;
                if (addedios >= MAXASDUIOS)
                {
                    dbg_syslog(LOG_INFO, "enqueue YX spoontaneous report with yxcount %d", addedios);
                    CS104_Slave_enqueueASDU(var->slave, newAsdu);
                    CS101_ASDU_destroy(newAsdu);
                    InformationObject_destroy(io);
                    io = NULL;
                    newAsdu = CS101_ASDU_create(alParams, false,
                        CS101_COT_SPONTANEOUS, var->lastoaddr, caddr > 0 ? caddr : (var->lastcaddr), false, false);
                    addedios = 0;
                }
            }
            if (addedios > 0)
            {
                dbg_syslog(LOG_INFO, "enqueue YX spoontaneous report with yxcount %d", addedios);
                CS104_Slave_enqueueASDU(var->slave, newAsdu);
                InformationObject_destroy(io);
                io = NULL;
            }
            CS101_ASDU_destroy(newAsdu);
        }
    }

    cJSON *yc = cJSON_GetObjectItem(root, "yc");
    if (!yc)
    {
        dbg_syslog(LOG_ERR, "get yc array failed.");
    }
    else
    {
        int yccount = cJSON_GetArraySize(yc);
        int addedios = 0;
        if (yccount > 0)
        {
            newAsdu = CS101_ASDU_create(alParams, false,
                    CS101_COT_SPONTANEOUS, var->lastoaddr, caddr > 0 ? caddr : (var->lastcaddr), false, false);
            for (int i = 0; i < yccount; i++)
            {
                cJSON *node = cJSON_GetArrayItem(yc, i);
                val_addr = atoi(node->string);
                val_float =  atof(node->valuestring);
                if (NULL == io)
                {
#ifdef CP56_SPON 
                    CP56Time2a_createFromMsTimestamp(&cp56, Hal_getTimeInMs());
                    io = (InformationObject)MeasuredValueShortWithCP56Time2a_create(NULL, val_addr, val_float, IEC60870_QUALITY_GOOD, &cp56);
                    CS101_ASDU_addInformationObject(newAsdu, io);
#else
                    io = (InformationObject)MeasuredValueShort_create(NULL, val_addr, val_float, IEC60870_QUALITY_GOOD);
                    CS101_ASDU_addInformationObject(newAsdu, io);
#endif
                }
                else
                {
#ifdef CP56_SPON
                    CP56Time2a_createFromMsTimestamp(&cp56, Hal_getTimeInMs());
                    CS101_ASDU_addInformationObject(newAsdu, (InformationObject)
                        MeasuredValueShortWithCP56Time2a_create((MeasuredValueShort)io, val_addr, val_float, IEC60870_QUALITY_GOOD, &cp56));
#else
                    CS101_ASDU_addInformationObject(newAsdu, (InformationObject)
                        MeasuredValueShort_create((MeasuredValueShort)io, val_addr, val_float, IEC60870_QUALITY_GOOD));
#endif
                }
                addedios++;
                if (addedios >= MAXASDUIOS)
                {
                    dbg_syslog(LOG_INFO, "enqueue YC spontaneous report with yccount %d", addedios);
                    CS104_Slave_enqueueASDU(var->slave, newAsdu);
                    CS101_ASDU_destroy(newAsdu);
                    InformationObject_destroy(io);
                    io = NULL;
                    newAsdu = CS101_ASDU_create(alParams, false,
                        CS101_COT_SPONTANEOUS, var->lastoaddr, caddr > 0 ? caddr : (var->lastcaddr), false, false);
                    addedios = 0;
                }
            }
            if (addedios > 0)
            {
                dbg_syslog(LOG_INFO, "enqueue YC spontaneous report with yccount %d", addedios);
                CS104_Slave_enqueueASDU(var->slave, newAsdu);
                InformationObject_destroy(io);
                io = NULL;
            }
            CS101_ASDU_destroy(newAsdu);
        }
    }
    
    cJSON_Delete(root);
    return;
}

static void iec104_server_new_mqtt_init(struct iec104_var *var)
{
    char clientId[MAX_CLIENT_ID_LEN] = {0};
    char topic[TOPIC_MAX_LEN] = {0};
    snprintf(clientId, MAX_CLIENT_ID_LEN, "INIT_IEC104_SERVER_NEW");
    var->session = ipc_session_new(clientId, (void*)var, IPC_DEFAULT);
    if (var->session == NULL)
        return -1;

    ipc_session_set_callbacks(var->session, iec104_server_new_mqtt_handle_recv_msg, NULL);

    snprintf(topic, TOPIC_MAX_LEN, "ipc/+/lnxall/device/changeReport");
    ipc_session_subscribe(var->session, topic);

    ipc_session_start(var->session);
    return 0;
}

int main(int argc, char** argv)
{
    int ret;
    unsigned int lasttime;
    struct iec104_var * var;

    running = true;
    /* Add Ctrl-C handler */
    signal(SIGINT, sigint_handler);

    var = iec104_server_new_init();
    if (var == NULL)
        return 1;

    var->context = redisConnect(REDIS_SERVER_IP, REDIS_SERVER_PORT);
    if (var->context->err)
    {
        redisFree(var->context);
        dbg_syslog(LOG_DEBUG, "%d connect redis server failed: %s\n", __LINE__, var->context->errstr);
        return false;
    }
    dbg_syslog(LOG_DEBUG, "connect redis server success.\n");
    
    iec104_server_new_mqtt_init(var);
    get_board_sn(var->sn_str);

    CS104_APCIParameters apciParams = CS104_Slave_getdefault_ConnectionParameters();
    apciParams->k = 320;
    /* create a new slave/server instance with default connection parameters and
     * default message queue size */
    var->slave = CS104_Slave_create(100, 100);

    CS104_Slave_setLocalAddress(var->slave, "0.0.0.0");

    /* set iec104 server listener port */
    CS104_Slave_setLocalPort(var->slave, 2404);
    
    /* Set mode to a single redundancy group
     * NOTE: library has to be compiled with CONFIG_CS104_SUPPORT_SERVER_MODE_SINGLE_REDUNDANCY_GROUP enabled (=1)
     */
    CS104_Slave_setServerMode(var->slave, CS104_MODE_MULTIPLE_REDUNDANCY_GROUPS);

    /* get the connection parameters - we need them to create correct ASDUs -
     * you can also modify the parameters here when default parameters are not to be used */
    /* CS101_AppLayerParameters alParams = CS104_Slave_getAppLayerParameters(slave); */

    /* when you have to tweak the APCI parameters (t0-t3, k, w) you can access them here */
    apciParams = CS104_Slave_getConnectionParameters(var->slave);

    dbg_syslog(LOG_DEBUG, "APCI parameters:\n");
    dbg_syslog(LOG_DEBUG, "  t0: %i\n", apciParams->t0);
    dbg_syslog(LOG_DEBUG, "  t1: %i\n", apciParams->t1);
    dbg_syslog(LOG_DEBUG, "  t2: %i\n", apciParams->t2);
    dbg_syslog(LOG_DEBUG, "  t3: %i\n", apciParams->t3);
    dbg_syslog(LOG_DEBUG, "  k: %i\n", apciParams->k);
    dbg_syslog(LOG_DEBUG, "  w: %i\n", apciParams->w);

    /* set the callback handler for the clock synchronization command */
    CS104_Slave_setClockSyncHandler(var->slave, clockSyncHandler, NULL);

    /* set the callback handler for the interrogation command */
    CS104_Slave_setInterrogationHandler(var->slave, interrogationHandler, var);

    /* set the callback handler for the counter interrogation command */
    CS104_Slave_setCounterInterrogationHandler(var->slave, counterinterrogationhandler, var);

    /* set handler for other message types */
    CS104_Slave_setASDUHandler(var->slave, asduHandler, var);

    /* set handler to handle connection requests (optional) */
    CS104_Slave_setConnectionRequestHandler(var->slave, connectionRequestHandler, NULL);

    /* set handler to track connection events (optional) */
    CS104_Slave_setConnectionEventHandler(var->slave, connectionEventHandler, NULL);

    /* uncomment to log messages */
    /* CS104_Slave_setRawMessageHandler(slave, rawMessageHandler, NULL); */

    CS104_Slave_start(var->slave);
    if (!CS104_Slave_isRunning(var->slave)) {
        dbg_syslog(LOG_DEBUG, "Starting server failed!\n");
        goto exit_program;
    }

    /* acquire mutex lock */
    ret = pthread_mutex_lock(&var->thread_lock);
    if (ret != 0) {
        fprintf(stderr, "Error, failed to acquire mutex: %d\n", ret);
        fflush(stderr);
        goto exit1_program;
    }

    lasttime = iec104_uptime(NULL, NULL);
    /* main loop */
    while (running) {
        unsigned int nowt;

        nowt = iec104_uptime(NULL, NULL);
        if (nowt >= lasttime) {
            if (CS104_Slave_getOpenConnections(var->slave) > 0)
            {
#ifdef FEATURE_PERIOD_REPORT
                for(int iasdu = 0; iasdu < asdus; iasdu++)
                    periodReport(asduca[iasdu], var, var->slave);
#endif
            }
            lasttime += IEC104_PERIODIC_INTERVAL;
            if (nowt >= lasttime)
                lasttime = nowt + IEC104_PERIODIC_INTERVAL;
        }
        ret = lnxall_condvar_timedwait(&var->thread_cond,
            &var->thread_lock, IEC104_PERIODIC_INTERVAL * 1000);
        if (ret == ETIMEDOUT)
            continue;
        if (ret != 0) {
            fprintf(stderr, "Error, cond_timedwait has failed: %d\n", ret);
            fflush(stderr);
            break;
        }
    }

    ret = pthread_mutex_unlock(&var->thread_lock);
    if (ret != 0) {
        fprintf(stderr, "Error, failed to release mutex: %d\n", ret);
        fflush(stderr);
    }
exit1_program:
    CS104_Slave_stop(var->slave);
    redisFree(var->context);

exit_program:
    CS104_Slave_destroy(var->slave);
    redisFree(var->context);
    free(var);
    Thread_sleep(500);
    return 2;
}
