#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 "scada_iec104.h"

#define SCADA_IEC104_APP_KEY "NRgYk9hJRt0RiEZU"
#define SCADA_IEC104_NODE_STATUS_BAK_FILE "/app/scada_iec104_status.bak"

scada_iec104_var_t gvar = {0};

#if 0
/* Callback handler to log sent or received messages (optional) */
static void rawMessageHandler(void *parameter, uint8_t *msg, int msgSize, bool sent)
{
    connect_config_t *connect_cfg = (connect_config_t *)parameter;
    if (sent)
    {
        dy_syslog_hex(LOG_DEBUG, msg, msgSize, "len %d			网关->[%s]		", msgSize, connect_cfg->socket_server_addr);
    }
    else
    {
        dy_syslog_hex(LOG_DEBUG, msg, msgSize, "len %d			网关<-[%s]		", msgSize, connect_cfg->socket_server_addr);
    }
}
#endif
/* Connection event handler */
static void connectionHandler(void *parameter, CS104_Connection connection, CS104_ConnectionEvent event)
{
    connect_config_t *connect_cfg = (connect_config_t *)parameter;
    switch (event)
    {
    case CS104_CONNECTION_OPENED:
        connect_cfg->connect_status = 1;
        dy_syslog(LOG_INFO, "Connection established");
        break;
    case CS104_CONNECTION_CLOSED:
        connect_cfg->connect_status = 0;
        dy_syslog(LOG_INFO, "Connection closed");
        break;
    case CS104_CONNECTION_STARTDT_CON_RECEIVED:
        dy_syslog(LOG_INFO, "Received STARTDT_CON");
        connect_cfg->connect_status = 2;
        break;
    case CS104_CONNECTION_STOPDT_CON_RECEIVED:
        connect_cfg->connect_status = 1;
        dy_syslog(LOG_INFO, "Received STOPDT_CON");
        break;
    }
}

/*
 * CS101_ASDUReceivedHandler implementation
 *
 * For CS104 the address parameter has to be ignored
 */
static bool asduReceivedHandler(void *parameter, int address, CS101_ASDU asdu)
{
    connect_config_t *connect_cfg = (connect_config_t *)parameter;
    int type = CS101_ASDU_getTypeID(asdu);
    int i;
    time_t now = time(NULL);

    // dy_syslog(LOG_INFO, "RECVD ASDU type: %s(%i) elements: %i",
    //           TypeID_toString(CS101_ASDU_getTypeID(asdu)),
    //           CS101_ASDU_getTypeID(asdu),
    //           CS101_ASDU_getNumberOfElements(asdu));

    for (i = 0; i < CS101_ASDU_getNumberOfElements(asdu); i++)
    {
        InformationObject io = CS101_ASDU_getElement(asdu, i);
        int addr = InformationObject_getObjectAddress(io);
        int L1_index = (addr >> L2_SHIFT) & L1_MASK;
        int L2_index = addr & L2_MASK;
        L2_data_table_t *l2 = connect_cfg->L1_data.L2_data[L1_index];

        if (l2 == NULL)
        {
            continue;
        }
        data_etnry_t *data = l2->data;
        if (data == NULL)
        {
            continue;
        }
        data[L2_index].update_ts = now;

        switch (type)
        {
        case M_SP_NA_1:
        {
            SinglePointInformation t = (SinglePointInformation)io;
            data[L2_index].val.val_bool = SinglePointInformation_getValue(t);
            data[L2_index].type = TYPE_BOOL;
            SinglePointInformation_destroy(t);
            //dy_syslog(LOG_DEBUG, "addr:%d L1_index:%d L2_index:%d val:%d", addr, L1_index, L2_index, data[L2_index].val.val_bool);
            break;
        }
        case M_ME_NC_1:
        {
            MeasuredValueShort t = (MeasuredValueShort)io;
            data[L2_index].val.val_float = MeasuredValueShort_getValue(t);
            data[L2_index].type = TYPE_FLOAT;
            MeasuredValueShort_destroy(t);
            //dy_syslog(LOG_DEBUG, "addr:%d L1_index:%d L2_index:%d val:%f", addr, L1_index, L2_index, data[L2_index].val.val_float);
            break;
        }
        case M_IT_NA_1:
        {
            IntegratedTotals t = (IntegratedTotals)io;
            data[L2_index].val.val_int = BinaryCounterReading_getValue(IntegratedTotals_getBCR(t));
            data[L2_index].type = TYPE_INT;
            IntegratedTotals_destroy(t);
            //dy_syslog(LOG_DEBUG, "addr:%d L1_index:%d L2_index:%d val:%d", addr, L1_index, L2_index, data[L2_index].val.val_int);
            break;
        }
        case M_ME_NB_1:
        {
            MeasuredValueScaled t = (MeasuredValueScaled)io;
            data[L2_index].val.val_int = MeasuredValueScaled_getValue(t);
            data[L2_index].type = TYPE_INT;
            MeasuredValueScaled_destroy(t);
            // dy_syslog(LOG_DEBUG, "addr:%d L1_index:%d L2_index:%d val:%d", addr, L1_index, L2_index, data[L2_index].val.val_int);
            break;
        }
        default:
            break;
        }
    }

    return true;
}

static data_etnry_t *scada_iec104_get_cache(scada_iec104_var_t *var, connect_config_t *connect_cfg, int addr)
{
    int L1_index = (addr >> L2_SHIFT) & L1_MASK;
    int L2_index = addr & L2_MASK;
    data_etnry_t *entry = NULL;

    L2_data_table_t *l2 = connect_cfg->L1_data.L2_data[L1_index];
    if (l2 != NULL)
    {
        entry = &l2->data[L2_index];
    }

    return entry;
}

static int scada_iec104_create_cache(scada_iec104_var_t *var, rglt_table_t *rt)
{
    int i = 0;

    for (i = 0; i < rt->rglt_cnt; i++)
    {
        int addr = rt->entry[i].addr;
        int L1_index = (addr >> L2_SHIFT) & L1_MASK;
        L2_data_table_t *l2 = rt->node->connect_config->L1_data.L2_data[L1_index];
        if (l2 == NULL)
        {
            l2 = calloc(1, sizeof(L2_data_table_t));
            rt->node->connect_config->L1_data.L2_data[L1_index] = l2;
        }
        data_etnry_t *data = l2->data;
        if (data == NULL)
        {
            data = calloc(1, sizeof(data_etnry_t) * L2);
            l2->cnt = L2;
            l2->data = data;
        }
    }
    rt->cache_created = true;

    return 0;
}

static int scada_iec104_handle_rglt(scada_iec104_var_t *var, pp_regulate_signal_t *regulate_data)
{
    rglt_table_t *rt = NULL;
    int i;
    int cnt = 0;

    if (regulate_data->data == NULL)
    {
        dy_syslog(LOG_ERR, "data is empty");
        return -1;
    }
    rt = (rglt_table_t *)regulate_data->data;
    dy_syslog(LOG_DEBUG, "send sn:%s rglt_cnt:%d", regulate_data->sn, rt->rglt_cnt);

    if (!rt->cache_created && regulate_data->period > 0)
    {
        scada_iec104_create_cache(var, rt);
    }

    if (rt->node->connect_config->connect_status != 2)
    {
        dy_syslog(LOG_DEBUG, "%s not connected, status:%d", rt->node->sn, rt->node->connect_config->connect_status);
        return -1;
    }

    cJSON *rt_data = cJSON_CreateObject();
    static char *buf = NULL;
    char topic[192] = {0};
    time_t now = time(NULL);
    cJSON_AddStringToObject(rt_data, "port", "IEC104");
    cJSON_AddStringToObject(rt_data, "sn", regulate_data->sn);
    cJSON_AddNumberToObject(rt_data, "ts", now);
    cJSON_AddNumberToObject(rt_data, "mi", regulate_data->mi);
    cJSON_AddStringToObject(rt_data, "src_identifier", regulate_data->src_identifier);

    cJSON *data_json = cJSON_CreateObject();

    for (i = 0; i < rt->rglt_cnt; i++)
    {
        data_etnry_t *data = NULL;
        data = scada_iec104_get_cache(var, rt->node->connect_config, rt->entry[i].addr);
        if (data != NULL)
        {
            if (data->update_ts != 0)
            {
                cnt++;
                switch (data->type)
                {
                case TYPE_INT:
                    cJSON_AddNumberToObject(data_json, rt->entry[i].id, data->val.val_int * rt->entry[i].k);
                    break;
                case TYPE_BOOL:
                    cJSON_AddBoolToObject(data_json, rt->entry[i].id, data->val.val_bool);
                    break;
                case TYPE_FLOAT:
                    cJSON_AddNumberToObject(data_json, rt->entry[i].id, data->val.val_float * rt->entry[i].k);
                    break;
                case TYPE_SHORT:
                    cJSON_AddNumberToObject(data_json, rt->entry[i].id, data->val.val_short * rt->entry[i].k);
                    break;
                default:
                    break;
                }
            }
        }
        else
        {
            dy_syslog(LOG_ERR, "addr:%d no data", rt->entry[i].addr);
            continue;
        }
    }

    if (cnt > 0)
    {
        char *tmp = cJSON_Print(data_json);
        if (buf == NULL)
        {
            buf = calloc(1, 8192);
        }
        b64_encode(tmp, strlen(tmp), buf, 8192);
        cJSON_AddNumberToObject(rt_data, "len", strlen(tmp));
        cJSON_AddStringToObject(rt_data, "data_b64", buf);
        char *data_tmp = cJSON_Print(rt_data);

        snprintf(topic, sizeof(topic), "ipc/%s/%s/device/%s/data/%s", var->sn_str, "IEC104",
                 regulate_data->sn, "raw_data");

        ipc_session_publish(var->session, topic, data_tmp, strlen(data_tmp));
        flush_device_status_by_sn(var->ds, ACTION_RCV, regulate_data->sn);
        free(data_tmp);
        free(tmp);
    }
    cJSON_Delete(data_json);
    cJSON_Delete(rt_data);
    return 0;
}

static int scada_iec104_heart_beat(scada_iec104_var_t *var)
{
    int ret;
    static int all_started = 0;
    static time_t last_trigger_ts = 0;
    time_t now = time(NULL);

    TIMER_CONFIRM(var->heartbeat_timer);

    while (1)
    {
        period_msg_t *msg = period_msg_next(var->period);
        if (msg && msg->started)
        {
            ret = scada_iec104_handle_rglt(var, msg->rs);
            if (ret < 0)
            {
                dy_syslog(LOG_ERR, "parse real data structure failed");
            }
        }
        else
        {
            break;
        }
    }

    connect_config_t *connect_cfg = NULL;
    list_for_each_entry(connect_cfg, &var->connect_list, list)
    {
        if (connect_cfg->connect_status == 0)
        {
            if (now - connect_cfg->last_con > 10)
            {
                dy_syslog(LOG_DEBUG, "try to reconnect");
                CS104_Connection_connectAsync(connect_cfg->con);
                connect_cfg->last_con = now;
            }
            continue;
        }
        else if (connect_cfg->connect_status == 1)
        {
            dy_syslog(LOG_DEBUG, "CS104_Connection_sendStartDT");
            CS104_Connection_sendStartDT(connect_cfg->con);
        }
        else if (connect_cfg->connect_status == 2 && now - connect_cfg->last_send > 60)
        {
            connect_cfg->last_send = now;
            dy_syslog(LOG_DEBUG, "CS104_Connection_sendCounterInterrogationCommand...");
            CS104_Connection_sendCounterInterrogationCommand(connect_cfg->con, CS101_COT_ACTIVATION, 1, IEC60870_QCC_FRZ_READ + IEC60870_QCC_RQT_GENERAL);
            CS104_Connection_sendInterrogationCommand(connect_cfg->con, CS101_COT_ACTIVATION, 1, IEC60870_QOI_STATION);
        }
    }

    if (now - last_trigger_ts > 20)
    {
        if (all_started == 0)
        {
            all_started = period_services_retrigger_by_appkey(var->session, var->period, var->nodes_cfg_table, var->template_table, SCADA_IEC104_APP_KEY);
        }
        last_trigger_ts = now;
    }

    return 0;
}

static int ds_flush_timer(scada_iec104_var_t *var)
{
    time_t now = time(NULL);

    TIMER_CONFIRM(var->node_status_timer);    

    if (now >= var->nodes_status_sync_date + 30)
    {
		size_t nlen;
        static time_t lasttime = 0;
        var->nodes_status_sync_date = now;
        char *node_status = ds_print(var->ds);

		nlen = node_status ? strlen(node_status) : 0;
        dy_syslog(LOG_INFO, "online_status_changed:%d strlen(node_status):%d, (now-lasttime):%d",
                  var->ds->online_status_changed, (int) nlen, (int) (now - lasttime));
        if (nlen > 20 && (var->ds->online_status_changed || now > lasttime + 3600)) //在线状态变化上报，1小时也会上报
        {
            char topic[TOPIC_MAX_LEN];
            int ret = 0;

            snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/%s/%s", var->sn_str, var->proc_name, TOPIC_EVT_NODESSTATUS);
            ret = ipc_session_publish(var->session, topic, node_status, strlen(node_status));
            if (ret == 0)
            {
                lasttime = now;
                var->ds->online_status_changed = false;
            }
            node_sta_backup(var->ds, SCADA_IEC104_NODE_STATUS_BAK_FILE);
        }
        if (node_status)
        {
            write_file_data(NODES_CACHE "/scada_iec104_status.json", node_status, strlen(node_status));
            free(node_status);
        }
    }

    return 0;
}

static void scada_iec104_loop(scada_iec104_var_t *var)
{
    int ret = -1, maxfd;
    fd_set rset;
    struct timeval timeout;

    while (1)
    {
        SELECT_INIT();
        SELECT_ADD_FD(var->heartbeat_timer);
        SELECT_ADD_FD(var->node_status_timer);
        timeout.tv_usec = 0;
        timeout.tv_sec = 10;
        ret = select(maxfd + 1, &rset, 0, 0, &timeout);
        if (ret < 0)
        {
            printf("errno %d\n", errno);
            if (errno == EINTR)
            {
                continue;
            }
            else
            {
                break;
            }
        }
        else if (ret == 0)
        {
        }
        else
        {
            if (var->heartbeat_timer > 0 && FD_ISSET(var->heartbeat_timer, &rset))
            {
                FD_CLR(var->heartbeat_timer, &rset);
                scada_iec104_heart_beat(var);
            }
            if (var->node_status_timer > 0 && FD_ISSET(var->node_status_timer, &rset))
            {
                FD_CLR(var->node_status_timer, &rset);
                ds_flush_timer(var);
            }
        }
    }
}

static void scada_iec104_subscribe_all(scada_iec104_var_t *var)
{
    ipc_session_t *ipc_session = var->session;
    char topic[TOPIC_MAX_LEN] = {0};

    snprintf(topic, TOPIC_MAX_LEN, "ipc/+/%s/device/+/data/%s", SCADA_IEC104_APP_KEY, TOPIC_SEND_RGLT_SIGNAL_RAW_DATA);
    ipc_session_subscribe(ipc_session, topic);
}

static int scada_iec104_mqtt_handle_recv_msg(void *obj, ipc_msg_t *mqtt_msg)
{
    scada_iec104_var_t *var = (scada_iec104_var_t *)obj;

    dy_syslog(LOG_DEBUG, " MQTT client: received MQTT topic:%s payload length:%d",
              mqtt_msg->topic, mqtt_msg->payloadLen);
    if (strstr(mqtt_msg->topic, TOPIC_SEND_RGLT_SIGNAL_RAW_DATA))
    {
        pp_regulate_signal_t regulate_data = {0};
        int ret;
        char *data_str = NULL;
        int node_cnt;
        int i;
        rglt_table_t *rt = NULL;
        int size;
        dev_node_t *dev = NULL;
        int find = 0;

        //先检查参数
        cJSON *root = cJSON_Parse(mqtt_msg->payload);
        if (!root)
        {
            return -1;
        }
        GET_JSON_VALUE_DY_STRING(root, "data_b64", data_str);
        GET_JSON_VALUE_INT(root, "period", regulate_data.period);
        GET_JSON_VALUE_INT(root, "mi", regulate_data.mi);
        GET_JSON_VALUE_STRING(root, "src_identifier", regulate_data.src_identifier);
        GET_JSON_VALUE_STRING(root, "sn", regulate_data.sn);
        cJSON_Delete(root);

        list_for_each_entry(dev, &var->node_list, list)
        {
            if (strcmp(dev->sn, regulate_data.sn) == 0)
            {
                find = 1;
                break;
            }
        }
        if (find == 0)
        {
            goto out;
        }

        root = cJSON_Parse(data_str);
        if (!root)
        {
            goto out;
        }
        // {"operation":"read","item_table":[{"id":"A","Addr":100,"k":1.2},{"id":"B","Addr":110}]}
        cJSON *nodes = cJSON_GetObjectItem(root, "item_table");
        if (!nodes)
        {
            dy_syslog(LOG_ERR, "get item_table failed");
            goto out;
        }

        node_cnt = cJSON_GetArraySize(nodes);
        size = sizeof(rglt_table_t) + node_cnt * sizeof(rglt_entry_t);
        rt = calloc(1, size);
        rt->rglt_cnt = node_cnt;
        rt->node = dev;
        for (i = 0; i < node_cnt; i++)
        {
            cJSON *node = cJSON_GetArrayItem(nodes, i);
            GET_JSON_VALUE_INT(node, "Addr", rt->entry[i].addr);
            GET_JSON_VALUE_STRING(node, "id", rt->entry[i].id);
            cJSON *k = cJSON_GetObjectItem(node, "k");
            if (k != NULL)
            {
                GET_JSON_VALUE_DOUBLE(node, "k", rt->entry[i].k);
            }
            else
            {
                rt->entry[i].k = 1;
            }
        }
        free(data_str);
        data_str = NULL;
        cJSON_Delete(root);
        regulate_data.data = (void *)rt;
        regulate_data.len = size;

        if (regulate_data.period > 0)
        {
            period_msg_t once;

            snprintf(once.key, 128, "%s_%s", regulate_data.sn, regulate_data.src_identifier);
            dy_syslog(LOG_INFO, "Recevied interval ctrl cmd periad %d key:%s", regulate_data.period, once.key);
            once.interval = regulate_data.period;
            once.last_poll = 0;
            once.rs = calloc(1, sizeof(pp_regulate_signal_t));
            memcpy(once.rs, &regulate_data, sizeof(pp_regulate_signal_t));
            once.started = 1;
            once.rs->data = calloc(1, regulate_data.len);
            memcpy(once.rs->data, regulate_data.data, regulate_data.len);

            period_msg_add(var->period, &once);
        }
        else
        {
            ret = scada_iec104_handle_rglt(var, &regulate_data);
            if (ret < 0)
            {
                dy_syslog(LOG_ERR, "parse real data structure failed");
                goto out;
            }
        }

    out:
        free(data_str);
        free(regulate_data.data);

        return 0;
    }

    return 0;
}

// 建立与内部broker之间的MQTT连接
static int scada_iec104_mqtt_client_init(scada_iec104_var_t *var)
{
    char clientId[MAX_CLIENT_ID_LEN] = {0};

    snprintf(clientId, MAX_CLIENT_ID_LEN, "INIT_SCADA_IEC104");
    var->session = ipc_session_new(clientId, (void*)var, IPC_DEFAULT);
    if (var->session == NULL)
        return -1;

    ipc_session_set_callbacks(var->session, scada_iec104_mqtt_handle_recv_msg, NULL);
    scada_iec104_subscribe_all(var);
    ipc_session_start(var->session);
    return 0;
}

static int scada_iec104_init(scada_iec104_var_t *var)
{
    connect_config_t *connect_cfg = NULL;
    dev_node_t *node = NULL;
    int i;

    INIT_LIST_HEAD(&var->connect_list);
    INIT_LIST_HEAD(&var->node_list);

    get_board_sn(var->sn_str);
    get_process_name(var->proc_name);
    time_t now = time(NULL);

    if (load_nodes_cfg(&var->nodes_cfg_table, NODES_CFG_PATH) == -1)
    {
        dy_syslog(LOG_ERR, "load nodes cfg fail");
        return -1;
    }
    if (load_templates_cfg(&var->template_table, TEMPLATES_CFG_PATH) == -1)
    {
        dy_syslog(LOG_ERR, "load template cfg fail");
        return -1;
    }

    dev_status_init_by_app_key(&var->ds, var->nodes_cfg_table, var->template_table, SCADA_IEC104_APP_KEY);
    node_sta_recovery(var->ds, SCADA_IEC104_NODE_STATUS_BAK_FILE);

    scada_iec104_mqtt_client_init(var);
    Lib60870_enableDebugOutput(true);
    for (i = 0; i < var->nodes_cfg_table->node_cnt; i++)
    {
        dy_syslog(LOG_DEBUG, "appkey:%s ip:%s port:%d", var->nodes_cfg_table->node[i].app_key, var->nodes_cfg_table->node[i].tcp_ip_addr, var->nodes_cfg_table->node[i].tcp_port);
        if (strcmp(var->nodes_cfg_table->node[i].app_key, SCADA_IEC104_APP_KEY) == 0)
        {
            int find = 0;
            list_for_each_entry(connect_cfg, &var->connect_list, list)
            {
                if (strcmp(var->nodes_cfg_table->node[i].tcp_ip_addr, connect_cfg->socket_server_addr) == 0 &&
                    var->nodes_cfg_table->node[i].tcp_port == connect_cfg->socket_server_port)
                {
                    find = 1;
                    break;
                }
            }

            if (find == 0)
            {
                connect_cfg = calloc(1, sizeof(connect_config_t));
                connect_cfg->var = var;

                strncpy(connect_cfg->socket_server_addr, var->nodes_cfg_table->node[i].tcp_ip_addr, sizeof(connect_cfg->socket_server_addr) - 1);
                connect_cfg->socket_server_port = var->nodes_cfg_table->node[i].tcp_port;

                dy_syslog(LOG_INFO, "==%d socket_addr:%s port:%d ==", i, connect_cfg->socket_server_addr, connect_cfg->socket_server_port);
                list_add_tail(&connect_cfg->list, &var->connect_list);
                var->conn_cnt++;

                connect_cfg->con = CS104_Connection_create(connect_cfg->socket_server_addr, connect_cfg->socket_server_port);

                // CS101_AppLayerParameters alParams = CS104_Connection_getAppLayerParameters(connect_cfg->con);
                // alParams->originatorAddress = 3;

                CS104_Connection_setConnectionHandler(connect_cfg->con, connectionHandler, (void *)connect_cfg);
                CS104_Connection_setASDUReceivedHandler(connect_cfg->con, asduReceivedHandler, (void *)connect_cfg);

                /* uncomment to log messages */
                // sCS104_Connection_setRawMessageHandler(connect_cfg->con, rawMessageHandler, (void *)connect_cfg);
                CS104_Connection_connectAsync(connect_cfg->con);
                connect_cfg->last_con = now;
            }

            node = calloc(1, sizeof(dev_node_t));
            node->var = var;
            strncpy(node->sn, var->nodes_cfg_table->node[i].sn, sizeof(node->sn) - 1);
            node->connect_config = connect_cfg;
            list_add_tail(&node->list, &var->node_list);
            var->node_cnt++;
        }
    }

    var->heartbeat_timer = my_timer_create();
    if (var->heartbeat_timer > 0)
    {
        if (var->conn_cnt)
            my_timer_set(var->heartbeat_timer, 3, 300);
        else
            my_timer_set(var->heartbeat_timer, 3, 10000);
    }
    var->node_status_timer = my_timer_create();
    if (var->node_status_timer > 0)
    {
        my_timer_set(var->node_status_timer, 1, 5000);
    }

    period_msg_init(&var->period, 64);
    // period_services_retrigger_by_appkey(var->session, var->period, var->nodes_cfg_table, var->template_table, SCADA_IEC104_APP_KEY);

    dy_syslog(LOG_INFO, "init done, board SN:%s", var->sn_str);

    return 0;
}

int main(int argc, char *argv[])
{
    scada_iec104_var_t *var = &gvar;

    openlog("IEC104", LOG_PID, LOG_DAEMON);

    memset(var, 0, sizeof(scada_iec104_var_t));

    scada_iec104_init(var);
    scada_iec104_loop(var);

    return 0;
}
