#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 <poll.h>

#include "rs485_common.h"
#include "uart.h"

extern gw_port_e ports[];

static int debug_data(rs485_thread_t *rs485_thread, char *identifier, unsigned char *data, unsigned short len)
{
    tag_table_t ptag = {0};
    char *buf;
    int buff_len = len * 2 + 512;

    if (rs485_thread->rs485_port != 1 || strlen(rs485_thread->rs485_ports_cfg->bridging_port) == 0)
        return 0;

    buf = (char *)malloc(buff_len);
    if (buf == NULL)
    {
        dbg_syslog(LOG_ERR, "malloc fail");
        return -1;
    }
    b64_encode(data, len, buf, buff_len);
    cJSON *root = cJSON_CreateObject();
    cJSON_AddStringToObject(root, "raw_data", buf);

    strcpy(ptag.sn, "debug_data");
    ptag.mi = 0;
    strncpy(ptag.identifier, identifier, sizeof(ptag.identifier) - 1);
    ptag.report_period = 0;
    ptag.data_type = DATA_TYPE_SERVICE;
    ptag.tag_node = cJSON_PrintUnformatted(root);
    ptag.time = time(NULL);

    dbg_syslog(LOG_DEBUG, "debug_data:%s", ptag.tag_node);
    if (ptag.tag_node != NULL)
    {
        ptag.port = rs485_thread->rs485_port;
        pp_send_tag_data(rs485_thread->session, &ptag);
        free(ptag.tag_node);
    }

    free(buf);
    cJSON_Delete(root);
    return 0;
}

static int modbus_sniffer_recv_msg(modbus_sniffer_session_t *sniffer, cmd_data_node_t *cmd_msg)
{
    int ret;
    const char * tempid;
    rs485_var_data_t *rs485_var_data = (rs485_var_data_t *)cmd_msg->sniffer_data.userParams;
    rs485_var_t *var = rs485_var_data->var;               // container_of(sniffer, rs485_var_t, modbus_sniffer[2]);
    rs485_data_list_t *rs485_data = rs485_var_data->data; //(rs485_data_list_t *)cmd_msg->sniffer_data.userParams;
    pp_real_time_data_t real_data = {0};

    dbg_syslog_hex(LOG_DEBUG, cmd_msg->sniffer_data.data, cmd_msg->sniffer_data.len, "RECV sniffer sn %s src_identifier %s len %d", rs485_data->rs.sn, rs485_data->rs.src_identifier, cmd_msg->sniffer_data.len);

    real_data.port = rs485_port_remap(rs485_data->rs.port);
    real_data.len = cmd_msg->sniffer_data.len;
    real_data.instruction_code = cmd_msg->sniffer_data.data[1];
    real_data.data = malloc(real_data.len);

    strcpy(real_data.sn, rs485_data->rs.sn);
    strncpy(real_data.src_identifier, rs485_data->rs.src_identifier, sizeof(real_data.src_identifier));

    tempid = var->rs485_ports_cfg.rs485_port_cfg[real_data.port].sub_template_id;
    if (tempid[0]) {
        strncpy(real_data.sub_template_id, tempid, sizeof(real_data.sub_template_id));
    }
    memcpy(real_data.data, cmd_msg->sniffer_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);
    status_list_t *status = calloc(1, sizeof(status_list_t));
    strcpy(status->sn, real_data.sn);
    status->action = ACTION_RCV;
    ret = pthread_mutex_trylock(&rs485_var_data->var->status_lock);
    if (ret == 0) {
        list_add_tail(&status->list, &rs485_var_data->var->status_list);
        pthread_mutex_unlock(&rs485_var_data->var->status_lock);
    } else {
        if (ret != EBUSY) {
            dbg_syslog(LOG_ERR, "Error, failed to lock status: %d", ret);
        }
        free(status);
    }

    free(real_data.data);

    return 0;
}

int rs485_send_regulate_signal(rs485_thread_t *rs485_thread, pp_regulate_signal_t *regulate_data, int no_base64)
{
    char *data_tmp = NULL;
    char *buf = NULL;
    cJSON *rs_data = cJSON_CreateObject();
    char topic[256];

    buf = (char *)calloc(1, MAXBUF);
    if (rs_data == NULL || buf == NULL) {
        if (buf != NULL)
            free(buf);
        if (rs_data != NULL)
            cJSON_Delete(rs_data);
        dbg_syslog(LOG_ERR, "malloc fail");
        return -1;
    }

    if (no_base64 == 0)
        b64_encode(regulate_data->data, regulate_data->len, buf, MAXBUF);
    else
        memcpy(buf, regulate_data->data, regulate_data->len);
    cJSON_AddStringToObject(rs_data, "data_b64", buf);
    cJSON_AddNumberToObject(rs_data, "len", regulate_data->len);
    cJSON_AddNumberToObject(rs_data, "period", regulate_data->period);
    cJSON_AddStringToObject(rs_data, "port", port_enum2char(regulate_data->port));
    cJSON_AddNumberToObject(rs_data, "mi", regulate_data->mi);
    cJSON_AddStringToObject(rs_data, "src_identifier", regulate_data->src_identifier);
    cJSON_AddStringToObject(rs_data, "sn", regulate_data->sn);
    cJSON_AddStringToObject(rs_data, "dtu_sn", regulate_data->dtu_sn);
    cJSON_AddNumberToObject(rs_data, "protocol", regulate_data->protocol);
    cJSON_AddNumberToObject(rs_data, "communication_timeout", regulate_data->communication_timeout);
    cJSON_AddStringToObject(rs_data, "term_addr", (const char *) regulate_data->term_addr);
    cJSON_AddStringToObject(rs_data, "tcp_ip_addr", regulate_data->tcp_ip_addr);
    cJSON_AddNumberToObject(rs_data, "tcp_port", regulate_data->tcp_port);
    if (strlen(regulate_data->requester) > 0)
        cJSON_AddStringToObject(rs_data, "requester", regulate_data->requester);

    data_tmp = cJSON_Print(rs_data);
    memset(topic, 0, sizeof(topic));
    snprintf(topic, sizeof(topic), "ipc/%s/%s/device/%s/data/%s", rs485_thread->var->sn_str, port_enum2char(regulate_data->port),
             (const char *) regulate_data->sn, TOPIC_SEND_RGLT_SIGNAL_RAW_DATA);

    ipc_session_publish(rs485_thread->session, topic, (unsigned char *) data_tmp, strlen(data_tmp));

    // 回收资源
    free(data_tmp);
    free(buf);
    cJSON_Delete(rs_data);

    return 0;
}

static int rs485_msg_tag_data_ctrl(rs485_thread_t *rs485_thread, ipc_msg_t *mqtt_msg)
{
    char find = 0;
    pp_regulate_signal_t regulate_data;
    rs485_data_list_t *rs485_data;
    int ret;
    int port;

    //先检查参数
    memset(&regulate_data, 0, sizeof(regulate_data));
    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;
    }
    port = rs485_port_map(regulate_data.port);
    if (rs485_thread->rs485_ports_cfg->disable)
    {
        dbg_syslog(LOG_INFO, "RS485[%d] DISABLED!!!", port);
        free(regulate_data.data);
        return 0;
    }

    rs485_data = calloc(sizeof(rs485_data_list_t), 1);
    if (rs485_data)
    {
        int port = rs485_port_map(regulate_data.port);
        memcpy(&rs485_data->rs, &regulate_data, sizeof(regulate_data));

        //周期命令需要考虑去重
        if (regulate_data.period > 0)
        {
            rs485_thread->period_triggered = 1;
            rs485_data_list_t *period_cmd = NULL;
            list_for_each_entry(period_cmd, &rs485_thread->data_list[RS485_DATA_TYPE_PERIOD_CMD], list)
            {
                if (period_cmd->rs.len == rs485_data->rs.len && (memcmp(period_cmd->rs.data, rs485_data->rs.data, period_cmd->rs.len) == 0))
                {
                    find = 1;
                    break;
                }
            }
            if (!find)
            {
                list_add_tail(&rs485_data->list, &rs485_thread->data_list[RS485_DATA_TYPE_PERIOD_CMD]);

                if (rs485_thread->rs485_ports_cfg->sniffer_mode == 1)
                {
                    sniffer_data_t pdata;
                    rs485_var_data_t *rs485_var_data = calloc(1, sizeof(rs485_var_data_t));

                    rs485_var_data->var = rs485_thread->var;
                    rs485_var_data->data = rs485_data;

                    pdata.userParams = rs485_var_data;
                    pdata.len = rs485_data->rs.len;
                    pdata.data = (char *) rs485_data->rs.data;

                    dy_modbus_sniffer_session_add_cmd(&rs485_thread->modbus_sniffer, &pdata);
                }
            }
        }
        else
        {
            list_add_tail(&rs485_data->list, &rs485_thread->data_list[RS485_DATA_TYPE_CMD]);
        }

        dbg_syslog(LOG_INFO, "sn %s data[0] %d identifier %s port %d period %d", rs485_data->rs.sn, rs485_data->rs.data[0], regulate_data.src_identifier, port, regulate_data.period);
    }

    if (find) {
        dbg_syslog(LOG_WARNING, "Repeated commands !!!");
        free(rs485_data);
        free(regulate_data.data);
    }
    return 0;
}

int rs485_tag_data_report(rs485_thread_t *rs485_thread, char *data, int len, unsigned char port, regulate_cmd_backup_t *fields,
                          pp_real_time_data_t *real_data, char *key_str, char *requester, time_t now)
{
    int rval;
    const char * tempid;
    // 485口数据异常处理
    if (len >= 256)
    {
        char unusual[256];
        char unusual_0[256] = {0};

        memset(unusual, 0xFF, sizeof(unusual));
        if (memcmp(unusual, data, 256) == 0 || memcmp(unusual_0, data, 256) == 0)
        {
            free(real_data->down_data);
            dbg_syslog(LOG_WARNING, "RS485[%d] data ERROR !!!", port);
            return 0;
            // system(RS485_RESTART_CMD);
        }
    }

    real_data->port = rs485_port_remap(port);
    real_data->len = len;
    real_data->instruction_code = data[1];
    real_data->data = malloc(real_data->len);
    real_data->ts = now;

    tempid = rs485_thread->rs485_ports_cfg->sub_template_id;
    if (tempid[0]) {
        strncpy(real_data->sub_template_id, tempid, sizeof(real_data->sub_template_id));
    }
    memcpy(real_data->data, data, real_data->len);
    //桥接模式
    if (strlen(rs485_thread->rs485_ports_cfg->bridging_port) && (requester == NULL || strcmp(requester, rs485_thread->rs485_ports_cfg->bridging_port) == 0))
    {
        pp_regulate_signal_t *rglt_data = calloc(1, sizeof(pp_regulate_signal_t));

        rglt_data->data = real_data->data;
        rglt_data->len = real_data->len;
        rglt_data->port = port_char2enum(rs485_thread->rs485_ports_cfg->bridging_port);
        strncpy(rglt_data->requester, port_enum2char(ports[rs485_thread->rs485_port]), sizeof(rglt_data->requester) - 1);
        dbg_syslog_hex(LOG_DEBUG, rglt_data->data, rglt_data->len, "regulate signal parser len %d port %d period %d mi %d src_identifier %s", rglt_data->len, rglt_data->port, rglt_data->period, rglt_data->mi, rglt_data->src_identifier);
        rs485_send_regulate_signal(rs485_thread, rglt_data, 0);
        debug_data(rs485_thread, "up", rglt_data->data, rglt_data->len);
        free(rglt_data);
    }
    else
    {
        send_to_proto_parser(rs485_thread->session, real_data);
    }

    // flush_device_status_by_sn(rs485_thread->var->ds, ACTION_RCV, real_data->sn);status_list_t *status = calloc(1, size(status_list_t));
    status_list_t *status = calloc(1, sizeof(status_list_t));
    strcpy(status->sn, real_data->sn);
    status->action = ACTION_RCV;
    rval = pthread_mutex_trylock(&rs485_thread->var->status_lock);
    if (rval == 0) {
        list_add_tail(&status->list, &rs485_thread->var->status_list);
        pthread_mutex_unlock(&rs485_thread->var->status_lock);
    } else {
        if (rval != EBUSY) {
            dbg_syslog(LOG_ERR, "Error, failed to lock status: %d", rval);
        }
        free(status);
    }

    free(real_data->data);
    free(real_data->down_data);

    return 0;
}

void rs485_bus_read_msg(rs485_thread_t *rs485_thread, unsigned char port, int isremote)
{
    char buf[4096] = {0};
    char read_nBytes[4096] = {0};
    int len = 0;
    int ret = 0;
    int count = 0;
    char key_str[64] = {0};
    data_info_t *data_info = NULL;
    regulate_cmd_backup_t *fields = NULL;
    pp_real_time_data_t real_data = {0};
    char requester[8];

    if (rs485_thread->fd <= 0)
    {
        dbg_syslog(LOG_WARNING, "RS485[%d] DISABLED!!!", port);
        return;
    }

    /* 接收客户端的消息 */
    while (1)
    {
        ret = read(rs485_thread->fd, read_nBytes + len, 4096 - len);
        if (ret > 0)
        {
            len += ret;
            count = 0;
        }
        else
        {
            if (ret == 0 && isremote) {
                dbg_syslog(LOG_ERR, "RS485 TCP has been closed!");
                close(rs485_thread->fd);
                rs485_thread->fd = -1;
                return;
            }
            count++;
            if (count < rs485_thread->rs485_ports_cfg->frame_timeout / 10)
                usleep(10 * 1000);
            else
                break;
        }
    }

	/* periodically check tty settings for CH341 chip */
	if (port >= 0x4) {
		int rval;
		unsigned int nwt;
		struct timespec spec;
		spec.tv_sec = 0;
		spec.tv_nsec = 0;
		rval = clock_gettime(CLOCK_MONOTONIC, &spec);
		nwt = (unsigned int) spec.tv_sec;
		if (rval == 0 && (nwt >= (rs485_thread->last_reset + 10) ||
			nwt < rs485_thread->last_reset)) {
			rs485_thread->last_reset = nwt;
			rval = uart_reset_termios(rs485_thread->fd);
			dbg_syslog(LOG_DEBUG, "RESET RS485[%u]: %d", (unsigned int) port, rval);
		}
	}

    time_t now = time(NULL);

    pthread_mutex_lock(&rs485_thread->send_lock);

    rs485_thread->timeout_cnt = 0;
    const char frame_cmp[12] = {0};
    if(len >= 10 && memcmp(read_nBytes,frame_cmp,sizeof(frame_cmp)) == 0) 
        rs485_thread->all_zero = true;
    // restore the identifier
    {
        sprintf(key_str, "RS485_%d", port);
        data_info = rs485_thread->identifier_backup.get(&rs485_thread->identifier_backup, key_str);
        if (data_info)
        {
            fields = (regulate_cmd_backup_t *)data_info->data;

            if ((now - fields->ts) < 10)
            {
                real_data.mi = fields->mi;
                strcpy(real_data.sn, fields->sn);
                strncpy(real_data.src_identifier, fields->src_identifier, sizeof(real_data.src_identifier));
                strcpy(real_data.requester, fields->requester);
                real_data.down_data = malloc(fields->len);
                memcpy(real_data.down_data, fields->data, fields->len);
                real_data.down_len = fields->len;
                strcpy(requester, fields->requester);

                rs485_thread->identifier_backup.delete(&rs485_thread->identifier_backup, key_str);
            }
        }
        else
        {
            strncpy(real_data.src_identifier, "default", sizeof(real_data.src_identifier));
        }
    }
    pthread_cond_signal(&rs485_thread->send_cond);
    pthread_mutex_unlock(&rs485_thread->send_lock);

    if (len < 3)
    {
        dbg_syslog_hex(LOG_DEBUG, read_nBytes, len, "len %d GW<-RS485[%d] dirty data!!!", len, port);
        return;
    }
    //判断是否有头尾配置
    int i = 0;
    int start = -1;
    int end = -1;
    // dbg_syslog(LOG_DEBUG,"message %02X config %02X %d", read_nBytes[0], rs485_thread->rs485_ports_cfg->msg_head[0],rs485_thread->rs485_ports_cfg->msg_head_len);
    // dbg_syslog_hex(LOG_DEBUG, read_nBytes, len, "raw rcv len %d	ms:%lld	id:%s	GW<-RS485[%d]		", len, get_time_stamp_ms(), real_data.src_identifier, port);
    if (rs485_thread->rs485_ports_cfg->msg_head_len)
    {
        for (i = 0; i < len - (rs485_thread->rs485_ports_cfg->msg_head_len - 1); i++)
        {
            if (!memcmp(&read_nBytes[i], rs485_thread->rs485_ports_cfg->msg_head, rs485_thread->rs485_ports_cfg->msg_head_len))
            {
                start = i;
                break;
            }
        }
    }
    if (rs485_thread->rs485_ports_cfg->msg_tail_len)
    {
        for (i = start + (rs485_thread->rs485_ports_cfg->msg_tail_len); i < len - (rs485_thread->rs485_ports_cfg->msg_tail_len - 1); i++)
        {
            if (!memcmp(&read_nBytes[i], rs485_thread->rs485_ports_cfg->msg_tail, rs485_thread->rs485_ports_cfg->msg_tail_len))
            {
                end = i + rs485_thread->rs485_ports_cfg->msg_tail_len;
                break;
            }
        }
    }
    // dbg_syslog(LOG_DEBUG,"start pos %d ,end pos %d", start, end);

    if (start >= 0 && end > rs485_thread->rs485_ports_cfg->msg_tail_len) //有头有尾..而且找到头尾
    {
        len = end - start;
        memcpy(buf, &read_nBytes[start], len);
    }
    else if (rs485_thread->rs485_ports_cfg->msg_head_len == 0 && rs485_thread->rs485_ports_cfg->msg_tail_len == 0) //无头无尾,直接拷贝
    {
        if (rs485_thread->protocol == GW_PROTC_MODBUS && rs485_thread->modbus_addr != (unsigned char)read_nBytes[0] &&
            rs485_thread->modbus_addr == (unsigned char)read_nBytes[1] && ((unsigned char)read_nBytes[0] > 0xF0 || (unsigned char)read_nBytes[0] == 0x00))
        {
            dbg_syslog(LOG_DEBUG, "modbus correction data");
            len = len - 1;
            memcpy(buf, read_nBytes + 1, len);
        }
        else
            memcpy(buf, read_nBytes, len);
    }
    else //有头尾,但是没有找到
    {
        return;
    }

    dbg_syslog_hex(LOG_NOTICE, buf, len > 176 ? 176 : len, "	len %d	ms:%lld	id:%s	GW<-RS485[%d]		", len, get_time_stamp_ms(), real_data.src_identifier, port);
    dbg_syslog_hex(LOG_DEBUG, real_data.down_data, real_data.down_len, "require data RS485[%d]", port);

    if (rs485_thread->var->cpu_busy)
    {
        dbg_syslog(LOG_WARNING, "CPU busy data will be discarded...");
        return;
    }

    if (rs485_thread->rs485_ports_cfg->frame_header_len && memcmp(buf, rs485_thread->rs485_ports_cfg->frame_header, rs485_thread->rs485_ports_cfg->frame_header_len))
    {
        dbg_syslog(LOG_WARNING, "frame_header invalid header data 0x%02x...", buf[0]);
        return;
    }
    if (rs485_thread->rs485_ports_cfg->frame_tail_len && memcmp(buf + len - rs485_thread->rs485_ports_cfg->frame_tail_len, rs485_thread->rs485_ports_cfg->frame_tail, rs485_thread->rs485_ports_cfg->frame_tail_len))
    {
        dbg_syslog(LOG_WARNING, "frame_tail invalid tail data 0x%02x...", buf[len - 1]);
        return;
    }

    if (rs485_thread->rs485_ports_cfg->sniffer_mode == 1)
        dy_modbus_sniffer_session_add_data(&rs485_thread->modbus_sniffer, buf, len);
    else
        rs485_tag_data_report(rs485_thread, buf, len, port, fields, &real_data, key_str, requester, now);
}

void rs485_bus_send(rs485_thread_t *rs485_thread, pp_regulate_signal_t *rs, unsigned char port, unsigned int communication_timeout)
{
    int ret;
	const char *data = (const char *) rs->data;
    int len = rs->len;

    if (rs485_thread->period_flag == 0 && memcmp("__", rs->src_identifier, 2) != 0) //如果关闭周期,应该是在升级了,这时收到非__开头的内部消息就丢弃
    {
        dbg_syslog(LOG_WARNING, "sn:%s  Send ignore [%s], device upgrading, ", rs->sn, rs->src_identifier);
        return;
    }
    if (rs485_thread->rs485_ports_cfg->sniffer_mode == 1)
        return;

    if (rs485_thread->fd <= 0)
    {
        dbg_syslog(LOG_WARNING, "RS485[%d] DISABLED!!!", port);
        return;
    }

    pthread_mutex_lock(&rs485_thread->send_lock);

    // store the identifier
    {
        char key_str[64] = {0};
        regulate_cmd_backup_t fields = {0};
        time_t now = time(NULL);

        strcpy(fields.src_identifier, rs->src_identifier);
        strcpy(fields.sn, rs->sn);
        strcpy(fields.requester, rs->requester);
        fields.mi = rs->mi;
        fields.ts = now;
        fields.protocol = rs->protocol;
        if (len < 256)
        {
            memcpy(fields.data, data, len);
            fields.len = len;
        }

        sprintf(key_str, "RS485_%d", port);
        rs485_thread->identifier_backup.insert(&rs485_thread->identifier_backup, key_str, &fields, sizeof(fields));
    }

    ret = write(rs485_thread->fd, data, len);
    if (ret == len)
    {
        dbg_syslog_hex(LOG_NOTICE, data, len, "sn:%s  len %d  identifier:%s  ms:%lld  GW->RS485[%d]", rs->sn, len, rs->src_identifier, get_time_stamp_ms(), port);
        debug_data(rs485_thread, "down", (unsigned char *) data, len);
        rs485_thread->modbus_addr = data[0];

        #if _MTK7628
            communication_timeout = 1000;   // Work around for 7628 which receives messages more than 400ms after sending messages(modbus RTU)
        #endif

        lnxall_condvar_timedwait(&rs485_thread->send_cond, &rs485_thread->send_lock, communication_timeout);
    }
    else
    {
        dbg_syslog(LOG_WARNING, "sn:%s	GW->RS485[%d] len %d !!!", rs->sn, port, len);
    }

    pthread_mutex_unlock(&rs485_thread->send_lock);

    // flush_device_status_by_sn(rs485_thread->var->ds, ACTION_SND, rs->sn);
    status_list_t *status = calloc(1, sizeof(status_list_t));
    strcpy(status->sn, rs->sn);
    status->action = ACTION_SND;
    ret = pthread_mutex_trylock(&rs485_thread->var->status_lock);
    if (ret == 0) {
        list_add_tail(&status->list, &rs485_thread->var->status_list);
        pthread_mutex_unlock(&rs485_thread->var->status_lock);
    } else {
        if (ret != EBUSY) {
            dbg_syslog(LOG_ERR, "Error, failed to lock status: %d", ret);
        }
        free(status);
    }
}

static int send_poll_period(rs485_thread_t *rs485_thread, struct list_head *data_list, int port)
{
    int ret, find = 0;
    unsigned int communication_timeout = 200;
    rs485_data_list_t *tmp = NULL;
    long long now = get_time_stamp_ms();

    //从头开始遍历
    if (rs485_thread->send_data == NULL)
    {
        rs485_thread->send_data = list_first_entry(data_list, typeof(*rs485_thread->send_data), list);
    }

    list_for_each_entry_safe_continue(rs485_thread->send_data, tmp, data_list, list)
    {
        if (rs485_thread->send_data->rs.period > 0)
        {
            if (rs485_thread->period_flag == 0)
                continue;
            ret = check_timeout_msecond(now, &rs485_thread->send_data->last_send, rs485_thread->send_data->rs.period);
            if (ret)
            {
                find = 1;
                if (rs485_thread->send_data->rs.communication_timeout >= 40)
                {
                    communication_timeout = rs485_thread->send_data->rs.communication_timeout;
                }
                rs485_thread->timeout_cnt++;
                rs485_bus_send(rs485_thread, &rs485_thread->send_data->rs, port, communication_timeout);
                break;
            }
        }
    }

    return find;
}

static int send_poll(rs485_thread_t *rs485_thread, struct list_head *data_list, int port)
{
    int ret, find = 0;
    unsigned int communication_timeout = 200;
    rs485_data_list_t *send_data = NULL;
    long long now = get_time_stamp_ms();

    list_for_each_entry(send_data, data_list, list)
    {
        if (!send_data->rs.period)
        {
            find = 1;
            if (send_data->rs.communication_timeout >= 40)
            {
                communication_timeout = send_data->rs.communication_timeout;
            }
            rs485_bus_send(rs485_thread, &send_data->rs, port, communication_timeout);
            list_del(&send_data->list);
            free(send_data->rs.data);
            free(send_data);
            break;
        }
        else
        {
            //第一次发送延时
            if (send_data->last_send < 5)
            {
                send_data->last_send++;
                find = 1;
                break;
            }
            if (rs485_thread->period_flag == 0)
                continue;
            ret = check_timeout_msecond(now, &send_data->last_send, send_data->rs.period);
            if (ret)
            {
                find = 1;
                if (send_data->rs.communication_timeout >= 40)
                {
                    communication_timeout = send_data->rs.communication_timeout;
                }
                rs485_bus_send(rs485_thread, &send_data->rs, port, communication_timeout);
                break;
            }
        }
    }

    return find;
}

static int rs485_send_poll(rs485_thread_t *rs485_thread, int time_trigger)
{
    int rval;

    rval = 0;
    if (time_trigger)
        TIMER_CONFIRM(rs485_thread->send_timer);

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

    // for (j = 0; j < RS485_DATA_TYPE_CMD_CNT; j++)
    {
        if (send_poll(rs485_thread, &rs485_thread->data_list[RS485_DATA_TYPE_CMD], rs485_thread->rs485_port) == 1) {
#if 0
            /* do not return, otherwise period data is lost */
            return 1;
#endif
            rval = 1;
        }
        if (send_poll_period(rs485_thread, &rs485_thread->data_list[RS485_DATA_TYPE_PERIOD_CMD], rs485_thread->rs485_port) == 1)
            return 1;
    }
    return rval;
}

static void rs485_subscribe_all(rs485_thread_t *rs485_thread)
{
    ipc_session_t *session = rs485_thread->session;
    char topic[TOPIC_MAX_LEN] = {0};

    snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/%s/device/+/data/%s",rs485_thread->var->sn_str, port_enum2char(ports[rs485_thread->rs485_port]), TOPIC_SEND_RGLT_SIGNAL_RAW_DATA);
    dbg_syslog(LOG_NOTICE, "sub topic: %s", topic);
    ipc_session_subscribe(session, topic);

    snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/%s/%s",rs485_thread->var->sn_str, port_enum2char(ports[rs485_thread->rs485_port]), TOPIC_NOTIFY_UPGRADE);
    dbg_syslog(LOG_NOTICE, "sub topic: %s", topic);
    ipc_session_subscribe(session, topic);
}

static int rs485_mqtt_handle_recv_msg(void *obj, ipc_msg_t *mqtt_msg)
{
    rs485_thread_t *rs485_thread = (rs485_thread_t *)obj;
    dbg_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))
    {
        rs485_msg_tag_data_ctrl(rs485_thread, mqtt_msg);    //发送下行数据
    }
    else if (strstr(mqtt_msg->topic, TOPIC_NOTIFY_UPGRADE))
    {
        if (strcmp(mqtt_msg->payload, START_UPGRADE) == 0)
        {
            dbg_syslog(LOG_NOTICE, "receive start upgrade, stop period poll");
            rs485_thread->period_flag = 0;
        }
        else
        {
            dbg_syslog(LOG_NOTICE, "receive stop upgrade, start period poll");
            rs485_thread->period_flag = 1;
        }
    }
    return 0;
}

// 建立与内部broker之间的MQTT连接
int rs485_mqtt_client_init(rs485_thread_t *rs485_thread)
{
    char clientId[256] = {0};

    snprintf(clientId, sizeof(clientId), "INT_rs485_%d_%s", rs485_thread->rs485_port, rs485_thread->var->sn_str);
    dbg_syslog(LOG_NOTICE, "clientId: %s", clientId);

    rs485_thread->session = ipc_session_new(clientId, (void *)rs485_thread, IPC_DEFAULT);
    if (rs485_thread->session == NULL)
        return -1;

    ipc_session_set_callbacks(rs485_thread->session, rs485_mqtt_handle_recv_msg, NULL);
    rs485_subscribe_all(rs485_thread);
    ipc_session_start(rs485_thread->session);
    return 0;
}

int rs485_period_rs_triger(rs485_thread_t *rs485_thread)
{
    int i, j, k;
    rs485_var_t *var = rs485_thread->var;

    for (i = 0; i < var->template_table->template_cnt; i++)
    {
        //解析周期命令
        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)
            {
                //节点信息
                for (k = 0; k < var->nodes_cfg_table->node_cnt; k++)
                {
                    if (var->nodes_cfg_table->node[k].port == rs485_thread->rs485_ports_cfg->port &&
                        strcmp(var->nodes_cfg_table->node[k].template_id, var->template_table->template[i].template_id) == 0)
                    {
                        char *data = service_rs_data(&var->template_table->template[i].service_tab->service[j], var->nodes_cfg_table->node[k].sn);
                        char topic[TOPIC_MAX_LEN] = {0};
                        snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/%s/device/%s/data/%s", var->sn_str, port_enum2char(var->nodes_cfg_table->node[k].port), var->nodes_cfg_table->node[k].sn, TOPIC_EVT_SET_RGLT);
                        ipc_session_publish(var->session, topic, (unsigned char *) data, strlen(data));
                        free(data);
                        rs485_thread->protocol = var->template_table->template[i].protocol;
                        if (rs485_thread->protocol == GW_PROTC_MODBUS)
                        {
                            rs485_thread->rs485_ports_cfg->sleep_before_read = 0;
                        }
                    }
                }
            }
        }
    }
    return 0;
}

static int rs485_get_period_cmd(rs485_thread_t *rs485_thread)
{
    TIMER_CONFIRM(rs485_thread->period_timer);
    if (rs485_thread->period_triggered == 0)
    {
        rs485_period_rs_triger(rs485_thread);
    }
    else if (rs485_thread->period_triggered_cnt == 0)
    {
        rs485_period_rs_triger(rs485_thread);
        rs485_thread->period_triggered_cnt = 1;
    }
    return 0;
}

static int check_uart_frame(rs485_thread_t *rs485_thread)
{
    TIMER_CONFIRM(rs485_thread->self_check_timer);
    rs485_port_cfg_t *rs485_port_cfg = rs485_thread->rs485_ports_cfg;
    if (!rs485_port_is_remote(rs485_port_cfg) && (rs485_thread->timeout_cnt >= TIMEOUT_CNT_MAX  ||  rs485_thread->all_zero)) {
        if(rs485_thread->all_zero)rs485_thread->all_zero = false;
        const char *buf = get_tty_from_port(rs485_port_cfg->port);
        dbg_syslog(LOG_INFO, "uart[%s] reload config with speed:%d", buf, rs485_port_cfg->speed);
        uart_uninit(rs485_thread->fd);
        rs485_thread->fd = uart_init(buf, rs485_port_cfg->speed, rs485_port_cfg->stop, rs485_port_cfg->parity, rs485_port_cfg->bits, 1);
    }
    return 0;
}

static void *rs485_thread_loop(void *param)
{
    rs485_thread_t *rs485_thread = (rs485_thread_t *)param;
    /* rs485_var_t *var = rs485_thread->var; */

    int ret = -1, need_send = 0;
    struct pollfd pfds[3];

    while (1)
    {
        pfds[0].fd = rs485_thread->send_timer;
        pfds[0].events = POLLIN;
        pfds[0].revents = 0;

        pfds[1].fd = rs485_thread->period_timer;
        pfds[1].events = POLLIN;
        pfds[1].revents = 0;

        pfds[2].fd = rs485_thread->self_check_timer;
        pfds[2].events = POLLIN;
        pfds[2].revents = 0;
        ret = poll(pfds, 0x3, need_send ? 10 : 5000);
        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 = rs485_send_poll(rs485_thread, 0);
        }
        else if (ret > 0)
        {
            if (rs485_thread->send_timer > 0 && pfds[0].revents)
            {
                need_send = rs485_send_poll(rs485_thread, 1);
            }
            if (rs485_thread->period_timer > 0 && pfds[1].revents)
            {
                rs485_get_period_cmd(rs485_thread);
            }
            if (rs485_thread->self_check_timer > 0 && pfds[2].revents)
            {
                check_uart_frame(rs485_thread);
            }
        }
    }
    return 0;
}

static void *rs485_thread_rx_loop(void *param)
{
    rs485_thread_t *rs485_thread = (rs485_thread_t *)param;
    /* rs485_var_t *var = rs485_thread->var; */
    rs485_port_cfg_t * portcfg;
    int ret = -1, isremote;
    struct timespec tspec;

    isremote = 0;
    portcfg = rs485_thread->rs485_ports_cfg;
    if (portcfg != NULL)
        isremote = rs485_port_is_remote(portcfg);
    while (1)
    {
        struct pollfd pfd;
        if (rs485_thread->fd < 0) {
            tspec.tv_sec = 0;
            tspec.tv_nsec = 0;
            if (isremote != 0) {
                dbg_syslog(LOG_ERR, "restart remote rs485/modbus: %s:%d",
                    portcfg->remote_ipaddr, portcfg->remote_port);
                rs485_thread->fd = uart_remote_init(
                    portcfg->remote_ipaddr, portcfg->remote_port);
                if (rs485_thread->fd < 0)
                    tspec.tv_sec = 10;
            } else {
                tspec.tv_sec = 10;
                dbg_syslog(LOG_ERR, "Error, rs485_thread->fd: %d", rs485_thread->fd);
            }
            if (tspec.tv_sec != 0) {
                nanosleep(&tspec, NULL);
                continue;
            }
        }

        pfd.fd = rs485_thread->fd;
        pfd.events = POLLIN;
        pfd.revents = 0;
        ret = poll(&pfd, 0x1, 5000);
        if (ret < 0)
        {
            int error = errno;	
            dbg_syslog(LOG_INFO, "errno %d\n", error);

            if (error == EINTR)
            {
                continue;
            }
            else
            {
                break;
            }
        }
        else if (ret > 0)
        {
            if (rs485_thread->fd > 0 && pfd.revents)
            {
                rs485_bus_read_msg(rs485_thread, rs485_thread->rs485_port, isremote);
            }
        }
    }
    return 0;
}

int rs485_start_thread(rs485_var_t *var, rs485_port_cfg_t *rs485_port_cfg, int rs485_port)
{
    int j;
    pthread_t thread_rs485;
    pthread_t thread_rs485_rx;
    rs485_thread_t *rs485_thread = NULL;
    const char *buf = NULL;
    int port = rs485_port_map(rs485_port_cfg->port);

    dbg_syslog(LOG_NOTICE, "=========rs485_port: %d=======", rs485_port);

    rs485_thread = calloc(1, sizeof(rs485_thread_t));
    if (rs485_thread == NULL)
    {
        dbg_syslog(LOG_ERR, "malloc rs485_thread failed");
        return -1;
    }

    rs485_thread->var = var;
    rs485_thread->rs485_ports_cfg = rs485_port_cfg;
    rs485_thread->rs485_port = rs485_port;
    rs485_thread->send_data = NULL;

    pthread_mutex_init(&rs485_thread->send_lock, NULL);
    lnxall_condvar_init(&rs485_thread->send_cond);

    for (j = 0; j < RS485_DATA_TYPE_CMD_CNT; j++)
    {
        INIT_LIST_HEAD(&rs485_thread->data_list[j]);
    }
    kv_array_init(&rs485_thread->identifier_backup, 32);

    if (strcmp(var->board_name, BOARD_WOOLINK_MT7621) == 0 || strcmp(var->board_name, BOARD_WOOLINK_MT7628) == 0)
    {
        if (rs485_port_cfg->sleep_before_read == 0)
        {
            rs485_port_cfg->sleep_before_read = 200; // for board use FTDI force set a sleep time before read
        }
    }

    buf = get_tty_from_port(rs485_port_cfg->port);
    if (rs485_port_cfg->bits == 0) {
        rs485_port_cfg->bits = 8;
    }

    if (rs485_port_is_remote(rs485_port_cfg)) {
        dbg_syslog(LOG_NOTICE, "remote RS485/modbus server: %s:%d", rs485_port_cfg->remote_ipaddr, rs485_port_cfg->remote_port);
        rs485_thread->fd = uart_remote_init(rs485_port_cfg->remote_ipaddr, rs485_port_cfg->remote_port);
    }
    else {
        rs485_thread->fd = uart_init(buf, rs485_port_cfg->speed, rs485_port_cfg->stop, rs485_port_cfg->parity, rs485_port_cfg->bits, 1);
    }
    if (rs485_thread->fd < 0)
    {
        dbg_syslog(LOG_ERR, "open tty device %s error!", buf);
        return -1;
    }
    set_gpio_for_port(var->product_name, rs485_port_cfg->port);

    dbg_syslog(LOG_NOTICE, "buf %s rs485[%d]_fd %d speed %d stop %d parity %d bits %d sniffer_mode %d", buf, port, rs485_thread->fd,
              rs485_port_cfg->speed, rs485_port_cfg->stop, rs485_port_cfg->parity, rs485_port_cfg->bits, rs485_port_cfg->sniffer_mode);

    if (rs485_port_cfg->sniffer_mode == 1)
    {
        rs485_thread->modbus_sniffer.handle_recv_msg = modbus_sniffer_recv_msg;
        rs485_thread->modbus_sniffer.protocol = rs485_port_cfg->protocol;
        rs485_thread->modbus_sniffer.store_info.len = SNIFFER_DATA_BUF_SIZE;
        rs485_thread->modbus_sniffer.store_info.buf = calloc(1, SNIFFER_DATA_BUF_SIZE);
        dy_modbus_sniffer_session_init(&rs485_thread->modbus_sniffer);
    }

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

    rs485_thread->period_timer = my_timer_create();
    if (rs485_thread->period_timer > 0)
    {
        my_timer_set(rs485_thread->period_timer, 1, 1000);
    }

    rs485_thread->self_check_timer = my_timer_create();
    if (rs485_thread->self_check_timer > 0)
    {
        my_timer_set(rs485_thread->self_check_timer, 1, 10000);
    }

    rs485_mqtt_client_init(rs485_thread);
    rs485_period_rs_triger(rs485_thread);

    rs485_thread->start_flag = 1;
    rs485_thread->period_flag = 1;

    pthread_create(&thread_rs485, NULL, rs485_thread_loop, (void *)rs485_thread);
    pthread_create(&thread_rs485_rx, NULL, rs485_thread_rx_loop, (void *)rs485_thread);
    return 0;
}
