#include "dr_common.h"
#include <lnxall_ubuslog.h>

static int dr_get_tag_data_from_str(dr_var_t *var, tag_table_t *tag, char *json_str)
{
    int ret = 0;

    cJSON *root = cJSON_Parse(json_str);
    if (root)
    {
        char buf[128];

        GET_JSON_VALUE_STRING(root, "sn", tag->sn);
        GET_JSON_VALUE_STRING(root, "requester", tag->requester);
        GET_JSON_VALUE_STRING(root, "identifier", tag->identifier);
        GET_JSON_VALUE_DY_STRING(root, "tag_node", tag->tag_node);
        tag->tag_len = tag->tag_node ? strlen(tag->tag_node) : 0;
        if (tag->tag_len == 0) {
            cJSON_Delete(root);
            return -1;
        }
        GET_JSON_VALUE_INT(root, "mi", tag->mi);
        GET_JSON_VALUE_INT(root, "time", tag->time);
        GET_JSON_VALUE_INT(root, "report", tag->report);
        GET_JSON_VALUE_INT(root, "res", tag->res);
        GET_JSON_VALUE_INT(root, "report_period", tag->report_period);
        GET_JSON_VALUE_INT(root, "data_type", tag->data_type);
        GET_JSON_VALUE_STRING(root, "port", buf);
        tag->port = port_char2enum(buf);
        GET_JSON_VALUE_STRING(root, "app_key", tag->app_key);
        GET_JSON_VALUE_DY_STRING(root, "ext", tag->ext);
    }
    else
    {
        dbg_syslog(LOG_ERR, "json parse error");
        ret = -1;
    }

    cJSON_Delete(root);
    return ret;
}

static int dr_condition_report(dr_var_t *var, tag_table_t *ptag, tag_table_t *ptag_change, tag_table_t *ptag_threshold, unsigned char *is_condition)
{
    int i, report_change = 0, report_threshold = 0;
    //根据设备与物模型关系，检索出SN、identifier
    //检索对应SN、identifier是否在条件上报列表中
    rtd_data_list_t *rtd_data = NULL;
    ptag_change->tag_node = NULL;
    ptag_threshold->tag_node = NULL;
    list_for_each_entry(rtd_data, &var->condition_list, list)
    {
        if (strcmp(rtd_data->tag.sn, ptag->sn) == 0) // && !strcmp(rtd_data->tag.identifier, ptag->identifier))
        {
            cJSON *jtag = cJSON_Parse(ptag->tag_node);
            cJSON *jtag_change = cJSON_CreateObject();
            cJSON *jtag_threshold = cJSON_CreateObject();
            cJSON *jtag_rtd = NULL;

            if (rtd_data->tag.tag_node)
            {
                jtag_rtd = cJSON_Parse(rtd_data->tag.tag_node);
            }
            else
                jtag_rtd = cJSON_CreateObject();

            //判断：变化上报\阈值上报,存储最新上报的值
            for (i = 0; i < rtd_data->poutput->propertyCnt; i++)
            {
                double val = 0xFFFFFF;
                char *identifier = rtd_data->poutput->property[i].identifier;
                if (!strstr(ptag->tag_node, identifier))
                {
                    continue;
                }

                GET_JSON_VALUE_DOUBLE(jtag, identifier, val);
                if (val == 0xFFFFFF)
                    continue;

                *is_condition = 1;
                if (!strcmp(rtd_data->poutput->property[i].condition.type, "change"))
                {
                    char buf[512] = {0};

                    //判断变化
                    double raw_val = 0xFFFFFF;
                    char raw_buf[512] = {0};
                    char identifier_raw[128] = {0};

                    sprintf(identifier_raw, "%s_raw", identifier);

                    if (rtd_data->poutput->property[i].raw_data_type != TAG_TYPE_VLENGTH)
                    {
                        GET_JSON_VALUE_DOUBLE(jtag_rtd, identifier, raw_val);
                    }
                    else
                    {
                        GET_JSON_VALUE_STRING(jtag, identifier, buf);
                        GET_JSON_VALUE_STRING(jtag_rtd, identifier, raw_buf);
                    }

                    if (rtd_data->poutput->property[i].raw_data_type != TAG_TYPE_VLENGTH && val != raw_val)
                    {
                        //检索对应物模型identifier,区分属性、服务、事件上报
                        cJSON_AddNumberToObject(jtag_change, identifier, val);

                        if (raw_val != 0xFFFFFF)
                        {
                            cJSON_AddNumberToObject(jtag_change, identifier_raw, raw_val);
                            cJSON_ReplaceItemInObject(jtag_rtd, identifier, cJSON_CreateNumber(val));
                        }
                        else
                            cJSON_AddNumberToObject(jtag_rtd, identifier, val);

                        report_change = 1;
                    }
                    else if (rtd_data->poutput->property[i].raw_data_type == TAG_TYPE_VLENGTH && strcmp(buf, raw_buf))
                    {
                        //检索对应物模型identifier,区分属性、服务、事件上报
                        cJSON_AddStringToObject(jtag_change, identifier, buf);

                        if (strlen(raw_buf))
                        {
                            cJSON_AddStringToObject(jtag_change, identifier_raw, raw_buf);
                            cJSON_ReplaceItemInObject(jtag_rtd, identifier, cJSON_CreateString(buf));
                        }
                        else
                            cJSON_AddStringToObject(jtag_rtd, identifier, buf);
                        report_change = 1;
                    }
                }
                else if (!strcmp(rtd_data->poutput->property[i].condition.type, "threshold"))
                {
                    //判断阈值
                    char identifier_threshold[128] = {0};
                    double thresholdVal = rtd_data->poutput->property[i].condition.thresholdVal;
                    sprintf(identifier_threshold, "%s_threshold", identifier);
                    if (rtd_data->poutput->property[i].condition.status != 2 && val >= thresholdVal)
                    {
                        //检索对应物模型identifier,区分属性、服务、事件上报
                        cJSON_AddNumberToObject(jtag_threshold, identifier, val);
                        cJSON_AddNumberToObject(jtag_threshold, identifier_threshold, thresholdVal);
                        rtd_data->poutput->property[i].condition.status = 2;
                        report_threshold = 1;
                    }
                    else if (rtd_data->poutput->property[i].condition.status != 1 && val < thresholdVal)
                    {
                        //检索对应物模型identifier,区分属性、服务、事件上报
                        cJSON_AddNumberToObject(jtag_threshold, identifier, val);
                        cJSON_AddNumberToObject(jtag_threshold, identifier_threshold, thresholdVal);
                        rtd_data->poutput->property[i].condition.status = 1;
                        report_threshold = 1;
                    }
                }
            }
            if (report_change || report_threshold)
            {
                //更新数据
                if (rtd_data->tag.tag_node)
                {
                    free(rtd_data->tag.tag_node);
                }

                if (jtag_rtd)
                {
                    rtd_data->tag.tag_node = cJSON_PrintUnformatted(jtag_rtd);
                }
            }

            if (report_change == 1)
            {
                ptag_change->tag_node = cJSON_PrintUnformatted(jtag_change);
                strcat(ptag_change->identifier, "_event_change");
                dbg_syslog(LOG_INFO, "condition report sn %s identifier %s tag_node %s ", rtd_data->tag.sn, ptag_change->identifier, ptag_change->tag_node);
            }
            if (report_threshold == 1)
            {
                ptag_threshold->tag_node = cJSON_PrintUnformatted(jtag_threshold);
                strcat(ptag_threshold->identifier, "_event_threshold");
                dbg_syslog(LOG_INFO, "condition report sn %s identifier %s tag_node %s ", rtd_data->tag.sn, ptag_threshold->identifier, ptag_threshold->tag_node);
            }
            cJSON_Delete(jtag_rtd);
            cJSON_Delete(jtag);
            cJSON_Delete(jtag_change);
            cJSON_Delete(jtag_threshold);

            break;
        }
    }

    return report_change + report_threshold;
}

static BOOL dr_dc_search(char * sn, char * identifier, cJSON * cjson, char * outdata)
{
	int i = 0;
	BOOL flag = FALSE;
	/* cJSON * node = NULL; */
	
	if (sn == NULL || identifier == NULL || g_dr_var_dc_t.db == NULL)
	{
		return flag;
	}
	
	for (i = 0; i < MAX_DC_NUM; i++)
	{
    	//dbg_syslog(LOG_INFO, "dr_cl_search, sn %s ,stored sn %s", sn, g_dr_var_cl_t.var_cl_s[i].sn_str);
    #if 0
		if (strcmp(sn, g_dr_var_cl_t.var_cl_s[i].sn_str) || strlen(sn) != strlen(g_dr_var_cl_t.var_cl_s[i].sn_str))
		{
			continue;
		}
	#endif
    	//dbg_syslog(LOG_INFO, "dr_cl_search, identifier %s ,stored identifier %s", identifier, g_dr_var_cl_t.var_cl_s[i].identifier);
		if (strcmp(identifier, g_dr_var_dc_t.var_dc_s[i].identifier) || strlen(identifier) != strlen(g_dr_var_dc_t.var_dc_s[i].identifier))
		{
			continue;
		}
    	//dbg_syslog(LOG_INFO, "dr_cl_search, found[%d]", i);
		
		flag = TRUE;
		break;
		/*忽略tag比对，目前先按identifier来做*/
	#if 0
		node = cJSON_GetObjectItem(cjson,g_dr_var_cl_t.var_cl_s[i].tag);
		if (node)
		{
			flag = TRUE;
			outdata = cJSON_Print(node);
			break;
		}
	#endif
	}

	return flag;
}

static int dr_data_report(dr_var_t *var, tag_table_t *ptag)
{
    cJSON *jdev = NULL;
    char topic[256];
    int ret = 0;
	char tagData[64] = {0};

    if (ptag->report == NONEED_REPORT)
    {
        dbg_syslog(LNXALL_LOGDEBUG, "NONEED_REPORT: %s, %s", ptag->sn, ptag->identifier);
        return 0;
    }
    jdev = cJSON_CreateObject();
    if (jdev == NULL)
    {
        dbg_syslog(LOG_ERR, "json create failed");
        return -1;
    }

    cJSON_AddStringToObject(jdev, "sn", ptag->sn);
    cJSON_AddNumberToObject(jdev, "mi", ptag->mi);
    cJSON_AddNumberToObject(jdev, "time", ptag->time);
    if (strlen(ptag->requester) > 0)
        cJSON_AddStringToObject(jdev, "requester", ptag->requester);

    if (strlen(ptag->identifier))
    {
        cJSON_AddStringToObject(jdev, "identifier", ptag->identifier);
    }
    cJSON *jtag = cJSON_Parse(ptag->tag_node);
    cJSON_AddItemToObject(jdev, "tags", jtag);
    char *data_tmp = cJSON_Print(jdev);
    //上报
	if (TRUE == dr_dc_search(ptag->sn, ptag->identifier, jtag, tagData))
	{
    	dbg_syslog(LOG_INFO, "store data to dc db, sn %s ,identifier %s, time:%ld", ptag->sn, ptag->identifier, ptag->time);
		dr_dc_store_tag_to_db(g_dr_var_dc_t.db, ptag->sn, ptag->identifier, cJSON_Print(jtag), ptag->time);
	}
	memset(topic, 0, sizeof(topic));
    dbg_syslog(LNXALL_LOGDEBUG, "report %s", data_tmp ? : "unknown");
    switch (ptag->data_type)
    {
    case DATA_TYPE_PROPERTY:
        snprintf(topic, sizeof(topic), "ipc/%s/%s/device/%s/data_filtered/property/%s", var->sn_str, strlen(ptag->app_key) > 0 ? ptag->app_key : port_enum2char(ptag->port), ptag->sn, ptag->identifier);
        ret = ipc_session_publish(var->session, topic, (unsigned char *) data_tmp, strlen(data_tmp));
        break;
    case DATA_TYPE_SERVICE:
        snprintf(topic, sizeof(topic), "ipc/%s/%s/device/%s/data_filtered/service/%s", var->sn_str, strlen(ptag->app_key) > 0 ? ptag->app_key : port_enum2char(ptag->port), ptag->sn, ptag->identifier);
        ret = ipc_session_publish(var->session, topic, (unsigned char *) data_tmp, strlen(data_tmp));
        break;
    case DATA_TYPE_EVENT:
        snprintf(topic, sizeof(topic), "ipc/%s/%s/device/%s/data_filtered/event/%s", var->sn_str, strlen(ptag->app_key) > 0 ? ptag->app_key : port_enum2char(ptag->port), ptag->sn, ptag->identifier);
        ret = ipc_session_publish(var->session, topic, (unsigned char *) data_tmp, strlen(data_tmp));
        break;
    case DATA_TYPE_PERIOD_PROPERTY:
        snprintf(topic, sizeof(topic), "ipc/%s/%s/device/%s/data_filtered/period_property/%s", var->sn_str, strlen(ptag->app_key) > 0 ? ptag->app_key : port_enum2char(ptag->port), ptag->sn, ptag->identifier);
        ret = ipc_session_publish(var->session, topic, (unsigned char *) data_tmp, strlen(data_tmp));
        break;
    default:
        ret = -1;
        dbg_syslog(LNXALL_LOGDEBUG, "unknown type %d: %s, %s",
            (int) ptag->data_type, ptag->sn, ptag->identifier);
        break;
    }
    if (ret == 0)
    {
        ptag->report = REPORTED;
    }

    free(data_tmp);
    cJSON_Delete(jdev);
    return ret;
}

static int process_period_report(dr_var_t *var, tag_table_t *ptag)
{
    if (ptag->report_period == 0)
    {
        ptag->report = NOT_REPORT;
        return 0;
    } else if (ptag->report_period == -1) {
        ptag->report = NONEED_REPORT;
        return 0;
    }

    int sn_index = hash_string(ptag->sn);
    int id_index = hash_string(ptag->identifier);
    dev_sn_node_t *dev_node;
    identifier_node_t *id_node = NULL;
    struct hlist_node *h1, *h2;
    time_t now = dr_uptime(NULL);

    hlist_for_each_entry_safe(dev_node, h1, &var->dev_sn_table[sn_index], list)
    {
        if (strcmp(dev_node->dev_sn, ptag->sn) == 0)
        {
            int find = 0;
            hlist_for_each_entry_safe(id_node, h2, &dev_node->identifier_table[id_index], list)
            {
                if (strcmp(id_node->identifier, ptag->identifier) == 0)
                {
                    // find matched identifier
                    find = 1;
                    if (now >= id_node->last_report + ptag->report_period)
                    {
                        ptag->report = NOT_REPORT;
                        id_node->last_report = now;
                    }
                    else
                    {
                        ptag->report = NONEED_REPORT;
                        dbg_syslog(LNXALL_LOGDEBUG, "Set NONEED_REPORT for %s, %s",
                            ptag->sn, ptag->identifier);
                    }
                    break;
                }
            }
            if (find == 0)
            {
                //对应的服务还没有在队列里面
                id_node = calloc(1, sizeof(identifier_node_t));
                strcpy(id_node->identifier, ptag->identifier);
                id_node->last_report = now;
                ptag->report = NOT_REPORT;
                hlist_add_head(&id_node->list, &dev_node->identifier_table[id_index]);
            }
            break;
        }
    }

    return 0;
}

static int dr_process_data(dr_var_t *var, ipc_msg_t *mqtt_msg)
{
    data_type_e type;
    tag_table_t tag = {0};
    unsigned char is_condition = 0;
    int insert = 1;

    dr_get_tag_data_from_str(var, &tag, mqtt_msg->payload);
    if (tag.tag_node == NULL)
    {
        dbg_syslog(LOG_WARNING, "tag node is NULL");
        return -1;
    }

    type = tag.data_type;
    tag.report = 0;

    if (var->condition_cnt)
    {
        // 有需要阈值上报或变化上报的物模型
        tag_table_t tag_change = tag;
        tag_table_t tag_threshold = tag;

        dr_condition_report(var, &tag, &tag_change, &tag_threshold, &is_condition);
        dbg_syslog(LOG_INFO, "identifier %s tag_change %p tag_threshold %p type %d is_condition:%d", tag.identifier, tag_change.tag_node, tag_threshold.tag_node, type, is_condition);
        if (tag_change.tag_node)
        {
            tag_change.data_type = DATA_TYPE_EVENT; // type;
            dr_data_report(var, &tag_change);
            free(tag_change.tag_node);
        }
        if (tag_threshold.tag_node)
        {
            tag_threshold.data_type = DATA_TYPE_EVENT;
            dr_data_report(var, &tag_threshold);
            free(tag_threshold.tag_node);
        }
    }

    if (type == DATA_TYPE_EVENT)
    {
        //事件立刻上报
        dr_data_report(var, &tag);
        free(tag.tag_node);
        free(tag.ext);
    }
    else if (type == DATA_TYPE_PROPERTY)
    {
        char *collect = NULL;
        int db_id = 0;
        struct list_head query_list;
        INIT_LIST_HEAD(&query_list);
        // 将采集上来的数据线和旧数据汇聚再插入数据库，DB线程会负责插入DB里
        dy_db_select_property_sn(var->db_session, tag.sn, &query_list);
        if (!list_empty(&query_list))
        {
            tag_list_node_t *tmp = NULL;
            tag_list_node_t *tag_list = NULL;

            list_for_each_entry_safe(tag_list, tmp, &query_list, list)
            {
                if (tag_list != 0)
                {
                    if (tag_list->tag.tag_node)
                    {
                        if (collect != NULL)
                            free(collect);
                        collect = strdup(tag_list->tag.tag_node);
                    }
                    db_id = tag_list->tag.id;
                    // 上报失败了，从list删除
                    list_del(&tag_list->list);
                    free(tag_list->tag.tag_node);
                    free(tag_list->tag.ext);
                    free(tag_list);
                }
            }
        }

        if (collect)
        {
            cJSON *root_collect = cJSON_Parse(collect);
            cJSON *root_single = cJSON_Parse(tag.tag_node);

            if (root_single && root_single->child)
            {
                cJSON *child = root_single->child;
                while (child)
                {
                    if (cJSON_GetObjectItem(root_collect, child->string))
                    {
                        cJSON_DeleteItemFromObject(root_collect, child->string);
                    }
                    if (cJSON_IsBool(child))
                    {
                        cJSON_AddBoolToObject(root_collect, child->string, child->valueint);
                    }
                    else if (cJSON_IsNumber(child))
                    {
                        cJSON_AddNumberToObject(root_collect, child->string, child->valuedouble);
                    }
                    else if (cJSON_IsString(child))
                    {
                        cJSON_AddStringToObject(root_collect, child->string, child->valuestring);
                    }
                    else
                    {
                        dbg_syslog(LOG_WARNING, "child json of then tag node type was error,collect failed");
                    }
                    child = child->next;
                }
            }
            free(collect);
            free(tag.tag_node);
            tag.tag_node = cJSON_PrintUnformatted(root_collect);
            tag.tag_len = strlen(tag.tag_node);
            tag.id = db_id;
            dr_data_report(var, &tag);
            dy_db_session_update_tag_content(var->db_session, &tag);
            cJSON_Delete(root_single);
            cJSON_Delete(root_collect);
        }
        else
        {
            dr_data_report(var, &tag);
            if (insert)
            {
                dy_db_session_insert_tag(var->db_session, &tag);
            }
            else
            {
                free(tag.tag_node);
                free(tag.ext);
            }
        }
    }
    else
    {
        process_period_report(var, &tag);
        if (is_condition == 0)
            dr_data_report(var, &tag);
        else {
            tag.report = REPORTED;
            dbg_syslog(LNXALL_LOGDEBUG, "is_condition = %d; %s, %s",
                (int) is_condition, tag.sn, tag.identifier);
        }

        if (0 /*insert*/)
        {
            // 将采集上来的数据插入到DB的list里面，DB线程会负责插入DB里
            dy_db_session_insert_tag(var->db_session, &tag);
        }
        else
        {
            free(tag.tag_node);
            free(tag.ext);
        }
    }
    return 0;
}

// 返回成功report的个数
int dr_report_query_result(dr_var_t *var)
{
    tag_list_node_t *tmp = NULL;
    tag_list_node_t *tag_list = NULL;
    dy_db_session_t *session = var->db_session;
    int ret = 0;
    int cnt = 0;

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

    list_for_each_entry_safe(tag_list, tmp, &session->query_list, list)
    {
        ret = dr_data_report(var, &tag_list->tag);
        if (ret != 0)
        {
            // 上报失败了，从list删除
            list_del(&tag_list->list);
            free(tag_list->tag.tag_node);
            free(tag_list->tag.ext);
            free(tag_list);
        }
        else
        {
            cnt++;
            //上报成功，保留在list里面，用于更新数据库的report字段
            tag_list->tag.report = REPORTED;
        }
    }

    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 cnt;
    }
}

static int dr_get_history_data(dr_var_t *var, ipc_msg_t *mqtt_msg)
{
#define COUNT_ONE_TIME 20
    cJSON *reply = NULL;
    cJSON *hist = NULL;
    struct list_head query_list;
    int mi = 0;
    char sn[SN_MAX_LEN] = {0};
    int ret = 0;
    int start_id = 0;
    char topic[256];
    cJSON *root = cJSON_Parse(mqtt_msg->payload);
    if (root == NULL)
    {
        ret = ERR_CODE_JSON_FORMAT;
        goto out;
    }

    GET_JSON_VALUE_INT(root, "mi", mi);
    GET_JSON_VALUE_STRING(root, "sn", sn);
    INIT_LIST_HEAD(&query_list);

	memset(topic, 0, sizeof(topic));
    while (1)
    {
        ret = dy_db_select_property_with_start_count(var->db_session, sn, start_id, COUNT_ONE_TIME, &query_list);
        if (ret == -1)
        {
            dbg_syslog(LOG_ERR, "db get error");
            goto out;
        }
        else if (ret == 0)
        {
            if (!list_empty(&query_list))
            {
                tag_list_node_t *tmp = NULL;
                tag_list_node_t *tag_list = NULL;
                int cnt = 0;

                if (reply == NULL)
                {
                    reply = cJSON_CreateObject();
                    if (reply == NULL)
                    {
                        dbg_syslog(LOG_ERR, "create json obj failed");
                        ret = -1;
                        goto out;
                    }
                    else
                    {
                        cJSON_AddNumberToObject(reply, "mi", mi);
                    }
                }

                hist = cJSON_CreateArray();
                cJSON_AddItemToObject(reply, "history_data", hist);

                list_for_each_entry_safe(tag_list, tmp, &query_list, list)
                {
                    if (tag_list != 0)
                    {
                        cJSON *item = cJSON_CreateObject();

                        cJSON_AddStringToObject(item, "sn", tag_list->tag.sn);
                        cJSON_AddStringToObject(item, "tags", tag_list->tag.tag_node);
                        cJSON_AddNumberToObject(item, "time", tag_list->tag.time);
                        start_id = tag_list->tag.id;

                        cJSON_AddItemToArray(hist, item);

                        list_del(&tag_list->list);
                        free(tag_list->tag.tag_node);
                        free(tag_list->tag.ext);
                        free(tag_list);
                        cnt++;
                    }
                }

                char *str = (cJSON_Print(reply));

                snprintf(topic, sizeof(topic), "G/%s/gateway/device/%s/service/%s", var->sn_str, var->sn_str, TOPIC_GET_HISTORY_DATA);
                ipc_session_publish(var->session, topic, (unsigned char *) str, strlen(str));
                free(str);

                cJSON_DeleteItemFromObject(reply, "history_data");

                if (cnt < COUNT_ONE_TIME)
                {
                    ret = 0;
                    dbg_syslog(LOG_DEBUG, "finished");
                    break;
                }
                else
                {
                    start_id++;
                    dbg_syslog(LOG_DEBUG, "not finished");
                }
            }
            else
            {
                dbg_syslog(LOG_DEBUG, "empty");
                break;
            }
        }
    }

out:

    cJSON_Delete(root);
    cJSON_Delete(reply);
    if (!list_empty(&query_list))
    {
        tag_list_node_t *tmp = NULL;
        tag_list_node_t *tag_list = NULL;

        list_for_each_entry_safe(tag_list, tmp, &query_list, list)
        {
            if (tag_list != 0)
            {
                // 上报失败了，从list删除
                list_del(&tag_list->list);
                free(tag_list->tag.tag_node);
                free(tag_list->tag.ext);
                free(tag_list);
            }
        }
    }

    return ret;
}

int dr_handle_mqtt_data(dr_var_t *var, ipc_msg_t *mqtt_msg)
{
    char topic[TOPIC_MAX_LEN] = {0};

    if (strstr(mqtt_msg->topic, TOPIC_GET_HISTORY_DATA))
    {
        sscanf(mqtt_msg->topic, "M/%*[^/]/gateway/device/%*[^/]/service/%" STR(TOPIC_MAX_LEN) "s", topic);
        if (strcmp(topic, TOPIC_GET_HISTORY_DATA) == 0)
        {
            return dr_get_history_data(var, mqtt_msg);
        }
    }

    if (strstr(mqtt_msg->topic, TOPIC_GET_CL_HISTORY_DATA))
    {
    	//ipc/[GW_SN]/gateway/device/[GW_SN]/service/Get_Cl_HistoryData
        sscanf(mqtt_msg->topic, "ipc/%*[^/]/gateway/device/%*[^/]/service/%" STR(TOPIC_MAX_LEN) "s", topic);
        if (strcmp(topic, TOPIC_GET_CL_HISTORY_DATA) == 0)
        {
            return dr_dc_get_history_data(var, mqtt_msg);
        }
    }

    return dr_process_data(var, mqtt_msg);
}

int dr_handle_period_property_report(dr_var_t *var)
{
    struct list_head query_list;
    time_t now = dr_uptime(NULL);

    int i;
    for (i = 0; i < var->nodes_cfg_table->node_cnt; i++)
    {
        // 将采集上来的数据线和旧数据汇聚再插入数据库，DB线程会负责插入DB里
        INIT_LIST_HEAD(&query_list);
        dy_db_select_property_by_sn(var->db_session, var->nodes_cfg_table->node[i].sn, &query_list);

        if (!list_empty(&query_list))
        {
            tag_list_node_t *tmp = NULL;
            tag_list_node_t *tag_list = NULL;

            list_for_each_entry_safe(tag_list, tmp, &query_list, list)
            {
                if (tag_list != 0)
                {
                    if (tag_list->tag.tag_node)
                    {
                        if ((var->period_property_info.hourly_report && now - tag_list->tag.time < 3600) ||
                            (var->period_property_info.period_report && now - tag_list->tag.time < var->period_property_info.period_report))
                        {
                            tag_list->tag.data_type = DATA_TYPE_PERIOD_PROPERTY;
                            dr_data_report(var, &tag_list->tag);
                        }
                        else
                            dbg_syslog(LOG_DEBUG, "invalid data collect time %d S", (int) (now - tag_list->tag.time));
                    }
                    // 上报失败了，从list删除
                    list_del(&tag_list->list);
                    free(tag_list->tag.tag_node);
                    free(tag_list->tag.ext);
                    free(tag_list);
                }
            }
        }
    }

    return 0;
}
