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

static int modbus_tcp_msg_tag_data_ctrl(modbus_tcp_thread_var_t *tcp_thread_var, ipc_msg_t *mqtt_msg)
{
    int ret;
    int find = 0;
    pp_regulate_signal_t regulate_data;
    modbus_tcp_data_list_t *modbus_tcp_data;
    modbus_tcp_data_list_t *tmp = NULL;
    modbus_tcp_data_list_t *tmp2 = NULL;

    //先检查参数
    ret = get_rglt_data_from_json(&regulate_data, mqtt_msg->payload);
    if (ret < 0 || regulate_data.len <= 0) {
        dbg_syslog(LOG_ERR, "parse real data structure failed");
        return -1;
    }

    modbus_tcp_data = calloc(sizeof(modbus_tcp_data_list_t), 1);
    if (modbus_tcp_data) {
        memcpy(&modbus_tcp_data->rs, &regulate_data, sizeof(regulate_data));
        if (regulate_data.period) {
            tcp_thread_var->period_triggered = 1;
            list_for_each_entry_safe(tmp, tmp2, &tcp_thread_var->data_list, list) {
                if (strcmp(regulate_data.sn, tmp->rs.sn) == 0 &&
                    regulate_data.period == tmp->rs.period &&
                    regulate_data.len == tmp->rs.len &&
                    memcmp(regulate_data.data, tmp->rs.data, tmp->rs.len) == 0)
                {
                    find = 1;
                    break;
                }
            }
            if (find == 0) {
                pthread_mutex_lock(&tcp_thread_var->data_lock);
                list_add_tail(&modbus_tcp_data->list, &tcp_thread_var->data_list);
                pthread_mutex_unlock(&tcp_thread_var->data_lock);
                dbg_syslog(LOG_DEBUG, "list_add_tail tcp_port(%d) tcp_ip_addr(%s) sn:%s\n", regulate_data.tcp_port, regulate_data.tcp_ip_addr, regulate_data.sn);
            }
        }
        else {
            pthread_mutex_lock(&tcp_thread_var->data_lock);
            list_add_tail(&modbus_tcp_data->list, &tcp_thread_var->data_list);
            pthread_mutex_unlock(&tcp_thread_var->data_lock);
            dbg_syslog(LOG_DEBUG, "list_add_tail tcp_port(%d) tcp_ip_addr(%s) sn:%s\n", regulate_data.tcp_port, regulate_data.tcp_ip_addr, regulate_data.sn);
        }

        dbg_syslog(LOG_DEBUG, "sn %s data[0] %d identifier %s port %d period %d", regulate_data.sn, regulate_data.data[0], regulate_data.src_identifier, MODBUS_TCP, regulate_data.period);
        if (find) {
            dbg_syslog(LOG_WARNING, "Repeated commands !!!");
            free(modbus_tcp_data);
            free(regulate_data.data);
        }
    }
    else {
        dbg_syslog(LOG_ERR, "modbus_tcp_data calloc failed!\n");
        free(regulate_data.data);
    }

    return 0;
}

static int modbus_tcp_mqtt_handle_recv_msg(void *obj, ipc_msg_t *mqtt_msg)
{
    modbus_tcp_thread_var_t *tcp_thread_var = (modbus_tcp_thread_t *)obj;
    dbg_syslog(LOG_INFO, "received topic  : %s", mqtt_msg->topic);
    dbg_syslog(LOG_INFO, "received payload: %s", mqtt_msg->payload);
    if (strstr(mqtt_msg->topic, TOPIC_SEND_RGLT_SIGNAL_RAW_DATA)) {
        modbus_tcp_msg_tag_data_ctrl(tcp_thread_var, mqtt_msg);     //发送下行数据
    }

    return 0;
}

static void modbus_tcp_subscribe_all(modbus_tcp_thread_var_t *tcp_thread_var)
{
    ipc_session_t *ipc_session = tcp_thread_var->session;
    char topic[TOPIC_MAX_LEN] = {0};
    tcp_node_t *tcp_node = NULL;
    tcp_node_t *tmp = NULL;

    list_for_each_entry_safe(tcp_node, tmp, &tcp_thread_var->node_list, list) {
        snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/%s/device/%s/data/%s",tcp_thread_var->var->sn_str, port_enum2char(MODBUS_TCP), tcp_node->node_cfg->sn, TOPIC_SEND_RGLT_SIGNAL_RAW_DATA);
        dbg_syslog(LOG_NOTICE, "sub topic: %s", topic);
        ipc_session_subscribe(ipc_session, topic);
    }
}

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

    snprintf(clientId, MAX_CLIENT_ID_LEN, "MT_%s_%s_%d",tcp_thread_var->var->sn_str,tcp_thread_var->tcp_ip_addr, tcp_thread_var->tcp_port);
    dbg_syslog(LOG_NOTICE, "clientId: %s", clientId);

    tcp_thread_var->session = ipc_session_new(clientId, (void *)tcp_thread_var, IPC_DEFAULT);
    if (tcp_thread_var->session == NULL) {
        dbg_syslog(LOG_ERR, "create mqtt session failed");
        return -1;
    }

    ipc_session_set_callbacks(tcp_thread_var->session, modbus_tcp_mqtt_handle_recv_msg, NULL);
    modbus_tcp_subscribe_all(tcp_thread_var);
    ipc_session_start(tcp_thread_var->session);

    return 0;
}

static int modbus_tcp_period_rs_triger(modbus_tcp_thread_var_t *tcp_thread_var)
{
    int i, j;
    modbus_tcp_var_t *var = tcp_thread_var->var;
    tcp_node_t *tcp_node = NULL;

    list_for_each_entry(tcp_node, &tcp_thread_var->node_list, list) {
        for (i = 0; i < var->template_table->template_cnt; i++) {
            if (strcmp(tcp_node->node_cfg->template_id, var->template_table->template[i].template_id)) {
                continue;
            }
            //解析周期命令
            for (j = 0; j < var->template_table->template[i].service_tab->serviceCnt; j++) {
                if (var->template_table->template[i].service_tab->service[j].server_period > 0 && 
                    !var->template_table->template[i].service_tab->service[j].direction)
                {
                    char *data = service_rs_data(&var->template_table->template[i].service_tab->service[j], tcp_node->node_cfg->sn);
                    char topic[TOPIC_MAX_LEN] = {0};
                    snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/%s/device/%s/data/%s", var->sn_str, port_enum2char(tcp_node->node_cfg->port), 
                            tcp_node->node_cfg->sn, TOPIC_EVT_SET_RGLT);
                    usleep(10 * 1000);
                    dbg_syslog(LOG_INFO, "pub topic  : %s", topic);
                    dbg_syslog(LOG_INFO, "pub payload: %s", data);
                    ipc_session_publish(tcp_thread_var->session, topic, data, strlen(data));
                    free(data);
                }
            }
            break;
        }
    }

    return 0;
}

static int modbus_tcp_tag_data_report(modbus_tcp_thread_var_t *tcp_thread_var, char *data, int len, unsigned char port, pp_regulate_signal_t *rs)
{
    pp_real_time_data_t real_data = {0};
    time_t now = time(NULL);

    real_data.port = port;
    real_data.len = len;
    real_data.instruction_code = data[1];
    real_data.data = malloc(real_data.len);
    real_data.mi = rs->mi;
    strcpy(real_data.sn, rs->sn);
    strncpy(real_data.src_identifier, rs->src_identifier, sizeof(real_data.src_identifier));
    strcpy(real_data.requester, rs->requester);
    memcpy(real_data.data, data, real_data.len);

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

    free(real_data.data);

    return 0;
}

static void modbus_tcp_bus_send(modbus_tcp_thread_var_t *tcp_thread_var, pp_regulate_signal_t *rs, unsigned char port)
{
    int ret;
    char *data = rs->data;
    int len = rs->len;

    flush_device_status_by_sn(tcp_thread_var->var->ds, ACTION_SND, rs->sn);

    if (port == MODBUS_TCP) {
        pthread_mutex_t * mutexp = NULL;
        // tcp发送数据,阻塞等待执行数据返回
        modbus_tcp_client_init(tcp_thread_var, &mutexp);

        if (tcp_thread_var->ctx && mutexp) {
            char read_nBytes[2048] = {0};

            ret = pthread_mutex_lock(mutexp);
            if (ret) {
                dbg_syslog(LOG_ERR, "Error, failed to lock mutex: %d", ret);
                return;
            }

            len = rs->len - 2;
            dbg_syslog_hex(LOG_NOTICE, data, len, "sn:%s, len %d, identifier:%s, GW->%s:%d, send_cnt[0x%llx]", rs->sn, len, rs->src_identifier, tcp_thread_var->tcp_ip_addr, tcp_thread_var->tcp_port, ++tcp_thread_var->send_cnt);
            int rc = modbus_communication_transparent(tcp_thread_var->ctx, data, len, read_nBytes);
            pthread_mutex_unlock(mutexp);
            if (rc > 0) {
                if (tcp_thread_var->try_cnt != 0) {
                    tcp_thread_var->try_cnt = 0;
                }
                dbg_syslog_hex(LOG_NOTICE, read_nBytes, rc, "sn:%s, len %d, GW<-%s:%d, recv_cnt[0x%llx]", rs->sn, rc, tcp_thread_var->tcp_ip_addr, tcp_thread_var->tcp_port, ++tcp_thread_var->recv_cnt);
                modbus_tcp_tag_data_report(tcp_thread_var, read_nBytes, rc, port, rs);
            }
            else {
                dbg_syslog(LOG_ERR, "ctx %p try_cnt %d rc %d, sn: %s", tcp_thread_var->ctx, tcp_thread_var->try_cnt, rc, rs->sn);
                if (++tcp_thread_var->try_cnt >= 8) {
                    modbus_freetcp(tcp_thread_var->tcp_ip_addr, tcp_thread_var->tcp_port, tcp_thread_var->ctx);
                    tcp_thread_var->ctx = NULL;
                    tcp_thread_var->try_cnt = 0;
                }
            }
        }
        else {
            dbg_syslog(LOG_WARNING, "tcp_ip_addr:%s port:%d ctx is NULL!!!", tcp_thread_var->tcp_ip_addr, tcp_thread_var->tcp_port);
        }
    }
}

static int send_poll(modbus_tcp_thread_var_t *tcp_thread_var, struct list_head *data_list, int port)
{
    int ret, find = 0;
    unsigned int communication_timeout = 200;
    modbus_tcp_data_list_t *send_data = NULL;
    modbus_tcp_data_list_t *tmp = NULL;
    long long now = uptime_stamp_ms();

    list_for_each_entry_safe(send_data, tmp, data_list, list) {
        if (!send_data->rs.period) {
            find = 1;
            pthread_mutex_lock(&tcp_thread_var->data_lock);
            list_del(&send_data->list);
            pthread_mutex_unlock(&tcp_thread_var->data_lock);
            modbus_tcp_bus_send(tcp_thread_var, &send_data->rs, port);
            free(send_data->rs.data);
            free(send_data);
        }
        else {
            if (tcp_thread_var->period_flag == 0) {
                continue;
            }
                
            ret = check_timeout_msecond(now, &send_data->last_send, send_data->rs.period);
            if (ret) {
                find = 1;
                modbus_tcp_bus_send(tcp_thread_var, &send_data->rs, port);
            }
        }
    }

    return find;
}

static int modbus_tcp_send_poll(modbus_tcp_thread_var_t *tcp_thread_var, int time_trigger)
{
    if (time_trigger) {
        TIMER_CONFIRM(tcp_thread_var->send_timer);
    }

    if (!tcp_thread_var->start_flag) {
        return 0;
    }

    return send_poll(tcp_thread_var, &tcp_thread_var->data_list, MODBUS_TCP);
}

static int modbus_tcp_get_period_cmd(modbus_tcp_thread_var_t *tcp_thread_var)
{
    TIMER_CONFIRM(tcp_thread_var->period_timer);

    if (tcp_thread_var->period_triggered == 0) {
        modbus_tcp_period_rs_triger(tcp_thread_var);
    }

    return 0;
}

static void *modbus_tcp_thread_loop(void *param)
{
    modbus_tcp_thread_var_t *tcp_thread_var = (modbus_tcp_thread_t *)param;
    int ret = -1, maxfd, i, need_send = 0;
    fd_set rset;
    struct timeval timeout;

    while (1) {
        SELECT_INIT();
        SELECT_ADD_FD(tcp_thread_var->send_timer);
        SELECT_ADD_FD(tcp_thread_var->period_timer);

        if (need_send) {
            timeout.tv_usec = 5000;
            timeout.tv_sec = 0;
        }
        else {
            timeout.tv_usec = 0;
            timeout.tv_sec = 5;
        }

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

            if (error == EINTR) {
                continue;
            }
            else {
                break;
            }
        }
        else if (ret == 0) {
            need_send = modbus_tcp_send_poll(tcp_thread_var, 0);
        }
        else if (ret > 0) {
            if (tcp_thread_var->send_timer > 0 && FD_ISSET(tcp_thread_var->send_timer, &rset)) {
                FD_CLR(tcp_thread_var->send_timer, &rset);
                need_send = modbus_tcp_send_poll(tcp_thread_var, 1);
            }
            if (tcp_thread_var->period_timer > 0 && FD_ISSET(tcp_thread_var->period_timer, &rset)) {
                FD_CLR(tcp_thread_var->period_timer, &rset);
                modbus_tcp_get_period_cmd(tcp_thread_var);
            }
        }
    }

    return NULL;
}

int modbus_tcp_start_thread(modbus_tcp_var_t *var, modbus_tcp_thread_var_t *tcp_thread_var)
{
    pthread_t thread_modbus_tcp;

    dbg_syslog(LOG_NOTICE, "start modbus_tcp_ip: %s, modbus_tcp_port: %d", tcp_thread_var->tcp_ip_addr, tcp_thread_var->tcp_port);

    modbus_tcp_mqtt_client_init(tcp_thread_var);
    modbus_tcp_period_rs_triger(tcp_thread_var);

    tcp_thread_var->send_timer = my_timer_create();
    if (tcp_thread_var->send_timer > 0) {
        my_timer_set(tcp_thread_var->send_timer, 1, 100);
    }

    tcp_thread_var->period_timer = my_timer_create();
    if (tcp_thread_var->period_timer > 0) {
        my_timer_set(tcp_thread_var->period_timer, 1, 5000);
    }

    tcp_thread_var->start_flag = 1;
    tcp_thread_var->period_flag = 1;

    pthread_create(&thread_modbus_tcp, NULL, modbus_tcp_thread_loop, (void *)tcp_thread_var);

    return 0;
}
