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

gw_port_e ports[] = {S7_PLC};

S7Object Client;
unsigned char Buffer[65536]; // 64 K buffer

/*

const byte S7AreaPE   =	0x81;
const byte S7AreaPA   =	0x82;
const byte S7AreaMK   =	0x83;
const byte S7AreaDB   =	0x84;
const byte S7AreaCT   =	0x1C;
const byte S7AreaTM   =	0x1D;


S7WLBit 1
S7WLByte 1
S7WLWord 2
S7WLDWord 4
S7WLReal 4
S7WLCounter 2
S7WLTimer 2

*/

int get_size_from_wordlen(int WordLen)
{
    if (WordLen == S7WLBit)
    {
        return 1;
    }
    else if (WordLen == S7WLByte)
    {
        return 1;
    }
    else if (WordLen == S7WLWord)
    {
        return 2;
    }
    else if (WordLen == S7WLDWord)
    {
        return 4;
    }
    else if (WordLen == S7WLReal)
    {
        return 4;
    }
    else if (WordLen == S7WLCounter)
    {
        return 2;
    }
    else if (WordLen = S7WLTimer)
    {
        return 2;
    }
    return -1;
}

//------------------------------------------------------------------------------
// Check error
//------------------------------------------------------------------------------
int Check(int Result, char *function)
{
    int ExecTime;
    char text[1024];
    if (Result == 0)
    {
        Cli_GetExecTime(Client, &ExecTime);
        dy_syslog(LOG_DEBUG, "%s Result: OK Execute time: %d ms", function, ExecTime);
    }
    else
    {
        if (Result < 0)
        {
            dy_syslog(LOG_DEBUG, "%s Library Error: %d", function, Result);
        }
        else
        {
            Cli_ErrorText(Result, text, 1024);
            dy_syslog(LOG_DEBUG, "%s %d %s", function, Result, text);
        }
    }
    return !Result;
}

static int s7plc_reconnect(s7_plc_var_t *var, plc_client_t *client)
{
    Cli_Disconnect(client->Client);
    sleep(1);
    Cli_Connect(client->Client);
    dy_syslog(LOG_DEBUG, "reconnect");

    return 0;
}

static int s7plc_send_request(s7_plc_var_t *var, plc_client_t *client, pp_regulate_signal_t *regulate_data, TS7DataItem *Items, int cnt, int opcode)
{
    int res = -1;
    int i, j;
    pp_real_time_data_t real_data = {0};
    char *rsp_str = NULL;
    char text[4096] = {0};
    int pos = 0;
    int IsConnected = 0;
    char *encode_str = NULL;

    // force opcode to 0
    opcode = 0;

    flush_device_status_by_sn(var->ds, ACTION_SND, regulate_data->sn);
    if (opcode == 0)
    {
        cJSON *nodes = cJSON_CreateArray();

        for (i = 0; i < cnt; i++)
        {
            cJSON *item = cJSON_CreateObject();
            cJSON_AddItemToArray(nodes, item);
            pos = 0;

            res = Cli_GetConnected(client->Client, &IsConnected);

            if (Check(res, "get connected"))
            {
                if (IsConnected == 0)
                {
                    s7plc_reconnect(var, client);
                }
                Cli_GetConnected(client->Client, &IsConnected);
                if (IsConnected != 0)
                {
                    Items[i].Result = Cli_ReadArea(client->Client, Items[i].Area, Items[i].DBNumber, Items[i].Start, Items[i].Amount, Items[i].WordLen, Items[i].pdata);
                    if (Check(Items[i].Result, "read Area"))
                    {
                        dy_syslog_hex(LOG_DEBUG, Items[i].pdata, get_size_from_wordlen(Items[i].WordLen), "got data");
                    }

                    if (Items[i].Result == 0)
                    {
                        cJSON_AddNumberToObject(item, "Result", Items[i].Result);
                        cJSON_AddNumberToObject(item, "Area", Items[i].Area);
                        if (Items[i].Area == 0x84)
                        {
                            cJSON_AddNumberToObject(item, "DBNumber", Items[i].DBNumber);
                        }
                        cJSON_AddNumberToObject(item, "Start", Items[i].Start);
                        cJSON_AddNumberToObject(item, "Amount", Items[i].Amount);
                        cJSON_AddNumberToObject(item, "WordLen", Items[i].WordLen);
                        if (encode_str == NULL)
                        {
                            encode_str = (char *)malloc(MAXBUF);
                        }
                        b64_encode(Items[i].pdata, Items[i].Amount * get_size_from_wordlen(Items[i].WordLen), encode_str, MAXBUF);
                        cJSON_AddStringToObject(item, "Data_B64", encode_str);

                        if (i == cnt -1)
                        {
                            rsp_str = cJSON_PrintUnformatted(nodes);
                            cJSON_Delete(nodes);
                        }
                    }
                    else
                    {
                        if (Items[i].Result & 0xFFFF)
                        {
                            dy_syslog(LOG_ERR, "Items[i].Result:%d", Items[i].Result);
                            s7plc_reconnect(var, client);
                        }
                    }
                }
            }
            else
            {
                dy_syslog(LOG_ERR, "wrong status for client");
            }
        }
    }
    else if (opcode == 1)
    {
        cJSON *nodes = cJSON_CreateArray();
        for (i = 0; i < cnt; i++)
        {
            cJSON *item = cJSON_CreateObject();
            cJSON_AddItemToArray(nodes, item);
            pos = 0;

            res = Cli_GetConnected(client->Client, &IsConnected);

            if (Check(res, "get connected"))
            {
                if (IsConnected == 0)
                {
                    s7plc_reconnect(var, client);
                }
                Cli_GetConnected(client->Client, &IsConnected);
                if (IsConnected != 0)
                {
                    Items[i].Result = Cli_WriteArea(client->Client, Items[i].Area, Items[i].DBNumber, Items[i].Start, Items[i].Amount, Items[i].WordLen, Items[i].pdata);
                    if (Check(Items[i].Result, "write Area"))
                    {
                        dy_syslog_hex(LOG_DEBUG, Items[i].pdata, get_size_from_wordlen(Items[i].WordLen), "got data");
                    }
                    cJSON_AddNumberToObject(item, "Result", Items[i].Result);
                    cJSON_AddNumberToObject(item, "Area", Items[i].Area);
                    if (Items[i].Area == 0x84)
                    {
                        cJSON_AddNumberToObject(item, "DBNumber", Items[i].DBNumber);
                    }
                    cJSON_AddNumberToObject(item, "Start", Items[i].Start);
                    cJSON_AddNumberToObject(item, "Amount", Items[i].Amount);
                    cJSON_AddNumberToObject(item, "WordLen", Items[i].WordLen);
                }
            }
        }
        rsp_str = cJSON_PrintUnformatted(nodes);
        cJSON_Delete(nodes);
    }

    if (rsp_str != NULL)
    {
        real_data.mi = regulate_data->mi;
        strcpy(real_data.sn, regulate_data->sn);
        strncpy(real_data.src_identifier, regulate_data->src_identifier, sizeof(real_data.src_identifier));

        real_data.port = regulate_data->port;
        real_data.len = strlen(rsp_str);
        real_data.data = rsp_str;
        real_data.ts = time(NULL);

        send_to_proto_parser(var->session, &real_data);
        flush_device_status_by_sn(var->ds, ACTION_RCV, real_data.sn);
    }

    free(rsp_str);
    free(encode_str);

    return 0;
}

static int s7plc_handle_request(s7_plc_var_t *var, plc_client_t *client, pp_regulate_signal_t *regulate_data)
{

    int ret = 0;
    cJSON *root = cJSON_Parse(regulate_data->data);
    if (root)
    {
        char ops[8];
        cJSON *areas = NULL;
        cJSON *Area = NULL;
        int cnt = 0;
        TS7DataItem *Items = NULL;
        int i = 0;
        int area_id = -1;
        int Start = -1;
        int Amount = 0;
        int WordLen = 0;
        int DBNumber = 0;
        int res;
        int opcode = -1;
        char buf[1024];
        char buf_str[1024];

        GET_JSON_VALUE_STRING(root, "operation", ops);

        if (strcmp(ops, "read") == 0)
        {
            opcode = 0;
        }
        else if (strcmp(ops, "write") == 0)
        {
            opcode = 1;
        }
        else
        {
            goto out;
        }

        areas = cJSON_GetObjectItem(root, "areas");
        if (!areas)
        {
            dy_syslog(LOG_ERR, "get areas error");
            ret = -1;
            goto out;
        }

        cnt = cJSON_GetArraySize(areas);

        Items = calloc(cnt, sizeof(TS7DataItem));

        for (i = 0; i < cnt; i++)
        {
            DBNumber = 0;
            Area = cJSON_GetArrayItem(areas, i);
            if (Area)
            {
                GET_JSON_VALUE_INT(Area, "Area", area_id);
                GET_JSON_VALUE_INT(Area, "Start", Start);
                GET_JSON_VALUE_INT(Area, "Amount", Amount);
                GET_JSON_VALUE_INT(Area, "DBNumber", DBNumber);
                GET_JSON_VALUE_INT(Area, "WordLen", WordLen);

                if (area_id == S7AreaDB && DBNumber == -1)
                {
                    dy_syslog(LOG_ERR, "need set the DB number when operate the DB Area");
                }

                Items[i].Area = area_id;
                Items[i].WordLen = WordLen;
                Items[i].DBNumber = DBNumber;
                Items[i].Start = Start;
                Items[i].Amount = Amount;
                if (opcode == 0)
                {
                    Items[i].pdata = calloc(Amount, get_size_from_wordlen(WordLen));
                }
                else
                {
                    Items[i].pdata = calloc(Amount, get_size_from_wordlen(WordLen));
                    GET_JSON_VALUE_STRING(Area, "Data", buf_str);
                    int len = b64_decode(buf_str, buf, 1024);
                    if (len <= (Amount * get_size_from_wordlen(WordLen)))
                    {
                        memcpy(Items[i].pdata, buf, len);
                    }
                }
            }
        }

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

            dy_syslog(LOG_INFO, "Recevied interval ctrl cmd periad %d", regulate_data->period);
            snprintf(once.key, 128, "%s_%s", regulate_data->sn, regulate_data->src_identifier);
            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 = (unsigned char *)Items;
            once.rs->len = cnt;
            memcpy(once.rs->term_addr, &client, sizeof(client));

            period_msg_add(var->period, &once);
        }
        else
        {
            s7plc_send_request(var, client, regulate_data, Items, cnt, opcode);
            for (i = 0; i < cnt; i++)
            {
                free(Items[i].pdata);
            }
            free(Items);
        }
    }
    else
    {
        dy_syslog(LOG_ERR, "json parse error");
        ret = -1;
    }

out:
    cJSON_Delete(root);

    return ret;
}

static int s7_plc_msg_tag_data_ctrl(s7_plc_var_t *var, ipc_msg_t *mqtt_msg)
{
    pp_regulate_signal_t regulate_data;
    int ret = 0;
    dev_node_t *dev = NULL;
    int find = 0;

    ret = get_rglt_data_from_json(&regulate_data, mqtt_msg->payload);
    if (ret < 0 || regulate_data.len <= 0)
    {
        dy_syslog(LOG_ERR, "parse real data structure failed");
        return -1;
    }

    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;
    }

    s7plc_handle_request(var, dev->client, &regulate_data);

out:
    free(regulate_data.data);

    return 0;
}

static int s7_plc_mqtt_handle_recv_msg(void *obj, ipc_msg_t *mqtt_msg)
{
    s7_plc_var_t *var = (s7_plc_var_t *)obj;
    dy_syslog(LOG_DEBUG, "received MQTT topic:%s payload length:%d", mqtt_msg->topic, mqtt_msg->payloadLen);
    if (strstr(mqtt_msg->topic, TOPIC_SEND_RGLT_SIGNAL_RAW_DATA))
    {
        //发送下行数据
        s7_plc_msg_tag_data_ctrl(var, mqtt_msg);
    }
}

static void s7_plc_subscribe_all(s7_plc_var_t *var)
{
    ipc_session_t *ipc_session = var->session;
    char topic[TOPIC_MAX_LEN] = {0};
    int i = 0;

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

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

    snprintf(clientId, MAX_CLIENT_ID_LEN, "INT_s7_plc_%s", var->sn_str);
    var->session = ipc_session_new(clientId, (void*)var, IPC_DEFAULT);
    if (var->session == NULL)
        return -1;
    ipc_session_set_callbacks(var->session, s7_plc_mqtt_handle_recv_msg, NULL);
    s7_plc_subscribe_all(var);
    ipc_session_start(var->session);
}

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

    TIMER_CONFIRM(var->node_status_timer);

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

        dy_syslog(LOG_INFO, "online_status_changed:%d strlen(node_status):%d, (now-lasttime):%d",
                  var->ds->online_status_changed, strlen(node_status), now - lasttime);
        if (node_status && strlen(node_status) > 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, S7_PLC_NODE_STATUS_BAK_FILE);
        }
        if (node_status)
        {
            write_file_data(NODES_CACHE "/s7_plc_car_status.json", node_status, strlen(node_status));
            free(node_status);
        }
    }
}

static int period_service_poll(s7_plc_var_t *var)
{
    TIMER_CONFIRM(var->period_poll_timer);
    int ret = 0;
    static int all_started = 0;
    time_t now = time(NULL);
    static time_t last_trigger_ts = 0;
    static time_t last_query = 0;
    plc_client_t *client;

    if (var->upgrading == 1)
    {
        return 0;
    }

    while (1)
    {
        period_msg_t *msg = period_msg_next(var->period);
        if (msg && msg->started)
        {
            memcpy(&client, msg->rs->term_addr, sizeof(client));
            s7plc_send_request(var, client, msg->rs, (TS7DataItem *)msg->rs->data, msg->rs->len, 0);
        }
        else
        {
            break;
        }
    }

    if (now - last_trigger_ts > 20)
    {
        if (all_started == 0)
        {
            all_started = period_services_retrigger(var->session, var->period, var->nodes_cfg_table, var->template_table, ports, ARRAY_SIZE(ports));
        }
        last_trigger_ts = now;
    }
    return 0;
}

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

    while (1)
    {
        SELECT_INIT();
        SELECT_ADD_FD(var->node_status_timer);
        SELECT_ADD_FD(var->period_poll_timer);

        timeout.tv_usec = 0;
        timeout.tv_sec = 5;

        ret = select(maxfd + 1, &rset, 0, 0, &timeout);
        if (ret < 0)
        {
            dy_syslog(LOG_INFO, "errno %d\n", errno);

            if (errno == EINTR)
            {
                continue;
            }
            else
            {
                break;
            }
        }
        else if (ret > 0)
        {
            if (var->node_status_timer > 0 && FD_ISSET(var->node_status_timer, &rset))
            {
                FD_CLR(var->node_status_timer, &rset);
                ds_flush_timer(var);
            }
            if (var->period_poll_timer > 0 && FD_ISSET(var->period_poll_timer, &rset))
            {
                FD_CLR(var->period_poll_timer, &rset);
                period_service_poll(var);
            }
        }
    }
}

static int s7_plc_create_conn(s7_plc_var_t *var)
{
    int i = 0;
    int cnt = 0;
    int j = 0;
    int Requested, Negotiated, res;
    int pdu_size = 1024;
    plc_client_t *client = NULL;
    dev_node_t *node = NULL;

    for (i = 0; i < var->nodes_cfg_table->node_cnt; i++)
    {
        if (var->nodes_cfg_table->node[i].port == S7_PLC)
        {

            if (var->nodes_cfg_table->node[i].port == S7_PLC)
            {
                if (var->nodes_cfg_table->node[i].ext_data)
                {
                    cJSON *ext = cJSON_Parse(var->nodes_cfg_table->node[i].ext_data);
                    if (ext)
                    {
                        char ip_addr[64] = {0};
                        int rack = -1, slot = -1;
                        int find = 0;

                        GET_JSON_VALUE_STRING(ext, "ip_addr", ip_addr);
                        GET_JSON_VALUE_INT(ext, "rack", rack);
                        GET_JSON_VALUE_INT(ext, "slot", slot);

                        if (strlen(ip_addr) == 0 || rack == -1 || slot == -1)
                        {
                            dy_syslog(LOG_ERR, "wrong ext data:%s", var->nodes_cfg_table->node[i].ext_data);
                            continue;
                        }

                        list_for_each_entry(client, &var->connect_list, list)
                        {
                            if (strcmp(ip_addr, client->ip_addr) == 0 && rack == client->rack && slot == client->slot)
                            {
                                find = 1;
                                break;
                            }
                        }

                        if (find == 0)
                        {
                            client = calloc(1, sizeof(plc_client_t));
                            if (client == NULL)
                            {
                                dy_syslog(LOG_ERR, "malloc failed");
                                continue;
                            }
                            client->ip_addr = strdup(ip_addr);
                            client->rack = rack;
                            client->slot = slot;
                            client->Client = Cli_Create();
                            Cli_SetParam(client->Client, p_i32_PDURequest, (void *)&pdu_size);
                            res = Cli_ConnectTo(client->Client, client->ip_addr, client->rack, client->slot);
                            if (Check(res, "UNIT Connection"))
                            {
                                Cli_GetPduLength(client->Client, &Requested, &Negotiated);
                                dy_syslog(LOG_DEBUG, "connect to %s %d %d", client->ip_addr, client->rack, client->slot);
                                dy_syslog(LOG_DEBUG, "PDU Requested  : %d bytes Negotiated : %d bytes", Requested, Negotiated);
                            };

                            list_add_tail(&client->list, &var->connect_list);
                        }

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

    return 0;
}

static int s7_plc_init(s7_plc_var_t *var)
{
    check_make_dir(NODES_CACHE);
    check_make_dir(NODES_CFG);
    get_board_sn(var->sn_str);

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

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

    node_sta_recovery(var->ds, S7_PLC_NODE_STATUS_BAK_FILE);

    kv_array_init(&var->identifier_backup, 32);

    var->node_status_timer = my_timer_create();
    if (var->node_status_timer > 0)
    {
        my_timer_set(var->node_status_timer, 1, 10000);
    }

    var->period_poll_timer = my_timer_create();
    if (var->period_poll_timer > 0)
    {
        my_timer_set(var->period_poll_timer, 1, 100);
    }

    s7_plc_mqtt_client_init(var);
    period_msg_init(&var->period, 512);
    period_services_trigger(var->session, var->period, var->nodes_cfg_table, var->template_table, ports, ARRAY_SIZE(ports));

    s7_plc_create_conn(var);
}

int main(int argc, char *argv[])
{
    s7_plc_var_t var = {0};

    openlog("s7_plc", LOG_PID, LOG_DAEMON);

    memset(&var, 0, sizeof(var));
    s7_plc_init(&var);
    s7_plc_loop(&var);

    return 0;
}
