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

static int send_data_in_query_list(channel_t *channel)
{
    tag_list_node_t *tmp = NULL;
    tag_list_node_t *tag_list = NULL;
    dy_db_session_t *session = channel->db_session;
    int ret = 0;

    if (list_empty(&session->query_list))
    {
        return 0;
    }

    list_for_each_entry_safe(tag_list, tmp, &session->query_list, list)
    {
        tag_table_t *p = &tag_list->tag;

        ret = mqtt_session_publish_sync(channel->session, p->ext, p->tag_node, p->tag_len);

        if (ret == 0)
        {
            // 发布成功保留在链表里，用于更新数据库
            p->report = REPORTED;
        }
        else
        {
            list_del(&tag_list->list);
            free(tag_list->tag.tag_node);
            free(tag_list->tag.ext);
            free(tag_list);
        }
    }

    if (list_empty(&session->query_list))
    {
        // 全部上报失败，不需要更新数据库
        dbg_syslog(LOG_WARNING, "all report failed");
        return 0;
    }
    else
    {
        list_move_tail_list(&session->query_list, &session->update_list);
        return 1;
    }
}

int data_report_poll(cloud_mqtt_var_t *var)
{
    TIMER_CONFIRM(var->report_timer);
    int i = 0;

    for (i = 0; i < var->channel_cnt; i++)
    {
        channel_t *channel = &var->channel[i];
        if (channel->type == LINK_TYPE_KAFKA)
        {
            continue;
        }

        if (channel->resume == 0)
            continue;

        if (access(CLOUD_MQTT_FORCE_WRITEDB, F_OK) == 0)
            continue;

        if (mqtt_session_get_state(channel->session) != MQTT_CONNECTED)
        {
            continue;
        }

        dy_db_session_select_not_repoted(channel->db_session, 720);
        send_data_in_query_list(channel);
        if (!list_empty(&channel->db_session->update_list))
        {
            dy_db_session_update_reported(channel->db_session);
            dy_db_session_delete_reported(channel->db_session);
        }
    }

    return 0;
}

static void cloud_mqtt_loop(cloud_mqtt_var_t *var)
{
    int ret = -1;
    struct pollfd pfds[2];

    while (1)
    {
        pfds[0].fd = var->status_poll_timer;
        pfds[0].events = POLLIN;
        pfds[0].revents = 0;
        pfds[1].fd = var->report_timer;
        pfds[1].events = POLLIN;
        pfds[1].revents = 0;
        ret = poll(pfds, 0x2, 10 * 1000);
        if (ret < 0)
        {
            int error;
            error = errno;
            printf("errno %d\n", error);
            if (error == EINTR)
            {
                continue;
            }
            else
            {
                break;
            }
        }
        else if (ret > 0)
        {
            if (var->status_poll_timer > 0 && pfds[0].revents)
            {
                status_poll(var);
            }
            if (var->report_timer > 0 && pfds[1].revents)
            {
                data_report_poll(var);
            }
        }
    }
}

static int cloud_mqtt_init(cloud_mqtt_var_t *var)
{
    get_board_sn(var->sn_str);
    get_process_name(var->proc_name);
    get_product_name(var->product_name);

    dbg_syslog(LOG_INFO, "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 (cloud_mqtt_load_cfg(var, CLOUD_MQTT_CFG) < 0)
    {
        dbg_syslog(LOG_ERR, "load config failed");
        return -1;
    }

    cloud_mqtt_channel_create(var);
    cloud_kafka_channel_create(var);

    var->status_poll_timer = my_timer_create();
    if (var->status_poll_timer <= 0)
    {
        dbg_syslog(LOG_INFO, "[ERR] create status check timer fail, %d", var->status_poll_timer);
    }

    if (var->status_poll_timer > 0)
    {
        my_timer_set(var->status_poll_timer, 10, 1000);
    }

    var->report_timer = my_timer_create();
    if (var->report_timer <= 0)
    {
        dbg_syslog(LOG_INFO, "[ERR] create status check timer fail, %d", var->report_timer);
    }

    if (var->report_timer > 0)
    {
        my_timer_set(var->report_timer, 10, 60 * 1000);
    }
    cloud_mqtt_internal_mqtt_init(var);

    dbg_syslog(LOG_INFO, "init done");

    return 0;
}

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

    lnxall_loglevel_set(LNXALL_LOGNOTICE, 1);

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

    cloud_mqtt_init(&var);
    cloud_mqtt_loop(&var);

    return 0;
}
