#include <errno.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <sys/stat.h>
#include <sys/syscall.h>
#include <sys/time.h>
#include <sys/types.h>
#include <unistd.h>
#include <poll.h>

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

dr_var_dc_t g_dr_var_dc_t = {0};

time_t dr_uptime(time_t * tp)
{
	struct timespec spec;

	spec.tv_sec = 0;
	spec.tv_nsec = 0;
	if (clock_gettime(CLOCK_BOOTTIME, &spec) == -1)
		clock_gettime(CLOCK_MONOTONIC, &spec);

	if (tp != NULL)
		*tp = spec.tv_sec;
	return spec.tv_sec;
}

static int dr_tag_select_and_report(dr_var_t *var)
{
    TIMER_CONFIRM(var->report_timer);

    dy_db_session_select_not_repoted(var->db_session, 100);
    int cnt = dr_report_query_result(var);
    if (cnt > 0)
    {
        dy_db_session_update_reported(var->db_session);
    }

    return 0;
}

static int dr_period_property_report(dr_var_t *var)
{
    TIMER_CONFIRM(var->period_property_timer);

    time_t now = dr_uptime(NULL);
    if ((!var->period_property_info.hourly_report && var->period_property_info.period_report && now - var->period_property_info.last_check >= var->period_property_info.period_report) ||
        (var->period_property_info.hourly_report && now % 3600 == 0))
    {
        dbg_syslog(LOG_INFO, "period_property_timer:%d", var->period_property_timer);
        var->period_property_info.last_check = now;
        dr_handle_period_property_report(var);
    }

    return 0;
}

static void *dr_period_property_loop(void *param)
{
    int ret = -1;
    struct pollfd pfd;
    dr_var_t *var = (dr_var_t *)param;

    while (1)
    {
        pfd.fd = var->period_property_timer;
        pfd.events = POLLIN;
        pfd.revents = 0;
        ret = poll(&pfd, 0x1, 10 * 1000);
        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 (var->period_property_timer > 0 && pfd.revents)
            {
                dr_period_property_report(var);
            }
        }
    }
    return NULL;
}

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

    while (1)
    {
        pfd.fd = var->report_timer;
        pfd.events = POLLIN;
        pfd.revents = 0;
        ret = poll(&pfd, 0x1, 5000);
        if (ret < 0)
        {
            int error = errno;
            printf("errno %d\n", error);
            if (error == EINTR)
            {
                continue;
            }
            else
            {
                break;
            }
        }

        else if (ret > 0)
        {
            if (var->report_timer > 0 && pfd.revents)
            {
                dr_tag_select_and_report(var);
            }
        }
    }
}

static void dr_subscribe_all(dr_var_t *var)
{
    ipc_session_t *ipc_session = var->session;
    char topic[256];

    memset(topic, 0, sizeof(topic));
    snprintf(topic, sizeof(topic) - 1, "ipc/%s/+/device/+/data/property/+",var->sn_str);
    ipc_session_subscribe(ipc_session, topic);
    snprintf(topic, sizeof(topic) - 1, "ipc/%s/+/device/+/data/event/+",var->sn_str);
    ipc_session_subscribe(ipc_session, topic);
    snprintf(topic, sizeof(topic) - 1, "ipc/%s/+/device/+/data/service/+",var->sn_str);
    ipc_session_subscribe(ipc_session, topic);

    snprintf(topic, sizeof(topic) - 1, "M/+/%s/device/%s/service/%s", TOPIC_GATEWAY_CONFIG, var->sn_str, TOPIC_GET_HISTORY_DATA);
    ipc_session_subscribe(ipc_session, topic);

    snprintf(topic, sizeof(topic) - 1, "ipc/%s/%s/device/%s/service/%s", var->sn_str, TOPIC_GATEWAY_CONFIG, var->sn_str, TOPIC_GET_CL_HISTORY_DATA);
    ipc_session_subscribe(ipc_session, topic);	
}

static int dr_mqtt_handle_recv_msg(void *obj, ipc_msg_t *mqtt_msg)
{
    dr_var_t *var = (dr_var_t *)obj;

     dbg_syslog(LOG_INFO, "DR MQTT client: received MQTT topic:%s payload:%s length:%d",
               mqtt_msg->topic, mqtt_msg->payload, mqtt_msg->payloadLen);
    dr_handle_mqtt_data(var, mqtt_msg);
    return 0;
}

// 建立与内部broker之间的MQTT连接
static int dr_mqtt_client_init(dr_var_t *var)
{
    char clientId[256];

    memset(clientId, 0, sizeof(clientId));
    snprintf(clientId, sizeof(clientId) - 1, "INT_dr_%s", 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, dr_mqtt_handle_recv_msg, NULL);
    dr_subscribe_all(var);
    ipc_session_start(var->session);
    return 0;
}

void pp_find_sn_by_object(dr_var_t *var, char *object_model_id, char *identifier, object_property_table_t *poutput)
{
    int i;

    for (i = 0; i < var->nodes_cfg_table->node_cnt; i++)
    {
        node_cfg_t *node = &var->nodes_cfg_table->node[i];
        if (strcmp(node->product_key, object_model_id) == 0)
        {
            rtd_data_list_t *rtd_data = calloc(sizeof(rtd_data_list_t), 1);
            strcpy(rtd_data->tag.sn, node->sn);
            memcpy(rtd_data->tag.identifier, identifier, sizeof(rtd_data->tag.identifier));
            rtd_data->poutput = poutput;
            var->condition_cnt++;
            list_add_tail(&rtd_data->list, &var->condition_list);
        }
    }
}

static int dr_condition_init(dr_var_t *var)
{
    int i, j, k;

    if (var->object_table == NULL)
    {
        return 0;
    }

    //节点信息
    for (i = 0; i < var->object_table->object_cnt; i++)
    {
        object_cfg_t *object = &var->object_table->object[i];
        for (j = 0; j < object->service_tab->serviceCnt; j++)
        {
            object_service_cfg_t *obj_service = &object->service_tab->service[j];
            for (k = 0; k < obj_service->poutput->propertyCnt; k++)
            {
                if (strlen(obj_service->poutput->property[k].condition.type))
                {
                    //检索物模型对应SN
                    pp_find_sn_by_object(var, object->object_model_id, obj_service->identifier, obj_service->poutput);
                    break;
                }
            }
        }
        for (j = 0; j < object->event_tab->eventCnt; j++)
        {
            object_event_cfg_t *obj_event = &object->event_tab->event[j];
            for (k = 0; k < obj_event->poutput->propertyCnt; k++)
            {
                if (strlen(obj_event->poutput->property[k].condition.type))
                {
                    //检索物模型对应SN
                    pp_find_sn_by_object(var, object->object_model_id, obj_event->identifier, obj_event->poutput);
                    break;
                }
            }
        }
    }

    return 0;
}

static int dr_db_init(dr_var_t *var)
{
    int ret = 0;
    const char *board_name = NULL;
    board_name = get_board_name();

    if (strcmp(board_name, BOARD_WOOLINK_MT7628) == 0)
    {
        ret = dy_db_session_init(&var->db_session, DATA_COLLECT_DB_TMPFS_NAME, (void *)var);
    }
    else
    {
        ret = dy_db_session_init(&var->db_session, DATA_COLLECT_DB_NAME, (void *)var);
    }

    return ret;
}


static int dr_dc_db_create_table(sqlite3 *db, char *table_name, char *table_item)
{
    char sql[512] = {0};
    int ret, nRow, nCol;
    char **pResult;
    char *errMsg;

    if (db == NULL)
    {
        return -1;
    }

    snprintf(sql, sizeof(sql) - 1, "select count(*)  from sqlite_master where type='table' and name = '%s';", table_name);
    ret = sqlite3_get_table(db, sql, &pResult, &nRow, &nCol, &errMsg);
    if (ret != 0 || atoi(pResult[1]) == 0)
    {
        memset(sql, 0, sizeof(sql));
        snprintf(sql, sizeof(sql) - 1, "CREATE TABLE %s (%s);", table_name, table_item);
        ret = sqlite3_exec(db, sql, 0, 0, &errMsg);
        if (errMsg)
        {
            dbg_syslog(LOG_INFO, "[ERR] CREATE TABLE %s errMsg %s ", table_name, errMsg);
        }
    }
    sqlite3_free_table(pResult);

    return ret;
}

int dr_dc_store_tag_to_db(sqlite3 *db, char * sn, char * identifier, char * data, time_t timestamp)
{
	int ret = 0;
	static sqlite3_stmt *stmt3 = NULL;
	static sqlite3_stmt *delete_stmt3 = NULL;
	static time_t last_delete = 0;
	static sqlite3_stmt * packdb_stmt3 = NULL;
	
    time_t now = dr_uptime(NULL);

	if (db == NULL || sn == NULL || identifier == NULL)
    {
        dbg_syslog(LOG_WARNING, "[EER] INSERT param is empty, sn:%s, id:%s", sn, identifier);
        return -1;
    }

    if (stmt3 == NULL)
    {
        ret = sqlite3_prepare_v2(db, DR_DC_INSERT_SQL, strlen(DR_DC_INSERT_SQL), &stmt3, 0);
        if (ret != SQLITE_OK)
        {
            dbg_syslog(LOG_ERR, "prepare error %s\n", sqlite3_errmsg(db));
            sqlite3_finalize(stmt3);
            stmt3 = NULL;
            return ret;
        }
    }
	
    //在绑定时，最左面的变量索引值是1。
    sqlite3_bind_text(stmt3, 1, sn, strlen(sn), SQLITE_TRANSIENT);
	//dbg_syslog(LOG_WARNING, "INSERT dr cl tag node:%s length:%d", data, strlen(data));
    sqlite3_bind_blob(stmt3, 2, data, strlen(data), SQLITE_TRANSIENT);
    sqlite3_bind_int(stmt3, 3, timestamp);
    if (strlen(identifier) > 0)
    {
        sqlite3_bind_text(stmt3, 4, identifier, strlen(identifier), SQLITE_TRANSIENT);
    }
    else
    {
        sqlite3_bind_null(stmt3, 4);
    }
	
    ret = sqlite3_step(stmt3);
    if (ret != SQLITE_DONE)
    {
        dbg_syslog(LOG_ERR, "Insert tag data error");
    }
    //重新初始化该sqlite3_stmt对象绑定的变量。
    sqlite3_reset(stmt3);


    if (delete_stmt3 == NULL)
    {
        ret = sqlite3_prepare_v2(db, DR_DC_DELETE_DATA_SQL, strlen(DR_DC_DELETE_DATA_SQL), &delete_stmt3, 0);
        if (ret != SQLITE_OK)
        {
            dbg_syslog(LOG_ERR, "prepare error %s\n", sqlite3_errmsg(db));
            sqlite3_finalize(delete_stmt3);
            delete_stmt3 = NULL;
        }
    }
	
    if (packdb_stmt3 == NULL) {
        ret = sqlite3_prepare_v2(db, DR_DC_PACK_DATABASE, strlen(DR_DC_PACK_DATABASE), &packdb_stmt3, 0);
        if (ret != SQLITE_OK) {
            dbg_syslog(LOG_ERR, "prepare error %s\n", sqlite3_errmsg(db));
            sqlite3_finalize(packdb_stmt3);
            packdb_stmt3 = NULL;
        }
    }

    if (now - last_delete > 3600 * 6)
    {
        time_t ts_threshold = now - 60 * 60 * 24 * g_dr_var_dc_t.limit_time;
        sqlite3_bind_int(delete_stmt3, 1, ts_threshold);

        ret = sqlite3_step(delete_stmt3);
        if (ret != SQLITE_DONE)
        {
            dbg_syslog(LOG_ERR, "delete data error: %d", ret);
        }

        //重新初始化该sqlite3_stmt对象绑定的变量。
        sqlite3_reset(delete_stmt3);
        last_delete = now;

        /* pack database to save storage space */
        ret = (packdb_stmt3 != NULL) ? sqlite3_step(packdb_stmt3) : -1;
        if (ret != SQLITE_DONE) {
            dbg_syslog(LOG_ERR, "Pack database error: %d", ret);
        }
        sqlite3_reset(packdb_stmt3);
    }

    return ret;

	
}

static int dr_dc_db_get_select_all_from_stmt(sqlite3_stmt *stmt, struct list_head *query_list)
{
    dr_dc_tag_table_t *tag;
    dr_dc_tag_list_node_t *tag_list = NULL;
    const unsigned char *str;
#define COLUMN_TEXT_ERRMSG "Fatal Error, sqlite3_column_text(...) at line: %d!\n"

    while (sqlite3_step(stmt) == SQLITE_ROW)
    {
        tag_list = calloc(1, sizeof(tag_list_node_t));
        if (tag_list == NULL)
        {
            dbg_syslog(LOG_ERR, "malloc failed");
            goto out;
        }
        tag = &tag_list->tag;
        if (SQLITE_INTEGER == sqlite3_column_type(stmt, 0))
        {
            tag->id = sqlite3_column_int(stmt, 0);
        }
        if (SQLITE_TEXT == sqlite3_column_type(stmt, 1))
        {
            str = sqlite3_column_text(stmt, 1);
            if (str == NULL) {
                dbg_syslog(LOG_ERR, COLUMN_TEXT_ERRMSG, __LINE__);
            } else
                strncpy(tag->sn, (const char *) str, sizeof(tag->sn) - 1);
        }
        if (SQLITE_TEXT == sqlite3_column_type(stmt, 2))
        {
            str = sqlite3_column_text(stmt, 2);
            if (str == NULL) {
                dbg_syslog(LOG_ERR, COLUMN_TEXT_ERRMSG, __LINE__);
            } else {
                size_t tlen = strlen((const char *) str);
                tag->tag_node = calloc(1, tlen + 1);
                strcpy(tag->tag_node, (const char *) str);
                tag->tag_len = tlen;
            }
        }
        else if (SQLITE_BLOB == sqlite3_column_type(stmt, 2))
        {
            const void *tag_node = sqlite3_column_blob(stmt, 2);
            tag->tag_len = sqlite3_column_bytes(stmt, 2);
            tag->tag_node = (char *) malloc(tag->tag_len + 1);
            memcpy(tag->tag_node, tag_node, tag->tag_len);
            tag->tag_node[tag->tag_len] = '\0';
        }
        if (SQLITE_INTEGER == sqlite3_column_type(stmt, 3))
        {
            tag->time = sqlite3_column_int(stmt, 3);
        }
        if (SQLITE_TEXT == sqlite3_column_type(stmt, 4))
        {
            str = sqlite3_column_text(stmt, 4);
            if (str == NULL) {
                dbg_syslog(LOG_ERR, COLUMN_TEXT_ERRMSG, __LINE__);
            } else
                strcpy(tag->identifier, (const char *) str);
        }

        list_add_tail(&tag_list->list, query_list);
        dbg_syslog(LOG_DEBUG, "dc data get, id %d sn %s tag %s time %ld identifier %s", tag->id,
                  tag->sn, tag->tag_node, tag->time, tag->identifier);
    }

    return 0;

out:
    return -1;
}


int dr_dc_get_from_db(sqlite3 *db, char *sn, time_t start, time_t end, struct list_head *query_list)
{
	int ret = 0;
	sqlite3_stmt * stmt = NULL;	
    char *errorMessage = NULL;
	
	if (!db)
	{
		return -1;
	}
	
    sqlite3_exec(db, "BEGIN TRANSACTION", NULL, NULL, &errorMessage);
    if (errorMessage)
    {
        sqlite3_free(errorMessage);
    }

	if (!sn)
	{
		ret = sqlite3_prepare_v2(db, DR_DC_GET_BETWEEN_TIME_WITHOUT_SN_SQL, strlen(DR_DC_GET_BETWEEN_TIME_WITHOUT_SN_SQL), &stmt, 0);
		if (ret != SQLITE_OK)
		{
			dbg_syslog(LOG_ERR, "prepare error %s", sqlite3_errmsg(db));
			sqlite3_finalize(stmt);
			stmt = NULL;
			sqlite3_exec(db, "COMMIT TRANSACTION", NULL, NULL, &errorMessage);
			if (errorMessage)
			{
				sqlite3_free(errorMessage);
			}
			return ret;
		}
		dbg_syslog(LOG_DEBUG, "select by time, start:%ld end:%ld", start, end);
		sqlite3_bind_int(stmt, 1, start);
		sqlite3_bind_int(stmt, 2, end);
	}
	else
	{
		ret = sqlite3_prepare_v2(db, DR_DC_GET_BETWEEN_TIME_WITH_SN_SQL, strlen(DR_DC_GET_BETWEEN_TIME_WITH_SN_SQL), &stmt, 0);
		if (ret != SQLITE_OK)
		{
			dbg_syslog(LOG_ERR, "prepare error %s", sqlite3_errmsg(db));
			sqlite3_finalize(stmt);
			stmt = NULL;
			sqlite3_exec(db, "COMMIT TRANSACTION", NULL, NULL, &errorMessage);
			if (errorMessage)
			{
				sqlite3_free(errorMessage);
			}
			return ret;
		}
		dbg_syslog(LOG_DEBUG, "select by sn and time, sn:%s start:%ld end:%ld", sn, start, end);
		sqlite3_bind_text(stmt, 1, sn, strlen(sn), SQLITE_TRANSIENT);
		sqlite3_bind_int(stmt, 2, start);
		sqlite3_bind_int(stmt, 3, end);
	}
	
	dr_dc_db_get_select_all_from_stmt(stmt, query_list);
	sqlite3_reset(stmt);

    ret = sqlite3_exec(db, "COMMIT TRANSACTION", NULL, NULL, &errorMessage);
    if (errorMessage)
    {
        sqlite3_free(errorMessage);
    }
	
    return ret;
}


int dr_dc_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;
	time_t start = 0, end = 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);
	GET_JSON_VALUE_INT(root, "start", start);
	GET_JSON_VALUE_INT(root, "end", end);
    INIT_LIST_HEAD(&query_list);

    memset(topic, 0, sizeof(topic));
   // while (1)
   // {
        ret = dr_dc_get_from_db(g_dr_var_dc_t.db, NULL, start, end, &query_list);
        if (ret == -1)
        {
            dbg_syslog(LOG_ERR, "db get error");
            goto out;
        }
        else if (ret == 0)
        {
            if (!list_empty(&query_list))
            {
                dr_dc_tag_list_node_t *tmp = NULL;
                dr_dc_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, "dc_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_AddStringToObject(item, "identifier", tag_list->tag.identifier);						
                        cJSON_AddNumberToObject(item, "time", tag_list->tag.time);

                        cJSON_AddItemToArray(hist, item);

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

                char *str = (cJSON_Print(reply));
				//ipc/[GW_SN]/gateway/device/[GW_SN]/service/Get_Cl_HistoryData_reply"
                snprintf(topic, sizeof(topic), "ipc/%s/%s/device/%s/service/%s", var->sn_str, TOPIC_GATEWAY_CONFIG, var->sn_str, TOPIC_GET_CL_HISTORY_DATA_REPLY);
				dbg_syslog(LOG_DEBUG, "dr_dc_get_history_data: topic:%s, str:%s",topic, str);
                ret = ipc_session_publish(var->session, topic, (unsigned char *) str, strlen(str));
				if (ret != 0)
				{
					dbg_syslog(LOG_DEBUG, "ipc_session_publish dr dc failed");
				}
                free(str);

                cJSON_DeleteItemFromObject(reply, "dc_history_data");
            }
            else
            {
                dbg_syslog(LOG_DEBUG, "empty");
                goto out;
            }
        }
  //  }

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_dc_db_init(dr_var_dc_t *var)
{
	char set_max_page_count[64] = {0};
    const char *board_name = NULL;

	/*init db*/
	board_name = get_board_name();
    if (strcmp(board_name, BOARD_WOOLINK_MT7628) == 0)
    {
    	var->db = db_open_db(DR_DC_DB_TMPFS_NAME);
    }
	else
	{
		var->db = db_open_db(DR_DC_DB_FILE);
	}
	
    sqlite3_exec(var->db, "PRAGMA synchronous = NORMAL;", NULL, NULL, NULL);
    sqlite3_exec(var->db, "PRAGMA journal_mode = WAL;", NULL, NULL, NULL);
    sqlite3_exec(var->db, "PRAGMA journal_size_limit = 16384;", NULL, NULL, NULL);
	//sqlite3_exec(var->db, "PRAGMA auto_vacuum = FULL;", NULL, NULL, NULL);
    sprintf(set_max_page_count, "PRAGMA max_page_count = %d;", (3 * 1024 * 1024) / 4);
    sqlite3_exec(var->db, set_max_page_count, NULL, NULL, NULL);
    dr_dc_db_create_table(var->db, DC_DATA_TABLE_NAME, DC_TAG_TABLE_CREATE);
    return 0;
}


static int dr_dc_load_cfg(dr_var_dc_t * dc_var)
{
    char *data = NULL;
    cJSON *root = NULL;
    int ret = 0;
	int i = 0;
	cJSON *filter = NULL;

    data = read_file_data(DR_DC_CFG_FILE);
    if (data == NULL)
    {
        dbg_syslog(LOG_ERR, "read cfg file %s error", DR_DC_CFG_FILE);
        ret = -1;
        goto out;
    }
    root = cJSON_Parse(data);
    if (!root)
    {
        dbg_syslog(LOG_ERR, "parse cfg file %s error", DR_DC_CFG_FILE);
        ret = -1;
        goto out;
    }
	
	GET_JSON_VALUE_INT(root, "limit_time", dc_var->limit_time);

	cJSON *dc_cfg = cJSON_GetObjectItem(root, "identifier_list");
    if (!dc_cfg)
    {
        ret = -1;
		goto out;
    }

    int count = cJSON_GetArraySize(dc_cfg);
	if (count > MAX_DC_NUM)
	{
		//return -1;
		count = MAX_DC_NUM;
	}

	for (i = 0; i < count; i++)
	{
		filter = cJSON_GetArrayItem(dc_cfg, i);
	#if 0
		GET_JSON_VALUE_STRING(filter, "sn", cl_var->var_cl_s[i].sn_str);

		if (strlen(cl_var->var_cl_s[i].sn_str) <= 0)
		{
			continue;
		}
	#endif
		GET_JSON_VALUE_STRING(filter, "identifier", dc_var->var_dc_s[i].identifier);
		if (strlen(dc_var->var_dc_s[i].identifier) <= 0)
		{
			continue;
		}
		//GET_JSON_VALUE_STRING(filter, "tag", cl_var->var_cl_s[i].tag);
		
		dbg_syslog(LOG_DEBUG, "dr_cl_load_cfg[%d] identifer:%s", i, dc_var->var_dc_s[i].identifier);
	}

out:
	if (root)
    	cJSON_Delete(root);

    return ret;	
	
}

static int dr_dc_init(dr_var_dc_t * dc_var)
{
	int ret = 0;
	ret = dr_dc_load_cfg(dc_var);	
	if (ret != 0)
	{
		dbg_syslog(LOG_ERR, "dr dc load cfg failed");
		return -1;
	}
	
	ret = dr_dc_db_init(dc_var);
	if (ret != 0)
	{
		dbg_syslog(LOG_ERR, "dr dc db init failed");
		return -1;
	}

	return 0;

}

static int dr_period_property_load_conf(dr_var_t *var, char *file)
{
    int ret = 0;

    char *json_str = NULL;
    json_str = read_file_data(file);
    if (!json_str)
    {
        return -1;
    }

    cJSON *root = cJSON_Parse(json_str);
    if (!root)
    {
        return -1;
    }

    GET_JSON_VALUE_INT(root, "period_report", var->period_property_info.period_report);
    GET_JSON_VALUE_INT(root, "hourly_report", var->period_property_info.hourly_report);
    dbg_syslog(LOG_DEBUG, "period_report %d hourly_report %d", var->period_property_info.period_report, var->period_property_info.hourly_report);

    cJSON_Delete(root);
    return ret;
}

static int dr_init(dr_var_t *var)
{
    int ret = 0;
    int i = 0, j = 0;
    dev_sn_node_t *dev_sn_node = NULL;

    get_board_sn(var->sn_str);
    get_process_name(var->proc_name);
    if (load_nodes_cfg(&var->nodes_cfg_table, NODES_CFG_PATH) == -1)
    {
        dbg_syslog(LOG_ERR, "load nodes cfg fail");
    }
    if (load_objects_cfg(&var->object_table, OBJECTS_CFG_PATH) == -1)
    {
        dbg_syslog(LOG_ERR, "load object cfg fail");
    }

    if (access(PERIOD_PROPERTY_CFG, F_OK) != -1)
    {
        dr_period_property_load_conf(var, PERIOD_PROPERTY_CFG);
    }

    for (i = 0; i < HASH_SIZE; i++)
    {
        INIT_HLIST_HEAD(&var->dev_sn_table[i]);
    }

    for (i = 0; i < var->nodes_cfg_table->node_cnt; i++)
    {
        node_cfg_t *node = &var->nodes_cfg_table->node[i];
        int hash_index = hash_string(node->sn);
        dev_sn_node = calloc(1, sizeof(dev_sn_node_t));
        strncpy(dev_sn_node->dev_sn, node->sn, SN_MAX_LEN - 1);
        for (j = 0; j < HASH_SIZE; j++)
        {
            INIT_HLIST_HEAD(&dev_sn_node->identifier_table[j]);
        }
        hlist_add_head(&dev_sn_node->list, &var->dev_sn_table[hash_index]);
    }

    INIT_LIST_HEAD(&var->condition_list);

    dr_condition_init(var);
    dr_db_init(var);

	dr_dc_init(&g_dr_var_dc_t);

    dr_mqtt_client_init(var);

    var->report_timer = my_timer_create();
    if (var->report_timer > 0)
    {
        my_timer_set(var->report_timer, 10, PERIOD_TIMER_INTERVAL);
    }
    else
    {
        dbg_syslog(LOG_ERR, "timer create failed");
        ret = -1;
    }

    if (var->period_property_info.period_report || var->period_property_info.hourly_report)
    {
        var->period_property_timer = my_timer_create();
        if (var->period_property_timer > 0)
        {
            my_timer_set(var->period_property_timer, 1, 1000);
        }
        else
        {
            dbg_syslog(LOG_ERR, "timer create failed");
            ret = -1;
        }

        pthread_t thread_period;
        pthread_create(&thread_period, NULL, dr_period_property_loop, (void *)var);
    }

    dbg_syslog(LOG_INFO, "init done, board SN:%s", var->sn_str);
    return ret;
}

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

    openlog("DATA REPORTER", LOG_PID, LOG_DAEMON);
    lnxall_loglevel_set(LNXALL_LOGNOTICE, 1);

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

    dr_init(&var);
    dr_loop(&var);

    return 0;
}
