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

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

gw_port_e ports[] = {RS485_1, RS485_2, RS485_3, RS485_4,
    RS485_5, RS485_6, RS485_7, RS485_8,
    REMOTE_CH_01,
    REMOTE_CH_02,
    REMOTE_CH_03,
    REMOTE_CH_04,
    REMOTE_CH_05,
    REMOTE_CH_06,
    REMOTE_CH_07,
    REMOTE_CH_08,
    REMOTE_CH_09,
    REMOTE_CH_10,
    REMOTE_CH_11,
    REMOTE_CH_12,
    REMOTE_CH_13,
    REMOTE_CH_14,
    REMOTE_CH_15,
    REMOTE_CH_16,
    REMOTE_CH_17,
    REMOTE_CH_18,
    REMOTE_CH_19,
    REMOTE_CH_20,
    RS485_9, RS485_10, RS485_11, RS485_12,
};

static int ds_flush(rs485_var_t *var)
{
	int rval;
	status_list_t *tmp;
	status_list_t *status = NULL;

	rval = pthread_mutex_lock(&var->status_lock);
	list_for_each_entry_safe(status, tmp, &var->status_list, list)
    {
    	flush_device_status_by_sn(var->ds, status->action, status->sn);
        list_del(&status->list);
        free(status);
    }
    if (rval == 0) {
        pthread_mutex_unlock(&var->status_lock);
    } else {
        dbg_syslog(LOG_ERR, "Error, failed to acquire status lock: %d", rval);
    }
    return 0;
}

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

    TIMER_CONFIRM(var->node_status_timer);

	ds_flush(var);
    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);

        dbg_syslog(LOG_DEBUG, "online_status_changed:%d strlen(node_status):%d, (now-lasttime):%d",
                  var->ds->online_status_changed, (int) strlen(node_status), (int) (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);
            dbg_syslog(LOG_INFO, "report dev status topic  : %s", topic);
            dbg_syslog(LOG_INFO, "report dev status payload: %s", node_status);

            ret = ipc_session_publish(var->session, topic, (unsigned char *) node_status, strlen(node_status));
            if (ret == 0)
            {
                lasttime = now;
                var->ds->online_status_changed = false;
            }
            node_sta_backup(var->ds, RS485_NODE_STATUS_BAK_FILE);
        }
        if (node_status)
        {
            write_file_data(NODES_CACHE "/rs485_status.json", node_status, strlen(node_status));
            free(node_status);
        }
    }
    return 0;
}

static void rs485_loop(rs485_var_t *var)
{
    int ret = -1;
    struct pollfd pfd;

    while (1)
    {
        pfd.fd = var->node_status_timer;
        pfd.events = POLLIN;
        pfd.revents = 0;
        ret = poll(&pfd, 0x1, 1000);
        if (ret < 0)
        {
            int error = errno;	
            dbg_syslog(LOG_WARNING, "node_status_timer %d errno %d\n", var->node_status_timer, error);

            if (error == EINTR) {
                continue;
            }
            else {
                break;
            }
        }
		else if(ret == 0)
		{
			ds_flush(var);
		}
        else if (ret > 0)
        {
            if (var->node_status_timer > 0 && pfd.revents)
            {
                ds_flush_timer(var);
            }
        }
    }
}

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

    snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/%s/device/+/data/%s",var->sn_str, "RS485", 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",var->sn_str, "RS485", 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_var_t *var = (rs485_var_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);
    }
    return 0;
}

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

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

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

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

int rs485_init(rs485_var_t *var)
{
    int i, portnum;

    check_make_dir(NODES_CACHE);
    check_make_dir(NODES_CFG);
    pthread_mutex_init(&var->status_lock, NULL);
	INIT_LIST_HEAD(&var->status_list);
    get_board_sn(var->sn_str);
    get_process_name(var->proc_name);
    get_product_name(var->product_name);
    var->board_name = get_board_name();

    dbg_syslog(LOG_NOTICE, "board name :%s", var->board_name);
    dbg_syslog(LOG_NOTICE, "board SN   :%s", var->sn_str);

    if (load_nodes_cfg(&var->nodes_cfg_table, NODES_CFG_PATH) == -1)
    {
        dbg_syslog(LOG_ERR, "load nodes cfg fail");
    }
    if (load_templates_cfg(&var->template_table, TEMPLATES_CFG_PATH) == -1)
    {
        dbg_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, RS485_NODE_STATUS_BAK_FILE);

	rs485_load_config(var);
	rs485_load_special_config(var);
	rs485_mqtt_client_init(var);

    portnum = var->rs485_ports_cfg.port_num;
    for (i = 0; i < portnum; i++)
    {
        int port = rs485_port_map(var->rs485_ports_cfg.rs485_port_cfg[i].port);

		if (var->rs485_ports_cfg.rs485_port_cfg[i].disable == 0)
		{
			rs485_start_thread(var, &var->rs485_ports_cfg.rs485_port_cfg[i], port);
		}
    }

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

    return 0;
}

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

    memset(&var, 0, sizeof(var));
    lnxall_loglevel_set(LNXALL_LOGNOTICE, 1);

    /* ignore SIGPIPE for RS485/modbus TCP connection */
    if (signal(SIGPIPE, SIG_IGN) == SIG_ERR) {
        dbg_syslog(LOG_ERR, "failed to ignore signal SIGPIPE: %s", strerror(errno));
    }

    if (argc > 1 && argv[1] != NULL) {
        int sec;
        struct timespec spec;
        sec = (int) strtol(argv[1], NULL, 0);
        if (sec > 0 && sec <= 600) {
            spec.tv_sec = (time_t) sec;
            spec.tv_nsec = 0;
            dbg_syslog(LOG_ERR, "delay %d seconds on start...", sec);
            nanosleep(&spec, NULL);
        }
    }
    
    rs485_init(&var);
    rs485_loop(&var);

    return 0;
}
