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


gw_port_e ports[] = {TCP_CLIENT};
static int get_rt_data_from_json_str(pp_real_time_data_t *real_data, char *json_str)
{
    int ret = 0;

    cJSON *root = cJSON_Parse(json_str);
    if (root)
    {
        char *payload = NULL;
        char *buf = NULL;
        int len;
        char *tmp;

        buf = (char *)malloc(MAXBUF);
        if (buf == NULL)
        {
            dy_syslog(LOG_ERR, "malloc failed");
            ret = -1;
            goto out;
        }

        GET_JSON_VALUE_STRING(root, "port", buf);
        real_data->port = port_char2enum(buf);
        GET_JSON_VALUE_INT(root, "instruction_code", real_data->instruction_code);
        GET_JSON_VALUE_INT(root, "mi", real_data->mi);
        GET_JSON_VALUE_INT(root, "ts", real_data->ts);
        GET_JSON_VALUE_INT(root, "len", real_data->len);
        GET_JSON_VALUE_INT(root, "sub_term_addr_len", real_data->sub_term_addr_len);
        GET_JSON_VALUE_INT(root, "start_addr", real_data->start_addr);
        GET_JSON_VALUE_STRING(root, "sub_template_id", real_data->sub_template_id);
        GET_JSON_VALUE_STRING(root, "src_identifier", real_data->src_identifier);
        GET_JSON_VALUE_STRING(root, "sn", real_data->sn);
        GET_JSON_VALUE_STRING(root, "dtu_sn", real_data->dtu_sn);

        GET_JSON_VALUE_DY_STRING(root, "data_b64", tmp);
        len = b64_decode(tmp, buf, MAXBUF);
        free(tmp);
        //dy_syslog(LOG_DEBUG, "len:%d, real_data->len:%d", len, real_data->len);
        ASSERT(len == real_data->len);
        real_data->data = malloc(len);
        if (real_data->data == NULL)
        {
            dy_syslog(LOG_ERR, "malloc fail len:%d", len);
            free(buf);
            goto out;
        }
        memcpy(real_data->data, buf, len);
        free(buf);
    }
    else
    {
        dy_syslog(LOG_ERR, "json parse error");
        ret = -1;
    }

out:
    cJSON_Delete(root);

    return ret;
}

static int tcp_client_heartbeat_timer(tcp_client_var_t *var)
{
    TIMER_CONFIRM(var->heartbeat_timer);

    return 0;
}

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

    while (1)
    {
        SELECT_INIT();
        SELECT_ADD_FD(var->heartbeat_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);
                tcp_client_heartbeat_timer(var);
            }
        }
    }
}

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

    for (i = 0; i < ARRAY_SIZE(ports); i++)
    {
        snprintf(topic, TOPIC_MAX_LEN, "ipc/+/%s/device/+/data/%s", port_enum2char(ports[i]), TOPIC_SEND_RGLT_SIGNAL_RAW_DATA);
        ipc_session_subscribe(ipc_session, topic);
    }
}

static int tcp_client_mqtt_handle_recv_msg(void *obj, ipc_msg_t *mqtt_msg)
{
    tcp_client_var_t *var = (tcp_client_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))
    {
        int ret = 0;
        pp_real_time_data_t real_data = {0};

        ret = get_rt_data_from_json_str(&real_data, mqtt_msg->payload);
        if (ret < 0)
        {
            dy_syslog(LOG_ERR, "parse real data structure failed");
        }

        connect_config_t *connect_cfg = NULL;
        list_for_each_entry(connect_cfg, &var->connect_list, list)
        {
            if (strcmp(connect_cfg->sn, real_data.sn) == 0)
            {
                dy_syslog_hex(LOG_DEBUG, real_data.data, real_data.len, "len %d			网关->[%s]		", real_data.len, connect_cfg->addr);
                dy_socket_session_snd(connect_cfg->socket_client, real_data.data, real_data.len);

                //store the identifier
                {
                    char key_str[64] = {0};
                    regulate_cmd_backup_t fields;

                    strcpy(fields.src_identifier, real_data.src_identifier);
                    strcpy(fields.sn, real_data.sn);
                    fields.mi = real_data.mi;

                    sprintf(key_str, "%s", real_data.sn);
                    var->identifier_backup.insert(&var->identifier_backup, key_str, &fields, sizeof(fields));
                }
                flush_device_status_by_sn(var->ds, ACTION_SND, real_data.sn);
            }
        }

        if (real_data.data)
        {
            free(real_data.data);
        }
    }
}

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

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

    ipc_session_set_callbacks(var->session, tcp_client_mqtt_handle_recv_msg, NULL);
    tcp_client_subscribe_all(var);
    ipc_session_start(var->session);
}

int tcp_client_data_report(tcp_client_var_t *var, char *sn, char *data, int len, unsigned char port)
{
    pp_real_time_data_t real_data = {0};
    regulate_cmd_backup_t *fields = NULL;

    //restore the identifier
    {
        char key_str[64] = {0};
        data_info_t *data_info = NULL;

        sprintf(key_str, "%s", sn);

        data_info = var->identifier_backup.get(&var->identifier_backup, key_str);
        if (data_info)
        {
            fields = (regulate_cmd_backup_t *)data_info->data;
        }
    }
    real_data.port = port;
    real_data.len = len;
    real_data.instruction_code = data[1];
    real_data.data = malloc(real_data.len);
    if (fields != NULL)
    {
        real_data.mi = fields->mi;
        strcpy(real_data.sn, fields->sn);
        strncpy(real_data.src_identifier, fields->src_identifier, sizeof(real_data.src_identifier));
    }

    memcpy(real_data.data, data, real_data.len);

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

    free(real_data.data);

    return 0;
}

static int socket_handle_rcv_msg(void *obj, char *buf, int len)
{
    connect_config_t *connect_cfg = (connect_config_t*)obj;
    tcp_client_var_t *var = connect_cfg->var;
    int mi = time(NULL);

    dy_syslog_hex(LOG_DEBUG, buf, len, "len %d			网关<-[%s]		", len, connect_cfg->addr);
    tcp_client_data_report(var, connect_cfg->sn, buf, len, TCP_CLIENT);

    flush_device_status_by_sn(var->ds, ACTION_SND, connect_cfg->sn);
    return 0;
}

static int connect_state_change(void *obj, socket_state_e state)
{
    connect_config_t *connect_cfg = (connect_config_t*)obj;
    tcp_client_var_t *var = connect_cfg->var;

    dy_syslog(LOG_DEBUG, "connect_state_change state %d", state);

    return 0;
}

int tcp_client_load_tcp_nodes(tcp_client_var_t *var)
{
    int i;
    connect_config_t *connect_cfg = NULL;

    for (i = 0; i < var->nodes_cfg_table->node_cnt; i++)
    {
        if (strlen(var->nodes_cfg_table->node[i].app_key) == 0 && var->nodes_cfg_table->node[i].port == TCP_CLIENT && strlen(var->nodes_cfg_table->node[i].tcp_ip_addr) > 5)
        {
            connect_cfg = calloc(1, sizeof(connect_config_t));
            connect_cfg->var = var;
            strcpy(connect_cfg->sn, var->nodes_cfg_table->node[i].sn);
            connect_cfg->addr = strdup(var->nodes_cfg_table->node[i].tcp_ip_addr);
            connect_cfg->port = var->nodes_cfg_table->node[i].tcp_port;

            connect_cfg->socket_client = dy_socket_session_init(SOCKET_TCP, var->nodes_cfg_table->node[i].tcp_ip_addr,
                                                                var->nodes_cfg_table->node[i].tcp_port,
                                                                socket_handle_rcv_msg,
                                                                connect_state_change,
                                                                (void *)connect_cfg);

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

static int tcp_client_init(tcp_client_var_t *var)
{
    INIT_LIST_HEAD(&var->connect_list);

    get_board_sn(var->sn_str);
    get_process_name(var->proc_name);

    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, TCP_CLIENT_NODE_STATUS_BAK_FILE);

    kv_array_init(&var->identifier_backup, 32);

    tcp_client_mqtt_client_init(var);
    tcp_client_load_tcp_nodes(var);

    var->heartbeat_timer = my_timer_create();
    if (var->heartbeat_timer > 0)
    {
        if (var->node_cnt)
            my_timer_set(var->heartbeat_timer, 3, 1000);
        else
            my_timer_set(var->heartbeat_timer, 3, 10000);
    }

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

    return 0;
}

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

    openlog("TCP-CLIENT", LOG_PID, LOG_DAEMON);

    memset(&var, 0, sizeof(tcp_client_var_t));

    tcp_client_init(&var);
    tcp_client_loop(&var);

    return 0;
}
