#include <errno.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <sys/select.h>
#include <sys/stat.h>
#include <sys/syscall.h>
#include <sys/time.h>
#include <sys/types.h>
#include <unistd.h>

#include "common_socket.h"

char heartbeat[64];
common_socket_var_t gvar = {0};

static int get_rt_data_from_json_str(pp_real_time_data_t *real_data, char *json_str)
{
    int ret = 0;
    char *b64_buf = NULL;

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

        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, "term_addr_b64", tmp);
        if (tmp != NULL)
        {
            len = b64_decode(tmp, buf, MAXBUF);
            ASSERT(len == sizeof(real_data->term_addr));
            memcpy(real_data->term_addr, buf, sizeof(real_data->term_addr));
            free(tmp);
        }

        GET_JSON_VALUE_DY_STRING(root, "sub_term_addr_b64", tmp);
        if (tmp != NULL)
        {
            len = b64_decode(tmp, buf, MAXBUF);
            ASSERT(len == sizeof(real_data->sub_term_addr));
            memcpy(real_data->sub_term_addr, buf, sizeof(real_data->sub_term_addr));
            free(tmp);
        }

        GET_JSON_VALUE_DY_STRING(root, "data_b64", tmp);
        if (tmp != NULL)
        {
            b64_buf = (char *)malloc(real_data->len * 2);
            if (b64_buf == NULL)
            {
                dy_syslog(LOG_ERR, "malloc failed");
                ret = -1;
                free(buf);
                goto out;
            }
            len = b64_decode(tmp, b64_buf, real_data->len * 2 + 1024);
            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);
            free(b64_buf);
            goto out;
        }

        memcpy(real_data->data, b64_buf, len);
        free(buf);
        free(b64_buf);
    }
    else
    {
        dy_syslog(LOG_ERR, "json parse error");
        ret = -1;
    }

out:
    cJSON_Delete(root);

    return ret;
}

static int common_socket_heart_beat(common_socket_var_t *var)
{
    int ret;
    time_t now = time(NULL);
    socket_state_e state;

    TIMER_CONFIRM(var->heartbeat_timer);

    connect_config_t *connect_cfg = NULL;
    list_for_each_entry(connect_cfg, &var->connect_list, list)
    {
        if (connect_cfg == NULL || connect_cfg->socket_client == NULL)
            continue;
        state = dy_socket_session_get_state(connect_cfg->socket_client);
        if (state != SOCKET_CONNECTED)
            continue;
        if (connect_cfg->last_send == 0 && connect_cfg->common_cfg.login_packet_len != 0)
        {
            connect_cfg->last_send = 1;
            dy_syslog_hex(LOG_DEBUG, connect_cfg->common_cfg.login_packet, connect_cfg->common_cfg.login_packet_len, "len %d			网关->[%s]		", connect_cfg->common_cfg.login_packet_len, connect_cfg->addr);
            dy_socket_session_snd(connect_cfg->socket_client, connect_cfg->common_cfg.login_packet, connect_cfg->common_cfg.login_packet_len);
            sleep(2);
        }
        if (connect_cfg->common_cfg.heartbeat_interval)
        {
            ret = check_timeout_second(now, &connect_cfg->last_send, connect_cfg->common_cfg.heartbeat_interval);
            //dy_syslog(LOG_DEBUG, "ret %d", ret);
            if (ret)
            {
                dy_syslog_hex(LOG_DEBUG, connect_cfg->common_cfg.heartbeat_packet, connect_cfg->common_cfg.heartbeat_packet_len, "len %d			网关->[%s]		", connect_cfg->common_cfg.heartbeat_packet_len, connect_cfg->addr);
                dy_socket_session_snd(connect_cfg->socket_client, connect_cfg->common_cfg.heartbeat_packet, connect_cfg->common_cfg.heartbeat_packet_len);
            }
        }
    }

    return 0;
}

static void common_socket_loop(common_socket_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)
        {
            int error = errno;
            printf("errno %d\n", error);
            if (error == 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);
                common_socket_heart_beat(var);
            }
        }
    }
}

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

    snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/+/device/+/data/%s", var->sn_str, TOPIC_UP_RAW_DATA);
    ipc_session_subscribe(session, topic);
}

static int common_socket_mqtt_handle_recv_msg(void *obj, ipc_msg_t *mqtt_msg)
{
    common_socket_var_t *var = (common_socket_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_UP_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 (real_data.port == connect_cfg->common_cfg.binding_port)
            {
                //dy_syslog_hex(LOG_DEBUG, real_data.data, real_data.len, "send tcp data len:%d", real_data.len);
                if (connect_cfg->last_send == 0 && connect_cfg->common_cfg.login_packet_len != 0)
                {
                    dy_syslog_hex(LOG_DEBUG, connect_cfg->common_cfg.login_packet, connect_cfg->common_cfg.login_packet_len, "len %d			网关->[%s]		", connect_cfg->common_cfg.login_packet_len, connect_cfg->addr);
                    dy_socket_session_snd(connect_cfg->socket_client, connect_cfg->common_cfg.login_packet, connect_cfg->common_cfg.login_packet_len);
                }
                char *data = calloc(1, real_data.len + connect_cfg->common_cfg.frame_header_len + connect_cfg->common_cfg.frame_tail_len);
                if (data)
                {
                    memcpy(data, connect_cfg->common_cfg.frame_header, connect_cfg->common_cfg.frame_header_len);
                    memcpy(data + connect_cfg->common_cfg.frame_header_len, real_data.data, real_data.len);
                    memcpy(data + connect_cfg->common_cfg.frame_header_len + real_data.len, connect_cfg->common_cfg.frame_tail, connect_cfg->common_cfg.frame_tail_len);
                    dy_syslog_hex(LOG_DEBUG, data, real_data.len + connect_cfg->common_cfg.frame_header_len + connect_cfg->common_cfg.frame_tail_len, "len %d			网关->[%s]		", real_data.len + connect_cfg->common_cfg.frame_header_len + connect_cfg->common_cfg.frame_tail_len, connect_cfg->addr);
                    dy_socket_session_snd(connect_cfg->socket_client, data, real_data.len + connect_cfg->common_cfg.frame_header_len + connect_cfg->common_cfg.frame_tail_len);
                    free(data);
                }
            }
            if (real_data.port == connect_cfg->mqtt_cfg.binding_port)
            {
                mqtt_session_publish(connect_cfg->session, connect_cfg->mqtt_cfg.pub_topic, real_data.data, real_data.len);
            }
        }

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

static int send_data_to_port(common_socket_var_t *var, gw_port_e port, unsigned char *buf, int len)
{
    cJSON *rs_data = cJSON_CreateObject();
    int mi = time(NULL);
    char topic[TOPIC_MAX_LEN] = {0};
    static char *data = NULL;

    if (data == NULL)
    {
        data = malloc(4096);
    }
    b64_encode(buf, len, data, 4096);
    cJSON_AddStringToObject(rs_data, "raw_data", data);

    cJSON_AddStringToObject(rs_data, "data_b64", data);
    cJSON_AddNumberToObject(rs_data, "len", len);
    cJSON_AddNumberToObject(rs_data, "period", 0);
    cJSON_AddStringToObject(rs_data, "port", port_enum2char(port));
    cJSON_AddNumberToObject(rs_data, "mi", mi);
    cJSON_AddStringToObject(rs_data, "src_identifier", "transparent");
    cJSON_AddNumberToObject(rs_data, "communication_timeout", 0);

    char *data_tmp = cJSON_Print(rs_data);
    dy_syslog(LOG_DEBUG, "data_tmp:%s", data_tmp);

    snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/%s/device/noneed/data/%s", var->sn_str, port_enum2char(port),
             TOPIC_SEND_RGLT_SIGNAL_RAW_DATA);
    ipc_session_publish(var->session, topic, data_tmp, strlen(data_tmp));
    // 回收资源
    free(data_tmp);
    cJSON_Delete(rs_data);
}

static int external_mqtt_handle_recv_msg(void *obj, mqtt_message_t *mqtt_msg)
{
    connect_config_t *connect_cfg = (connect_config_t *)obj;
    common_socket_var_t *var = connect_cfg->var;

    // send to port
    send_data_to_port(var, connect_cfg->mqtt_cfg.binding_port, mqtt_msg->payload, mqtt_msg->payloadLen);
}

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

    snprintf(clientId, MAX_CLIENT_ID_LEN, "INT_COMMON_SOCKET_%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, common_socket_mqtt_handle_recv_msg, NULL);
    common_socket_subscribe_all(var);
    ipc_session_start(var->session);
}

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

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

    if (len == connect_cfg->common_cfg.login_packet_len && connect_cfg->common_cfg.login_packet && memcmp(connect_cfg->common_cfg.login_packet, buf, len) == 0)
    {
        dy_syslog(LOG_DEBUG, "login packet...");
        return 0;
    }

    if (len == connect_cfg->common_cfg.heartbeat_packet_len && connect_cfg->common_cfg.heartbeat_packet && memcmp(connect_cfg->common_cfg.heartbeat_packet, buf, len) == 0)
    {
        dy_syslog(LOG_DEBUG, "heartbeat packet...");
        return 0;
    }

    send_data_to_port(var, connect_cfg->common_cfg.binding_port, buf, len);
}

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

    dy_syslog(LOG_DEBUG, "connect_state_change state %d", state);
    if (state == SOCKET_DISCONNECTED)
        connect_cfg->last_send = 0;

    return 0;
}

int del_space(char *src)
{
    char *pTmp = src;
    unsigned int iSpace = 0;

    while (*src != '\0')
    {
        if (*src != ' ')
        {
            *pTmp++ = *src;
        }
        else
        {
            iSpace++;
        }

        src++;
    }

    *pTmp = '\0';
    return iSpace;
}

static int common_socket_load_conf(common_socket_var_t *var, char *file)
{
    char buff[128] = {0};
    int ret = -1;
    int i = 0;

    char *json_str = NULL;
    json_str = read_file_data(file);
    if (!json_str)
    {
        return -1;
    }

    cJSON *root = cJSON_Parse(json_str);
    if (!root)
    {
        return -1;
    }

    cJSON *nodes = cJSON_GetObjectItem(root, "common_socket_cfg");
    if (!nodes)
    {
        ret = -1;
        goto __cleanup;
    }

    int size = cJSON_GetArraySize(nodes);

    for (i = 0; i < size; i++)
    {
        cJSON *node = cJSON_GetArrayItem(nodes, i);
        if (!node)
        {
            ret = i + 1;
            goto __cleanup;
        }

        connect_config_t *connect_cfg = calloc(1, sizeof(connect_config_t));

        connect_cfg->var = var;
        GET_JSON_VALUE_STRING(node, "protocol_type", buff);
        connect_cfg->type = socket_type_char2enum(buff);
        if (connect_cfg->type == SOCKET_MQTT)
        {
            char addr[MAX_HOST_LEN]; ///<MQTT地址
            unsigned short port;     ///<MQTT端口
            char user[MAX_USER_LEN]; ///<MQTT
            char pass[MAX_PASS_LEN];
            char clientId[MAX_CLIENT_ID_LEN];

            memset(clientId, 0, sizeof(clientId));
            GET_JSON_VALUE_STRING(node, "client_id", clientId);
            if (clientId[0] == '\0')
                snprintf(clientId, MAX_CLIENT_ID_LEN, "%s_TRS_%d", var->sn_str, i);
            connect_cfg->session = mqtt_session_new(clientId, (void *)connect_cfg);
            if (connect_cfg->session == NULL)
            {
                continue;
            }
            GET_JSON_VALUE_STRING(node, "binding_port", buff);
            connect_cfg->mqtt_cfg.binding_port = port_char2enum(buff);

            GET_JSON_VALUE_STRING(node, "host", addr);
            GET_JSON_VALUE_INT(node, "port", port);
            GET_JSON_VALUE_STRING(node, "user", user);
            GET_JSON_VALUE_STRING(node, "pass", pass);

            char cafile_url[128] = {0}, certfile_url[128] = {0}, keyfile_url[128] = {0};
            GET_JSON_VALUE_STRING(node, "cafile", cafile_url);
            GET_JSON_VALUE_STRING(node, "certfile", certfile_url);
            GET_JSON_VALUE_STRING(node, "keyfile", keyfile_url);
            if (strlen(cafile_url) && strlen(certfile_url) && strlen(keyfile_url))
            {
                char cafile[128] = {0};
                char certfile[128] = {0};
                char keyfile[128] = {0};
                char path[128] = {0};
                sprintf(path, "%s/%d", COMMON_SOCKET_TLS_FILE_PATH, i);
                if (strlen(cafile_url))
                {
                    mqtt_tls_file_overlow(cafile_url, path, cafile, 128);
                }
                if (strlen(certfile_url))
                {
                    mqtt_tls_file_overlow(certfile_url, path, certfile, 128);
                }
                if (strlen(keyfile_url))
                {
                    mqtt_tls_file_overlow(keyfile_url, path, keyfile, 128);
                }
                dy_syslog(LOG_DEBUG, "cafile:%s,certfile:%s,keyfile:%s", cafile, certfile, keyfile);
                mqtt_session_set_tls(connect_cfg->session, cafile, certfile, keyfile);
            }

            mqtt_session_set_address(connect_cfg->session, addr, port, user, pass);
            mqtt_session_set_opts(connect_cfg->session, DEFAULT_QOS, KEEPALIVEINTERVAL);
            GET_JSON_VALUE_DY_STRING(node, "pub_topic", connect_cfg->mqtt_cfg.pub_topic);
            GET_JSON_VALUE_DY_STRING(node, "sub_topic", connect_cfg->mqtt_cfg.sub_topic);
        }
        else
        {
            GET_JSON_VALUE_DY_STRING(node, "socket_addr", connect_cfg->addr);
            GET_JSON_VALUE_INT(node, "socket_port", connect_cfg->port);
            GET_JSON_VALUE_STRING(node, "binding_port", buff);
            connect_cfg->common_cfg.binding_port = port_char2enum(buff);

            GET_JSON_VALUE_STRING(node, "frame_header", buff);
            del_space(buff);
            connect_cfg->common_cfg.frame_header_len = strlen(buff) / 2;
            if (connect_cfg->common_cfg.frame_header_len > 0)
            {
                connect_cfg->common_cfg.frame_header = calloc(1, connect_cfg->common_cfg.frame_header_len);
                str2hex(connect_cfg->common_cfg.frame_header, buff, connect_cfg->common_cfg.frame_header_len);
            }

            GET_JSON_VALUE_STRING(node, "frame_tail", buff);
            del_space(buff);
            connect_cfg->common_cfg.frame_tail_len = strlen(buff) / 2;
            if (connect_cfg->common_cfg.frame_tail_len > 0)
            {
                connect_cfg->common_cfg.frame_tail = calloc(1, connect_cfg->common_cfg.frame_tail_len);
                str2hex(connect_cfg->common_cfg.frame_tail, buff, connect_cfg->common_cfg.frame_tail_len);
            }

            GET_JSON_VALUE_STRING(node, "login_packet", buff);
            del_space(buff);
            connect_cfg->common_cfg.login_packet_len = strlen(buff) / 2;
            if (connect_cfg->common_cfg.login_packet_len > 0)
            {
                connect_cfg->common_cfg.login_packet = calloc(1, connect_cfg->common_cfg.login_packet_len);
                str2hex(connect_cfg->common_cfg.login_packet, buff, connect_cfg->common_cfg.login_packet_len);
            }

            GET_JSON_VALUE_STRING(node, "heartbeat_packet", buff);
            del_space(buff);
            connect_cfg->common_cfg.heartbeat_packet_len = strlen(buff) / 2;
            if (connect_cfg->common_cfg.heartbeat_packet_len > 0)
            {
                connect_cfg->common_cfg.heartbeat_packet = calloc(1, connect_cfg->common_cfg.heartbeat_packet_len);
                str2hex(connect_cfg->common_cfg.heartbeat_packet, buff, connect_cfg->common_cfg.heartbeat_packet_len);
            }

            GET_JSON_VALUE_INT(node, "heartbeat_interval", connect_cfg->common_cfg.heartbeat_interval);
        }

        dy_syslog(LOG_INFO, "==%d %d socket_addr:%s port:%d protocol_type:%d==", size, i, connect_cfg->addr, connect_cfg->port, connect_cfg->type);
        list_add_tail(&connect_cfg->list, &var->connect_list);
    }
    ret = 0;

__cleanup:
    cJSON_Delete(root);
    return ret;
}

static int common_socket_init(common_socket_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");
    }

    common_socket_mqtt_client_init(var);

    if (access(COMMON_SOCKET_CFG, F_OK) != -1)
    {
        common_socket_load_conf(var, COMMON_SOCKET_CFG);
    }

    connect_config_t *connect_cfg = NULL;
    list_for_each_entry(connect_cfg, &var->connect_list, list)
    {
        var->node_cnt++;
        if (connect_cfg->type == SOCKET_MQTT)
        {
            char topic[TOPIC_MAX_LEN] = {0};

            mqtt_session_set_callbacks(connect_cfg->session, external_mqtt_handle_recv_msg, NULL);
            mqtt_session_subscribe(connect_cfg->session, connect_cfg->mqtt_cfg.sub_topic);
            mqtt_session_start(connect_cfg->session);
        }
        else
        {
            connect_cfg->socket_client = dy_socket_session_init(connect_cfg->type, connect_cfg->addr,
                                                                connect_cfg->port,
                                                                socket_handle_rcv_msg,
                                                                connect_state_change,
                                                                (void *)connect_cfg);



        }
    }

    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[])
{
    common_socket_var_t *var = &gvar;

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

    common_socket_init(var);
    common_socket_loop(var);

    return 0;
}
