#include "dy_utils/dy_common.h"
#include "lnxall_list.h"
#include "mqtt_session.h"
#include "ipc_session.h"
#include "dy_utils/dy_ipc.h"

#define MQTT_SYSTOPIC_RULESENGINE_RULES "$SYS/rulesengine/rules"
#define MQTT_SYSTOPIC_RULESENGINE_ERROR "$SYS/rulesengine/error"
#define MQTT_SYSTOPIC_RULESENGINE_VERSION "$SYS/rulesengine/version"
#define MQTT_SYSTOPIC_RULESENGINE_VARIABLES "$SYS/rulesengine/variables"
#define MQTT_SYSTOPIC_RULESENGINE_STACKS "$SYS/rulesengine/stacks"

#define SERVICES_TOPIC_PART 5 ///<下行服务的topic有几部分要解析

#define INTERNAL_BROKER_ADDR "localhost"
#define INTERNAL_BROKER_USER ""
#define INTERNAL_BROKER_PASS ""
#define INTERNAL_RULE_ENG_PORT 1884
#define KEEP_ALIVE_MAX  3600
#define DEFAULT_IPC_QOS 0

typedef struct
{
    struct list_head list;
    char *topic;
    char *payload;
    int payloadlen;
} mqtt_msg_queue_t;

typedef struct
{
    struct list_head list;
    char *topic;
} subscrib_key_t;

typedef struct _rules_engin_var_t
{
    int module_slow_timer_fd;
    ipc_session_t *session;
    mqtt_session_t *rule_session;
    char sn_str[SN_MAX_LEN];
    char proc_name[PROGRAM_NAME_LEN];
    //此处根据模块需要添加
    nodes_cfg_table_t *nodes_cfg_table; //节点信息
} rules_engin_var_t;
static bool rules_debug = false;
static time_t rules_debug_stop = 0;

static struct list_head mqtt_mq;
static struct list_head subscrib_list;

static void careful_topic_clear()
{
    int i = 0;
    while (!list_empty(&subscrib_list))
    {
        subscrib_key_t *tmp = list_get_first(&subscrib_list, subscrib_key_t, list);
        if (tmp)
        {
            list_del(&tmp->list);
            free(tmp->topic);
            free(tmp);
        }
    }
}

static bool careful_topic_lookup(char *topic)
{
    bool result = false;
    subscrib_key_t *tmp = NULL;
    list_for_each_entry(tmp, &subscrib_list, list)
    {
        if (tmp && !strcmp(topic, tmp->topic))
        {
            result = true;
            break;
        }
    }
    return result;
}

static void careful_topic_insert(char *topic)
{
    if (careful_topic_lookup(topic))
    {
        return; ///如果存在，直接返回
    }
    dy_syslog(LOG_INFO, "subscribe topic :%s", topic);
    subscrib_key_t *tmp = calloc(sizeof(subscrib_key_t), 1);
    list_add_tail(&tmp->list, &subscrib_list);
    if (topic)
    {
        tmp->topic = malloc((strlen(topic) + 1));
        strcpy(tmp->topic, topic);
    }
}

static int rules_deamon_senser_api(rules_engin_var_t *var, mqtt_message_t *mqtt_msg, int type)
{
    char *data = mqtt_msg->payload;

    cJSON *api = cJSON_Parse(data);
    if (api)
    {
        cJSON *sn = cJSON_GetObjectItem(api, "sn");
        cJSON *identifier = cJSON_GetObjectItem(api, "identifier");
        cJSON *tags = cJSON_GetObjectItem(api, "tags");

        if (sn && sn->valuestring && tags && tags->child && identifier && identifier->valuestring)
        {
            cJSON *child = NULL;
            for (child = tags->child; child; child = child->next)
            {
                {
                    if (child->string)
                    {
                        int len_topic;
                        char topic[100] = {0};

                        if (type == DATA_TYPE_SERVICE)
                        {
                            len_topic = snprintf(topic, 100, "%s/%s/%s/%s/output/%s", var->sn_str, sn->valuestring, "services", identifier->valuestring, child->string);
                        }
                        else if (type == DATA_TYPE_EVENT)
                        {
                            len_topic = snprintf(topic, 100, "%s/%s/%s/%s/%s", var->sn_str, sn->valuestring, "events", identifier->valuestring, child->string);
                        }
                        else if (type == DATA_TYPE_PROPERTY)
                        {
                            len_topic = snprintf(topic, 100, "%s/%s/%s/%s/%s", var->sn_str, sn->valuestring, "properties", identifier->valuestring, child->string);
                        }
                        else
                        {
                            /* code */
                        }

                        if (strlen(topic) > 0)
                        {
                            char *p = cJSON_Print(child);
                            if (p != NULL)
                            {
                                if (careful_topic_lookup(topic))
                                {
                                    mqtt_session_publish(var->rule_session, topic, p, strlen(p));
                                }
                                free(p);
                            }
                        }
                    }
                }
            }
        }
    }
    cJSON_Delete(api);

    return 0;
}

static int rules_engins_action_api(rules_engin_var_t *var, char *topic, int topicLen, mqtt_message_t *message)
{
    int i;
    char *s[SERVICES_TOPIC_PART];
    if (topic && topicLen > 0)
    {
        if (strcmp(topic, MQTT_SYSTOPIC_RULESENGINE_STACKS) == 0)
        {
            write_file_data("/var/rules_stacks.json", message->payload, message->payloadLen);
            if (rules_debug)
            {
                char *str = strdup(message->payload);
                snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/%s/%s", var->sn_str, var->proc_name, TOPIC_EVT_RULESSTACKS);
                ipc_session_publish(var->session, topic, str, strlen(str));
                free(str);
            }
        }
        else
        {
            for (i = 0; i < SERVICES_TOPIC_PART; i++)
            {
                s[i] = malloc(topicLen) + 1;
            }

            int count = sscanf(topic, "%*[^/]/%[^/]/%[^/]/%[^/]/%s", s[0], s[1], s[2], s[3]);
            if (count >= 4) ///< 丢掉网关SN
            {
                cJSON *root = cJSON_CreateObject();
                if (!strcmp(s[1], "services") && !strcmp(s[3], "input") && message)
                {
                    cJSON *child = cJSON_Parse(message->payload);
                    cJSON_AddStringToObject(root, "sn", s[0]);
                    cJSON_AddStringToObject(root, "identifier", s[2]);

                    if (child && child->child)
                    {
                        cJSON *c = child->child;
                        while (c)
                        {
                            if (cJSON_IsBool(c))
                            {
                                cJSON_AddBoolToObject(root, c->string, c->valueint);
                            }
                            else if (cJSON_IsNumber(c))
                            {
                                cJSON_AddNumberToObject(root, c->string, c->valuedouble);
                            }
                            else if (cJSON_IsString(c))
                            {
                                cJSON_AddStringToObject(root, c->string, c->valuestring);
                            }
                            else
                            {
                                dy_syslog(LOG_WARNING, "child json of then tag node type was error,collect failed");
                            }
                            c = c->next;
                        }
                    }
                    cJSON_Delete(child);
                    char *data = cJSON_Print(root);

                    char topic[TOPIC_MAX_LEN] = {0};
                    char port[64] = "UNKNOWN";

                    for (i = 0; i < var->nodes_cfg_table->node_cnt; i++)
                    {
                        dy_syslog(LOG_DEBUG, "sn:%s var->nodes_cfg_table->node[%d].sn:%s", s[0], i, var->nodes_cfg_table->node[i].sn);
                        if (strcmp(s[0], var->nodes_cfg_table->node[i].sn) == 0)
                        {
                            if (strlen(var->nodes_cfg_table->node[i].app_key) != 0)
                            {
                                strcpy(port, var->nodes_cfg_table->node[i].app_key);
                            }
                            else
                            {
                                strcpy(port, port_enum2char(var->nodes_cfg_table->node[i].port));
                            }
                            break;
                        }
                    }
                    snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/%s/device/%s/data/%s", var->sn_str, port, s[0], TOPIC_EVT_SET_RGLT);
                    ipc_session_publish(var->session, topic, data, strlen(data));
                    dy_syslog(LOG_DEBUG, "rule topic:%s data:%s", topic, data);
                    free(data);

                    return 0;
                }
            }

            for (i = 0; i < SERVICES_TOPIC_PART; i++)
            {
                free(s[i]);
            }
        }
    }

    return -1;
}

static int rules_status_poll(rules_engin_var_t *var)
{
    TIMER_CONFIRM(var->module_slow_timer_fd);
    time_t now = time(NULL);
    if (now >= rules_debug_stop)
    {
        rules_debug = false;
    }
    return 0;
}

static void rules_loop(rules_engin_var_t *var)
{
    int ret = -1, maxfd = 0;
    fd_set rset;
    struct timeval timeout;
    while (1)
    {
        SELECT_INIT();
        SELECT_ADD_FD(var->module_slow_timer_fd);

        timeout.tv_usec = 0;
        timeout.tv_sec = 2;
        ret = select(maxfd + 1, &rset, 0, 0, &timeout);

        if (ret < 0)
        {
            dy_syslog(LOG_INFO, "%s_%d errno %d\n", __FUNCTION__, __LINE__, errno);

            if (errno == EINTR)
            {
                continue;
            }
            else
            {
                break;
            }
        }
        else if (ret > 0)
        {

            if (var->module_slow_timer_fd > 0 && FD_ISSET(var->module_slow_timer_fd, &rset))
            {
                FD_CLR(var->module_slow_timer_fd, &rset);

                rules_status_poll(var);
            }
        }
    }
}

static void rules_engins_config_paser(rules_engin_var_t *var, char *data)
{
    int ret = 1;
    int mi = 0;
    int m;
    careful_topic_clear();
    cJSON *root = cJSON_Parse(data);
    if (root)
    {
        cJSON *nodes = cJSON_GetObjectItem(root, "nodes");
        for (m = 0; m < cJSON_GetArraySize(nodes); m++)
        {
            cJSON *node = cJSON_GetArrayItem(nodes, m);
            cJSON *topic = cJSON_GetObjectItem(node, "topic");
            cJSON *value_topic = cJSON_GetObjectItem(node, "value-topic");
            cJSON *topics[2] = {topic, value_topic};
            int n;
            for (n = 0; n < 2; n++)
            {
                if (!topics[n])
                {
                    continue;
                }
                int i;
                char *s[SERVICES_TOPIC_PART];
                char *t = topics[n]->valuestring;
                dy_syslog(LOG_DEBUG, "rules nodes [%d] topic[%d] = %s", m, n, t);
                for (i = 0; i < SERVICES_TOPIC_PART; i++)
                {
                    s[i] = malloc(strlen(t) + 1);
                }
                int error = 0;
                int count = sscanf(t, "%*[^/]/%[^/]/%[^/]/%[^/]/%[^/]/%s", s[0], s[1], s[2], s[3], s[4]);

                if (error == 0 && count >= 3 && count <= SERVICES_TOPIC_PART)
                {
                    if (!(!strcmp(s[count - 3], "services") && !strcmp(s[count - 1], "input"))) //如果不是控制服务input就是控制服务的标记
                    {
                        careful_topic_insert(t); //为了标记规则引擎订阅了哪些数据
                    }
                }
                for (i = 0; i < SERVICES_TOPIC_PART; i++)
                {
                    free(s[i]);
                }
            }
        }
        mqtt_session_publish(var->rule_session, MQTT_SYSTOPIC_RULESENGINE_RULES, data, strlen(data));
        ret = 0;
    }
    cJSON_Delete(root);
}

static int rules_engins_debug(rules_engin_var_t *var, mqtt_message_t *mqtt_msg)
{
    cJSON *root = cJSON_Parse(mqtt_msg->payload);
    cJSON *keepalive = cJSON_GetObjectItem(root, "keep_alive");
    rules_debug_stop = keepalive->valueint + time(NULL);
    rules_debug = true;
    cJSON_Delete(root);
    gw_common_result(var->session, json_pkt_json_mi(mqtt_msg->payload), 0, NULL);
    return 0;
}

static int rules_mqtt_handle_recv_msg(void *obj, mqtt_message_t *mqtt_msg)
{
    rules_engin_var_t *var = (rules_engin_var_t *)obj;
    int ret = 0;
    bool matched;
    dy_syslog(LOG_DEBUG, "MSG MQTT client: received MQTT topic:%s payload length:%d",
              mqtt_msg->topic, mqtt_msg->payloadLen);

    // BUG: topicLen always equals 0
    ret = mosquitto_topic_matches_sub("ipc/+/+/device/+/data_filtered/property/+", mqtt_msg->topic, &matched);
    if (ret == 0 && matched)
    {
        //实时数据
        rules_deamon_senser_api(var, mqtt_msg, DATA_TYPE_PROPERTY);
    }
    else
    {
        //事件
        ret = mosquitto_topic_matches_sub("ipc/+/+/device/+/data_filtered/event/+", mqtt_msg->topic, &matched);
        if (ret == 0 && matched)
        {
            rules_deamon_senser_api(var, mqtt_msg, DATA_TYPE_EVENT);
        }
        else
        {
            ret = mosquitto_topic_matches_sub("ipc/+/+/device/+/data_filtered/service/+", mqtt_msg->topic, &matched);
            if (ret == 0 && matched)
            {
                rules_deamon_senser_api(var, mqtt_msg, DATA_TYPE_SERVICE);
            }
        }
    }

    if (strstr(mqtt_msg->topic, TOPIC_SET_RULESDEBUG))
    {
        rules_engins_debug(var, mqtt_msg);
    }
}

static int rules_mqtt_handle_rule_msg(void *obj, mqtt_message_t *mqtt_msg)
{
    rules_engin_var_t *var = (rules_engin_var_t *)obj;
    dy_syslog(LOG_DEBUG, "rule MQTT client: received MQTT topic:%s payload length:%d",
              mqtt_msg->topic, mqtt_msg->payloadLen);

    rules_engins_action_api(var, mqtt_msg->topic, strlen(mqtt_msg->topic), mqtt_msg);
}

static void rules_engins_load_cfg(rules_engin_var_t *var)
{
    char *json_str = NULL;
    json_str = read_file_data(RULES_CFG_LOCAL);
    if (!json_str)
    {
        return;
    }
    rules_engins_config_paser(var, json_str);
    free(json_str);
}

static void mqtt_subscribe_rules(rules_engin_var_t *var)
{
    char topic[TOPIC_MAX_LEN];
    ipc_session_t *session = var->session;
    sprintf(topic, "%s/+/+/+/input/#", var->sn_str); //订阅下行服务
    // TODO: which topic will be sub in the light of rules config

    mqtt_session_subscribe(var->rule_session, topic);
    mqtt_session_subscribe(var->rule_session, MQTT_SYSTOPIC_RULESENGINE_STACKS);

    snprintf(topic, TOPIC_MAX_LEN, "ipc/+/+/device/+/data_filtered/property/+");
    ipc_session_subscribe(session, topic);

    snprintf(topic, TOPIC_MAX_LEN, "ipc/+/+/device/+/data_filtered/event/+");
    ipc_session_subscribe(session, topic);

    snprintf(topic, TOPIC_MAX_LEN, "ipc/+/+/device/+/data_filtered/service/+");
    ipc_session_subscribe(session, topic);

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

static int mqtt_connect(rules_engin_var_t *var)
{
    char clientId[MAX_CLIENT_ID_LEN] = {0};

    snprintf(clientId, MAX_CLIENT_ID_LEN, "%s_%s", "rule_", var->sn_str);
    var->rule_session = mqtt_session_new(clientId, (void *)var);
    if (var->rule_session == NULL)
        return -1;

    mqtt_session_set_address(var->rule_session, INTERNAL_BROKER_ADDR, INTERNAL_RULE_ENG_PORT, INTERNAL_BROKER_USER, INTERNAL_BROKER_PASS);
    mqtt_session_set_opts(var->rule_session, DEFAULT_IPC_QOS, KEEP_ALIVE_MAX);
    mqtt_session_set_callbacks(var->rule_session, rules_mqtt_handle_rule_msg, NULL);
    mqtt_session_start(var->rule_session);

    snprintf(clientId, MAX_CLIENT_ID_LEN, "%s_%s", "internal_", 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, rules_mqtt_handle_recv_msg, NULL);
    ipc_session_start(var->session);

    return 0;
}

static int rule_eng_init(rules_engin_var_t *var)
{
    get_board_sn(var->sn_str);
    get_process_name(var->proc_name);
    int i = 0;

    if (load_nodes_cfg(&var->nodes_cfg_table, NODES_CFG_PATH) == -1)
    {
        dy_syslog(LOG_ERR, "load nodes cfg fail");
    }
    for (i = 0; i < var->nodes_cfg_table->node_cnt; i++)
    {
        dy_syslog(LOG_DEBUG, "var->nodes_cfg_table->node[%d].sn:%s", i, var->nodes_cfg_table->node[i].sn);
    }

    INIT_LIST_HEAD(&mqtt_mq);
    INIT_LIST_HEAD(&subscrib_list);
    mqtt_connect(var);
    mqtt_subscribe_rules(var);
    var->module_slow_timer_fd = my_timer_create();
    if (var->module_slow_timer_fd > 0)
    {
        my_timer_set(var->module_slow_timer_fd, 2, 200); //
    }
    rules_engins_load_cfg(var);
}

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

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

    rule_eng_init(&var);
    rules_loop(&var);

    return 0;
}
