/* This work is licensed under a Creative Commons CCZero 1.0 Universal License.
 * See http://creativecommons.org/publicdomain/zero/1.0/ for more information. */

/*#include <open62541/client_config_default.h>
#include <open62541/client_highlevel.h>
#include <open62541/client_subscriptions.h>
#include <open62541/plugin/log_stdout.h>
*/
#include <stdio.h>
#include <stdlib.h>
#include "ua_mqttv1.h"
#include <stdbool.h>
#include <string.h>
#include <signal.h>

ua_mqttv1_var_t gvar = {0};

static void ua_mqttv1_connredis(redisContext **connContext)
{
    while(true)
    {
        *connContext = redisConnect(REDIS_SERVER_IP, REDIS_SERVER_PORT);
        if ((*connContext)->err)
        {
            dbg_syslog(LOG_ERR, "connect redis server failure: %s\n", (*connContext)->errstr);
            redisFree(*connContext);
            usleep(10 * 1000 * 1000); //wait 10s for redis server start
        }
        else
        {
            dbg_syslog(LOG_NOTICE, "connect redis server success");
            break;
        }
    }
}

static void ua_mqttv1_read_devtype(ua_mqttv1_var_t *var, char *path)
{
    char typenam[MAX_DEVTYPE_LEN] = {0};
    int addr;

    FILE *fp = fopen(path, "r");
    if (NULL == fp)
    {
        dbg_syslog(LOG_ERR, "fopen %s error", path);
        return;
    }
    while(fscanf(fp, "%s %d", typenam, &addr) != EOF)
    {
        devtype_t *curdev;
        HASH_FIND_STR(var->devtypes, typenam, curdev);
        if (NULL == curdev)
        {
            curdev = malloc(sizeof(devtype_t));
            snprintf(curdev->typenam, MAX_DEVTYPE_LEN, "%s", typenam);
            curdev->addr = addr;
            HASH_ADD_STR(var->devtypes, typenam, curdev);
        }
    }
    fclose(fp);
}
/*
static UA_INLINE UA_Boolean
UA_StatusCode_isEqualTop (UA_StatusCode s1, UA_StatusCode s2) {
  return ((s1 & 0xFFFF0000) == (s2 & 0xFFFF0000));
}
*/
static UA_ByteString loadFile(const char *const path) {
    UA_ByteString fileContents = UA_STRING_NULL;

    /* Open the file */
    FILE *fp = fopen(path, "rb");
    if(!fp) {
        errno = 0; /* We read errno also from the tcp layer... */
        return fileContents;
    }

    /* Get the file length, allocate the data and read */
    fseek(fp, 0, SEEK_END);
    fileContents.length = (size_t)ftell(fp);
    fileContents.data = (UA_Byte *)UA_malloc(fileContents.length * sizeof(UA_Byte));
    if(fileContents.data) {
        fseek(fp, 0, SEEK_SET);
        size_t read = fread(fileContents.data, sizeof(UA_Byte), fileContents.length, fp);
        if(read != fileContents.length)
            UA_ByteString_clear(&fileContents);
    } else {
        fileContents.length = 0;
    }
    fclose(fp);

    return fileContents;
}

/*
static UA_StatusCode nodeIter(UA_NodeId childId, UA_Boolean isInverse, UA_NodeId referenceTypeId, void *handle) 
{
    if(isInverse)
        return UA_STATUSCODE_GOOD;
    UA_NodeId *parent = (UA_NodeId *)handle;
    dbg_syslog(LOG_ERR, "%u, %u --- %u ---> NodeId %u, %u\n",
           parent->namespaceIndex, parent->identifier.numeric,
           referenceTypeId.identifier.numeric, childId.namespaceIndex,
           childId.identifier.numeric);
    return UA_STATUSCODE_GOOD;
}
*/

static uint UA_Client_ReadSingleId(ua_mqttv1_var_t *var, UA_NodeId nid, void *result)
{
    uint ret = 0;
    UA_ReadValueId rid;
    UA_ReadValueId_init(&rid);
    rid.nodeId = nid;
    rid.attributeId = UA_ATTRIBUTEID_VALUE;
    UA_ReadRequest request;
    UA_ReadRequest_init(&request);
    request.nodesToRead = &rid;
    request.nodesToReadSize = 1;
    UA_ReadResponse response = UA_Client_Service_read(var->client, request);

    UA_StatusCode retval = response.responseHeader.serviceResult;
    if(retval == UA_STATUSCODE_GOOD) {
        UA_DataValue res = response.results[0];
        if (res.hasValue)
        {
            if (res.value.type == &UA_TYPES[UA_TYPES_STRING])
            {
                UA_String svalue = *(UA_String *)res.value.data;
                strncpy((char*)result, svalue.data, (int)svalue.length);
                ret = 1;
            }
        }
    }

    UA_ReadResponse_clear(&response);
    return ret;
}

static void UA_Client_WriteSingleNode(ua_mqttv1_var_t *var, char *devsn, char *datatag, void* value)
{
    bool found = false;
    for (int i = 0; i < var->dev_table.devcount; i++)
    {
        if (strncmp(var->dev_table.devtab[i].devsn, devsn, strlen(var->dev_table.devtab[i].devsn)) == 0)
        {
            found = true;
        }
    }
    if(found == false)
    {
        dbg_syslog(LOG_ERR, "device %s not found", devsn);
        return;
    }
    UA_WriteRequest wReq;
    UA_WriteRequest_init(&wReq);
    wReq.nodesToWrite = UA_WriteValue_new();
    wReq.nodesToWriteSize = 1;
    wReq.nodesToWrite[0].nodeId = UA_NODEID_STRING_ALLOC(1, "the.answer");
    wReq.nodesToWrite[0].attributeId = UA_ATTRIBUTEID_VALUE;
    wReq.nodesToWrite[0].value.hasValue = true;
    wReq.nodesToWrite[0].value.value.type = &UA_TYPES[UA_TYPES_INT32];
    wReq.nodesToWrite[0].value.value.storageType = UA_VARIANT_DATA_NODELETE; //do not free the integer on deletion 
    wReq.nodesToWrite[0].value.value.data = &value;
    UA_WriteResponse wResp = UA_Client_Service_write(var->client, wReq);
    if(wResp.responseHeader.serviceResult == UA_STATUSCODE_GOOD)
            printf("the new value is: %p\n", value);
    UA_WriteRequest_clear(&wReq);
    UA_WriteResponse_clear(&wResp);
}

static void UA_Client_ReadMultiValueAttribute(ua_mqttv1_var_t *var, devserv_t *dserv)
{
    char cmd[CMD_MAX_LENGTH] = {0};
    time_t now;

    if (dserv->datanum <= 0)
    {
        dbg_syslog(LOG_ERR, "service no data to read");
        return;
    }
    UA_ReadValueId *item = malloc(sizeof(UA_ReadValueId) * (dserv->datanum));
    devdata_t *pos;
    int tagcount = 0;
    for(pos = dserv->devdata; pos != NULL; pos = pos->hh.next)
    {
        UA_ReadValueId_init(&item[tagcount]);
        item[tagcount].nodeId = pos->nodeId;
        item[tagcount].attributeId = UA_ATTRIBUTEID_VALUE;
        tagcount++;
    }
    
    UA_ReadRequest request;
    UA_ReadRequest_init(&request);
    request.nodesToRead = item;
    request.nodesToReadSize = dserv->datanum;
    UA_ReadResponse response = UA_Client_Service_read(var->client, request);
    UA_StatusCode retval = response.responseHeader.serviceResult;
    now = time(NULL);

    if(retval == UA_STATUSCODE_GOOD) {
        if(response.resultsSize != tagcount)
        {
            dbg_syslog(LOG_ERR, "resultsize %d tagcount %d", (int)response.resultsSize, tagcount);
            UA_ReadResponse_clear(&response);
            return;
        }
        int i = 0;
        for(pos = dserv->devdata; pos != NULL; pos = pos->hh.next)
        {
            UA_DataValue res = response.results[i++];
            if (res.hasValue && UA_Variant_isScalar(&res.value))
            {  
                if (pos->addr == 0)
                {
                    if (strncmp(dserv->service, SERV_ALARM, strlen(SERV_ALARM)) == 0)
                    {
                        pos->dtype = REDIS_DATA_U8;
                        if (var->proto == 1)
                            pos->addr = dserv->dev->yxaddr++;
                        else 
                            pos->addr = dserv->dev->modbus_addr++;
                    }
                    else
                    {
                        if (var->proto == 1)
                            pos->addr = dserv->dev->ycaddr++;
                        if (res.value.type == &UA_TYPES[UA_TYPES_INT32])
                        {
                            pos->dtype = REDIS_DATA_I32;
                            if (var->proto == 2)
                            {
                                pos->addr = dserv->dev->modbus_addr;
                                dserv->dev->modbus_addr += 2;
                            }
                        }
                        else if (res.value.type == &UA_TYPES[UA_TYPES_DOUBLE])
                        {
                            pos->dtype = REDIS_DATA_F32;
                            if (var->proto == 2)
                            {
                                pos->addr = dserv->dev->modbus_addr;
                                dserv->dev->modbus_addr += 2;
                            }
                        }
                    }
                }
                else 
                {
                    if (var->proto == 1)
                        snprintf(cmd, CMD_MAX_LENGTH, "zremrangebyscore %d %u %u", dserv->dev->common_addr, pos->addr, pos->addr);
                    else 
                        snprintf(cmd, CMD_MAX_LENGTH, "zremrangebyscore %d_modbus %u %u", dserv->dev->common_addr, pos->addr, pos->addr);
                    
                    redisReply *reply_del = (redisReply *)redisCommand(var->redisContext, cmd);
                    freeReplyObject(reply_del);
                }
                if (res.value.type == &UA_TYPES[UA_TYPES_INT32])
                {
                    UA_Int32 mvalue = *(UA_Int32 *)res.value.data;
                    dbg_syslog(LOG_DEBUG, "the value is %i", mvalue);
                    if (var->proto == 1)
                        snprintf(cmd, CMD_MAX_LENGTH, "zadd %d %d %d|%d|%s|%lu", dserv->dev->common_addr, pos->addr, mvalue, pos->dtype, pos->datatag, now);
                    else 
                        snprintf(cmd, CMD_MAX_LENGTH, "zadd %d_modbus %d %d|%d|%s|%lu", dserv->dev->common_addr, pos->addr, mvalue, pos->dtype, pos->datatag, now);
                }
                else if (res.value.type == &UA_TYPES[UA_TYPES_DOUBLE])
                {
                    UA_Double dvalue = *(UA_Double *)res.value.data;
                    dbg_syslog(LOG_ERR, "the value is %lf", dvalue);
                    if (var->proto == 1)
                        snprintf(cmd, CMD_MAX_LENGTH, "zadd %d %d %lf|%d|%s|%lu", dserv->dev->common_addr, pos->addr, dvalue, pos->dtype, pos->datatag, now);
                    else 
                        snprintf(cmd, CMD_MAX_LENGTH, "zadd %d_modbus %d %lf|%d|%s|%lu", dserv->dev->common_addr, pos->addr, dvalue, pos->dtype, pos->datatag, now);
                }
                redisReply *reply = (redisReply *)redisCommand(var->redisContext, cmd);
                freeReplyObject(reply);
            }
        }
    }

    UA_ReadResponse_clear(&response);
    free(item);
    return;
}

static int node_browse(ua_mqttv1_var_t *var, UA_NodeId browseId, UA_NodeId* retNodes, int layer, void *aptr)
{   
    int nodecnt = 0;
    dev_entry_t *dentry = NULL;
    devserv_t *dserv = NULL;

    if (UA_NodeId_isNull(&browseId))
    {
        dbg_syslog(LOG_ERR, "browseId is null, return");
        return 0;
    }

    if (1 == layer)
    {
        dentry = (dev_entry_t *)aptr;
    }
    else if (2 == layer)
    {
        dserv = (devserv_t *)aptr;
    }
    dbg_syslog(LOG_ERR, "browsing node");
    UA_BrowseRequest bReq;
    UA_BrowseRequest_init(&bReq);
    bReq.requestedMaxReferencesPerNode = 0;
    bReq.nodesToBrowse = UA_BrowseDescription_new();
    bReq.nodesToBrowseSize = 1;
    bReq.nodesToBrowse[0].nodeId = browseId; /* browse objects folder */
    bReq.nodesToBrowse[0].resultMask = UA_BROWSERESULTMASK_ALL; /* return everything */
    dbg_syslog(LOG_ERR, "browsing node2");
    UA_BrowseResponse bResp = UA_Client_Service_browse(var->client, bReq);
    dbg_syslog(LOG_ERR, "bResp.resultsSize %d", (int)bResp.resultsSize);
    //dbg_syslog(LOG_ERR, "layer %d %-9s %-16s %-16s %-16s\n", layer, ":", "NODEID", "BROWSE NAME", "DISPLAY NAME");
    for(size_t i = 0; i < bResp.resultsSize; ++i) {
        for(size_t j = 0; j < bResp.results[i].referencesSize; ++j) {
            UA_ReferenceDescription *ref = &(bResp.results[i].references[j]);
            if (0 == layer)
                UA_NodeId_copy(&ref->nodeId.nodeId, retNodes+(nodecnt++));
            else if (1 == layer)
            {
                if (strncmp(SERV_STATUS, ref->browseName.name.data, strlen(SERV_STATUS)) == 0)
                {
                    UA_NodeId_copy(&ref->nodeId.nodeId, &dentry->statserv.nodeId);
                    strncpy(dentry->statserv.service, SERV_STATUS, strlen(SERV_STATUS));
                    dentry->statserv.dev = dentry;
                }
                else if (strncmp(SERV_DATA, ref->browseName.name.data, strlen(SERV_DATA)) == 0)
                {
                    UA_NodeId_copy(&ref->nodeId.nodeId, &dentry->dataserv.nodeId);
                    strncpy(dentry->dataserv.service, SERV_DATA, strlen(SERV_DATA));
                    dentry->dataserv.dev = dentry;
                }
                else if (strncmp(SERV_SETTING, ref->browseName.name.data, strlen(SERV_SETTING)) == 0)
                {
                    UA_NodeId_copy(&ref->nodeId.nodeId, &dentry->setserv.nodeId);
                    strncpy(dentry->setserv.service, SERV_SETTING, strlen(SERV_SETTING));
                    dentry->setserv.dev = dentry;
                }
                else if (strncmp(SERV_ALARM, ref->browseName.name.data, strlen(SERV_ALARM)) == 0)
                {
                    UA_NodeId_copy(&ref->nodeId.nodeId, &dentry->alarmserv.nodeId);
                    strncpy(dentry->alarmserv.service, SERV_ALARM, strlen(SERV_ALARM));
                    dentry->alarmserv.dev = dentry;
                }
            }
            else if (2 == layer)
            {
                if (strncmp(ref->browseName.name.data, "FolderType", 10) == 0)
                    continue;

                char datatag[TAG_MAX_LENGTH] = {0};
                strncpy(datatag, ref->browseName.name.data, (int)ref->browseName.name.length);
                devdata_t *cdata = NULL;
                HASH_FIND_STR(dserv->devdata, datatag, cdata);
                if (NULL == cdata)
                {
                    cdata = malloc(sizeof(devdata_t));
                    UA_NodeId_copy(&ref->nodeId.nodeId, &cdata->nodeId);
                    snprintf(cdata->datatag, TAG_MAX_LENGTH, "%s", datatag);
                    cdata->addr = 0;
                    cdata->dtype = 0;
                    dbg_syslog(LOG_ERR, "%s ```` %s ```` %s ```` %u", cdata->datatag, datatag, ref->browseName.name.data, (int)ref->browseName.name.length);
                    HASH_ADD_STR(dserv->devdata, datatag, cdata);
                    dbg_syslog(LOG_DEBUG, "add %s ```` %s ```` %s ```` %u", cdata->datatag, datatag, ref->browseName.name.data, (uint)ref->browseName.name.length);
                    dserv->datanum++;
                }
            }
        }
    }
    UA_BrowseRequest_clear(&bReq);
    UA_BrowseResponse_clear(&bResp);
    return nodecnt;
}

static void ua_mqttv1_parse_system(ua_mqttv1_var_t *var)
{
    char cmd[CMD_MAX_LENGTH] = {0};
    UA_BrowseRequest bReq;
    UA_BrowseRequest_init(&bReq);
    bReq.requestedMaxReferencesPerNode = 0;
    bReq.nodesToBrowse = UA_BrowseDescription_new();
    bReq.nodesToBrowseSize = 1;
    bReq.nodesToBrowse[0].nodeId = UA_NODEID_NUMERIC(0, UA_NS0ID_ROOTFOLDER); /* browse objects folder */
    bReq.nodesToBrowse[0].resultMask = UA_BROWSERESULTMASK_ALL; /* return everything */
    dbg_syslog(LOG_ERR, "browse system");
    UA_BrowseResponse bResp = UA_Client_Service_browse(var->client, bReq);
    dbg_syslog(LOG_ERR, "bResp.resultsSize %d", (int)bResp.resultsSize);
    for(size_t i = 0; i < bResp.resultsSize; ++i) {
        dbg_syslog(LOG_ERR, "bResp.results[i].referencesSize %d", (int)bResp.results[i].referencesSize);
        for(size_t j = 0; j < bResp.results[i].referencesSize; ++j) {
            UA_ReferenceDescription *ref = &(bResp.results[i].references[j]);
            dbg_syslog(LOG_ERR, "ref->browseName.name.data %s", ref->browseName.name.data);
            if (strncmp(FOLDER_SYSTEM, ref->browseName.name.data, strlen(FOLDER_SYSTEM)) == 0)
            {
                UA_BrowseRequest bReq2;
                UA_BrowseRequest_init(&bReq2);
                bReq2.requestedMaxReferencesPerNode = 0;
                bReq2.nodesToBrowse = UA_BrowseDescription_new();
                bReq2.nodesToBrowseSize = 1;
                bReq2.nodesToBrowse[0].nodeId = ref->nodeId.nodeId; /* browse objects folder */
                bReq2.nodesToBrowse[0].resultMask = UA_BROWSERESULTMASK_ALL; /* return everything */
                UA_BrowseResponse bResp2 = UA_Client_Service_browse(var->client, bReq2);
                dbg_syslog(LOG_ERR, "bResp2.resultsSize %d", (int)bResp2.resultsSize);
                for(size_t k = 0; k < bResp2.resultsSize; ++k) {
                    dbg_syslog(LOG_ERR, "bResp2.results[k].referencesSize %d", (int)bResp2.results[k].referencesSize);
                    for(size_t m = 0; m < bResp2.results[k].referencesSize; ++m) {
                        UA_ReferenceDescription *ref2 = &(bResp2.results[k].references[m]);
                        dbg_syslog(LOG_ERR, "ref2->browseName.name.data %s", ref2->browseName.name.data);
                        if (strncmp(SYSTEM_ABOUT, ref2->browseName.name.data, strlen(SYSTEM_ABOUT)) == 0)
                        {
                            UA_BrowseRequest bReq3;
                            UA_BrowseRequest_init(&bReq3);
                            bReq3.requestedMaxReferencesPerNode = 0;
                            bReq3.nodesToBrowse = UA_BrowseDescription_new();
                            bReq3.nodesToBrowseSize = 1;
                            bReq3.nodesToBrowse[0].nodeId = ref2->nodeId.nodeId; /* browse objects folder */
                            bReq3.nodesToBrowse[0].resultMask = UA_BROWSERESULTMASK_ALL; /* return everything */
                            UA_BrowseResponse bResp3 = UA_Client_Service_browse(var->client, bReq3);
                            dbg_syslog(LOG_ERR, "bResp3.resultsSize %d", (int)bResp3.resultsSize);
                            for (size_t p = 0; p < bResp3.resultsSize; ++p){
                                dbg_syslog(LOG_ERR, "bResp3.results[p].referencesSize %d", (int)bResp3.results[p].referencesSize);
                                for (size_t q = 0; q < bResp3.results[p].referencesSize; ++q)
                                {
                                    UA_ReferenceDescription *ref3 = &(bResp3.results[p].references[q]);
                                    dbg_syslog(LOG_ERR, "ref3->browseName.name.data %s", ref3->browseName.name.data);
                                    if ((strncmp(SYSTEM_STANO, ref3->browseName.name.data, strlen(SYSTEM_STANO)) == 0)
                                        || (ref3->nodeId.nodeId.identifierType == UA_NODEIDTYPE_STRING && (strncmp(SYSTEM_STANO, ref3->nodeId.nodeId.identifier.string.data, strlen(SYSTEM_STANO)) == 0 )))
                                    {
                                        if (1 == UA_Client_ReadSingleId(var, ref3->nodeId.nodeId, var->stano))
                                            dbg_syslog(LOG_INFO, "Obtain stano %s", var->stano);
                                        else 
                                            dbg_syslog(LOG_INFO, "Obtain stano failed");
                                    }
                                }
                            }
                            UA_BrowseRequest_clear(&bReq3);
                            UA_BrowseResponse_clear(&bResp3);
                        }
                        else if (strncmp(SYSTEM_CONFIG, ref2->browseName.name.data, strlen(SYSTEM_CONFIG)) == 0)
                        {
                            UA_BrowseRequest bReq3;
                            UA_BrowseRequest_init(&bReq3);
                            bReq3.requestedMaxReferencesPerNode = 0;
                            bReq3.nodesToBrowse = UA_BrowseDescription_new();
                            bReq3.nodesToBrowseSize = 1;
                            bReq3.nodesToBrowse[0].nodeId = ref2->nodeId.nodeId; /* browse objects folder */
                            bReq3.nodesToBrowse[0].resultMask = UA_BROWSERESULTMASK_ALL; /* return everything */
                            UA_BrowseResponse bResp3 = UA_Client_Service_browse(var->client, bReq3);
                            for (size_t p = 0; p < bResp3.resultsSize; ++p)
                                for (size_t q = 0; q < bResp3.results[p].referencesSize; ++q)
                                {
                                    UA_ReferenceDescription *ref3 = &(bResp3.results[p].references[q]);
                                    if (strncmp(SYSTEM_STANAME, ref3->browseName.name.data, strlen(SYSTEM_STANAME)) == 0)
                                    {
                                        if (1 == UA_Client_ReadSingleId(var, ref3->nodeId.nodeId, var->staname))
                                            dbg_syslog(LOG_INFO, "Obtain staname %s", var->staname);
                                        else 
                                            dbg_syslog(LOG_INFO, "Obtain staname failed");
                                    }
                                    else if (strncmp(SYSTEM_CNTID, ref3->browseName.name.data, strlen(SYSTEM_CNTID)) == 0)
                                    {
                                        if (1 == UA_Client_ReadSingleId(var, ref3->nodeId.nodeId, var->ccsn))
                                            dbg_syslog(LOG_INFO, "Obtain ccsn %s", var->ccsn);
                                        else 
                                            dbg_syslog(LOG_INFO, "Obtain ccsn failed");
                                    }
                                }

                            UA_BrowseRequest_clear(&bReq3);
                            UA_BrowseResponse_clear(&bResp3);
                        }
                    }
                }
                UA_BrowseRequest_clear(&bReq2);
                UA_BrowseResponse_clear(&bResp2);
            }
        }
    }
    
    dbg_syslog(LOG_ERR, "stano %s staname %s ccsn %s", var->stano, var->staname, var->ccsn);
    snprintf(cmd, CMD_MAX_LENGTH, "hset sys stano %s staname %s ccsn %s", var->stano, var->staname, var->ccsn);
    redisReply *reply = (redisReply *)redisCommand(var->redisContext, cmd);
    freeReplyObject(reply);
    UA_BrowseRequest_clear(&bReq);
    UA_BrowseResponse_clear(&bResp);
}

static int ds_flush_timer(ua_mqttv1_var_t *var)
{
    TIMER_CONFIRM(var->ua_sync_timer);

    for (int i = 0; i < var->dev_table.devcount; i++)
    {
        dbg_syslog(LOG_ERR, "read request on %s %s", var->dev_table.devtab[i].devsn, var->dev_table.devtab[i].devtype);
        UA_Client_ReadMultiValueAttribute(var, &var->dev_table.devtab[i].setserv);
        UA_Client_ReadMultiValueAttribute(var, &var->dev_table.devtab[i].statserv);
        UA_Client_ReadMultiValueAttribute(var, &var->dev_table.devtab[i].dataserv);
        UA_Client_ReadMultiValueAttribute(var, &var->dev_table.devtab[i].alarmserv);
    }
    return 0;
}

static void ua_mqttv1_loop(ua_mqttv1_var_t *var)
{
    int ret = -1;
    struct pollfd pfd;
    while(1)
    {
        pfd.fd = var->ua_sync_timer;
        pfd.events = POLLIN;
        pfd.revents = 0;
        ret = poll(&pfd, 0x1, 10000);
        if (ret < 0)
        {
            int error = errno;	
            dbg_syslog(LOG_WARNING, "ua_sync_timer %d errno %d\n", var->ua_sync_timer, error);
            if (error == EINTR) {
                continue;
            }
            else {
                break;
            }
        }
        else if(ret == 0)
        {
            /* do nothing */
        }
        else if(ret > 0)
        {
            if (var->ua_sync_timer > 0 && pfd.revents)
            {
                ds_flush_timer(var);
            }
        }
    }
}

static void *ua_mqttv1_redis_subthread(void *param)
{
    ua_mqttv1_var_t *var = (ua_mqttv1_var_t *)param;
    redisReply *reply;
    while(true)
    {
        void *_reply = NULL;
        if (redisGetReply(var->subContext, &_reply) != REDIS_OK)
        {
            continue;
        }
        reply = (redisReply *)_reply;
        if(reply->elements == 3)
        {
            if (strstr(reply->element[0]->str, "message") && strstr(reply->element[1]->str, UA_SETPTS))
            {
                dbg_syslog(LOG_DEBUG, "redis rcv %s %s %s", reply->element[0]->str, reply->element[1]->str, reply->element[2]->str);
                cJSON *root=cJSON_Parse(reply->element[2]->str); 
                if (root)
                {
                    cJSON *pts = cJSON_GetObjectItem(root, "pts");
                    if (pts)
                    {
                        char devsn[DEV_SN_LEN] = {0};
                        GET_JSON_VALUE_STRING(pts, "devsn", devsn);
                        UA_Client_WriteSingleNode(var, devsn, "TAG", (void *)1);
                    }
                }
                cJSON_Delete(root);
            }
        }
        freeReplyObject(reply);
    }
    return NULL;
}

static int ua_mqttv1_redis_subscribe(ua_mqttv1_var_t *var)
{
    char cmd[CMD_MAX_LENGTH];
    pthread_t subthread;

    var->subContext = redisConnect(REDIS_SERVER_IP, REDIS_SERVER_PORT);
    if (var->subContext->err)
    {
        dbg_syslog(LOG_ERR, "connect redis server failure: %s\n", var->subContext->errstr);
        redisFree(var->subContext);
        return false;
    }
    dbg_syslog(LOG_DEBUG, "connect redis server success.");

    snprintf(cmd, CMD_MAX_LENGTH, "SUBSCRIBE %s", UA_SETPTS);
    redisReply *reply = (redisReply *)redisCommand(var->subContext, cmd);
    if (NULL == reply)
    {
        dbg_syslog(LOG_ERR, "subscribe %s failed", cmd);
        return false;
    }
    dbg_syslog(LOG_DEBUG, "subscribe success");
    freeReplyObject(reply);

    pthread_create(&subthread, NULL, ua_mqttv1_redis_subthread, (void *)var);
    return true;
}

static void ua_mqttv1_loadcfg(ua_mqttv1_var_t *var, char *cfgpath)
{
    char *data = NULL;
    cJSON *root = NULL;
    int lens;

    if (cfgpath == NULL)
    {
        dbg_syslog(LOG_ERR, "cfg_file:%s not exist, use defaut MQTT server", cfgpath);
        return;
    }

    data = read_file_data_and_lens(cfgpath, &lens);
    if (data == NULL)
    {
        dbg_syslog(LOG_ERR, "read cfg file %s error, load default", cfgpath);
        return;
    }

    root = cJSON_Parse(data);
    if (!root)
    {
        dbg_syslog(LOG_ERR, "parse cfg file %s error, load default", cfgpath);
        return;
    }

    GET_JSON_VALUE_STRING(root, "endpointurl", var->endpointUrl);
    GET_JSON_VALUE_STRING(root, "username", var->usrname);
    GET_JSON_VALUE_STRING(root, "password", var->userpass);
    GET_JSON_VALUE_INT(root, "syncinterval", var->syncinterval);
    GET_JSON_VALUE_INT(root, "proto", var->proto);

    dbg_syslog(LOG_INFO, "url:%s name:%s pass:%s interval:%d", var->endpointUrl, var->usrname, var->userpass, var->syncinterval);
    cJSON_Delete(root);
}

static void ua_mqttv1_init(ua_mqttv1_var_t *var)
{
    char cmd[CMD_MAX_LENGTH] = {0};
    ua_mqttv1_connredis(&var->redisContext);
    ua_mqttv1_loadcfg(var, UA_MQTTV1_CFG_PATH);
    ua_mqttv1_read_devtype(var, UA_MQTTV1_DEVTYPE);

    snprintf(cmd, CMD_MAX_LENGTH, "del uadevs");
    redisReply *reply = (redisReply *)redisCommand(var->redisContext, cmd);
    freeReplyObject(reply);

    var->client = UA_Client_new();
    UA_ClientConfig *config = UA_Client_getConfig(var->client);

    UA_ByteString certificate = loadFile(CERTFILE);
    UA_ByteString privateKey  = loadFile(KEYFILE);
    UA_ClientConfig_setDefaultEncryption(config, certificate, privateKey,
                                            NULL, 0, NULL, 0);
    UA_String securityPolicyUri = UA_String_fromChars("http://opcfoundation.org/UA/SecurityPolicy#None");
    UA_MessageSecurityMode securityMode = UA_MESSAGESECURITYMODE_INVALID;
    UA_ByteString_clear(&certificate);
    UA_ByteString_clear(&privateKey);

    UA_ClientConfig_setDefault(config);

    config->securityMode = securityMode;
    config->securityPolicyUri = securityPolicyUri;

    UA_StatusCode retval = UA_Client_connectUsername(var->client, var->endpointUrl, var->usrname, var->userpass);
    if(retval != UA_STATUSCODE_GOOD) {
        dbg_syslog(LOG_ERR, "Could not connect with retval %x\n", retval);
        UA_Client_delete(var->client);
        return;
    }

    dbg_syslog(LOG_ERR, "Connect to ua server succ");
    ua_mqttv1_parse_system(var);

    /* Browse all data objects */
    dbg_syslog(LOG_ERR, "Browsing nodes in objects folder:\n");
    UA_NodeId topnode[MAX_L1_NODES] = {0};
    UA_NodeId rootnode = UA_NODEID_NUMERIC(0, UA_NS0ID_OBJECTSFOLDER);
    int l1nodecnt = node_browse(var, rootnode, topnode, 0, NULL);

    for(int i = 0; i < l1nodecnt; i++)
    {
        if (topnode[i].identifierType == UA_NODEIDTYPE_STRING)
        {
            dbg_syslog(LOG_ERR, "nodeid: %-16.*s", (int)topnode[i].identifier.string.length, topnode[i].identifier.string.data);
            devtype_t *curdev;
            char typenam[MAX_DEVTYPE_LEN] = {0};
            strncpy(typenam, topnode[i].identifier.string.data, (int)topnode[i].identifier.string.length);
            for(int k = strlen(typenam) - 1; k >= 0; k--)
            {
                if (typenam[k]>= '0' && typenam[k] <= '9')
                    typenam[k] = 0;
                else 
                    break;
            }
            HASH_FIND_STR(var->devtypes, typenam, curdev);
            if (curdev != NULL)
            {
                int devno = 0;
                int caddr = curdev->addr++;
                dev_entry_t *dentry = &var->dev_table.devtab[var->dev_table.devcount];
                snprintf(dentry->devsn, DEV_SN_LEN, "%-16.*s", (int)topnode[i].identifier.string.length, topnode[i].identifier.string.data);
                snprintf(dentry->devtype, MAX_DEVTYPE_LEN, "%s", typenam);
                dentry->dev_no = devno;
                dentry->common_addr = caddr;
                dentry->yxaddr = 1;
                dentry->ycaddr = 16385;
                dentry->modbus_addr = 1;
                var->dev_table.devcount++;
                
                snprintf(cmd, CMD_MAX_LENGTH, "hset uadevs %-16.*s %s|0|%d", (int)topnode[i].identifier.string.length, topnode[i].identifier.string.data, typenam, caddr);
                reply = (redisReply *)redisCommand(var->redisContext, cmd);
                freeReplyObject(reply);
            
                node_browse(var, topnode[i], NULL, 1, dentry);

                node_browse(var, dentry->statserv.nodeId, NULL, 2, &dentry->statserv);
                node_browse(var, dentry->dataserv.nodeId, NULL, 2, &dentry->dataserv);
                node_browse(var, dentry->alarmserv.nodeId, NULL, 2, &dentry->alarmserv);
                node_browse(var, dentry->setserv.nodeId, NULL, 2, &dentry->setserv);
            }
        }
    }

    var->ua_sync_timer = my_timer_create();
    if (var->ua_sync_timer > 0)
    {
        my_timer_set(var->ua_sync_timer, 1, var->syncinterval * 1000);
    }

    ua_mqttv1_redis_subscribe(var);
}

int main(int argc, char *argv[]) {
    ua_mqttv1_var_t *var = &gvar;
    lnxall_loglevel_set(LNXALL_LOGDEBUG, 1);
    ua_mqttv1_init(var);
    ua_mqttv1_loop(var);

    UA_Client_disconnect(var->client);
    UA_Client_delete(var->client);
    return EXIT_SUCCESS;
}
