#include <sys/select.h>
#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 <pthread.h>

#include "aiot_state_api.h"
#include "aiot_sysdep_api.h"
#include "aiot_mqtt_api.h"
#include "aiot_dm_api.h"

#include "alink_common.h"

alink_var_t g_alink_var;

/*
 * 这个例程适用于`Linux`这类支持pthread的POSIX设备, 它演示了用SDK配置MQTT参数并建立连接, 之后创建2个线程
 *
 * + 一个线程用于保活长连接
 * + 一个线程用于接收消息, 并在有消息到达时进入默认的数据回调, 在连接状态变化时进入事件回调
 *
 * 接着演示了在MQTT连接上进行属性上报, 事件上报, 以及处理收到的属性设置, 服务调用, 取消这些代码段落的注释即可观察运行效果
 *
 * 需要用户关注或修改的部分, 已经用 TODO 在注释中标明
 *
 */

/* 位于portfiles/aiot_port文件夹下的系统适配函数集合 */
extern aiot_sysdep_portfile_t g_aiot_sysdep_portfile;

/* 位于external/ali_ca_cert.c中的服务器证书 */
extern const char *ali_ca_cert;

/* TODO: 如果要关闭日志, 就把这个函数实现为空, 如果要减少日志, 可根据code选择不打印
 *
 * 例如: [1577589489.033][LK-0317] mqtt_basic_demo&a13FN5TplKq
 *
 * 上面这条日志的code就是0317(十六进制), code值的定义见core/aiot_state_api.h
 *
 */

/* 日志回调函数, SDK的日志会从这里输出 */
int32_t alink_state_logcb(int32_t code, char *message)
{
    //dy_syslog(LOG_DEBUG, "%s", message);
    return 0;
}

/* MQTT事件回调函数, 当网络连接/重连/断开时被触发, 事件定义见core/aiot_mqtt_api.h */
void alink_mqtt_event_handler(void *handle, const aiot_mqtt_event_t *event, void *userdata)
{
    switch (event->type) {
        /* SDK因为用户调用了aiot_mqtt_connect()接口, 与mqtt服务器建立连接已成功 */
        case AIOT_MQTTEVT_CONNECT: {
            dy_syslog(LOG_DEBUG, "AIOT_MQTTEVT_CONNECT");
        }
        break;

        /* SDK因为网络状况被动断连后, 自动发起重连已成功 */
        case AIOT_MQTTEVT_RECONNECT: {
            dy_syslog(LOG_DEBUG, "AIOT_MQTTEVT_RECONNECT");
        }
        break;

        /* SDK因为网络的状况而被动断开了连接, network是底层读写失败, heartbeat是没有按预期得到服务端心跳应答 */
        case AIOT_MQTTEVT_DISCONNECT: {
            char *cause = (event->data.disconnect == AIOT_MQTTDISCONNEVT_NETWORK_DISCONNECT) ? ("network disconnect") :
                          ("heartbeat disconnect");
            dy_syslog(LOG_DEBUG, "AIOT_MQTTEVT_DISCONNECT: %s", cause);
        }
        break;

        default: {

        }
    }
}

/* 执行aiot_mqtt_process的线程, 包含心跳发送和QoS1消息重发 */
void *alink_mqtt_process_thread(void *args)
{
    int32_t res = STATE_SUCCESS;
	iotx_dev_meta_info_t *pmeta = (iotx_dev_meta_info_t *)args;

    while (pmeta->thread_running) {
        res = aiot_mqtt_process(pmeta->mqtt_handle);
        if (res == STATE_USER_INPUT_EXEC_DISABLED) {
            break;
        }
        sleep(1);
    }
    return NULL;
}

/* 执行aiot_mqtt_recv的线程, 包含网络自动重连和从服务器收取MQTT消息 */
void *alink_mqtt_recv_thread(void *args)
{
    int32_t res = STATE_SUCCESS;
	iotx_dev_meta_info_t *pmeta = (iotx_dev_meta_info_t *)args;

    while (pmeta->thread_running) {
        res = aiot_mqtt_recv(pmeta->mqtt_handle);
        if (res < STATE_SUCCESS) {
            if (res == STATE_USER_INPUT_EXEC_DISABLED) {
                break;
            }
            sleep(1);
        }
    }
    return NULL;
}

/* 用户数据接收处理回调函数 */
static void alink_dm_recv_handler(void *dm_handle, const aiot_dm_recv_t *recv, void *userdata)
{
	iotx_dev_meta_info_t *pmeta = (iotx_dev_meta_info_t *)userdata;
	
    switch (recv->type) {

        /* 属性上报, 事件上报, 获取期望属性值或者删除期望属性值的应答 */
        case AIOT_DMRECV_GENERIC_REPLY: {
			pmeta->post_cnt--;
            dy_syslog(LOG_DEBUG, "sn = %s post_cnt %d msg_id = %d, code = %d, data = %.*s, message = %.*s",
				   pmeta->alias_sn,pmeta->post_cnt,
                   recv->data.generic_reply.msg_id,
                   recv->data.generic_reply.code,
                   recv->data.generic_reply.data_len,
                   recv->data.generic_reply.data,
                   recv->data.generic_reply.message_len,
                   recv->data.generic_reply.message);
			if(pmeta->post_cnt < 0)
				pmeta->post_cnt = 0;
        }
        break;

        /* 属性设置 */
        case AIOT_DMRECV_PROPERTY_SET: {
            dy_syslog(LOG_DEBUG, "msg_id = %ld, params = %.*s",
                   (unsigned long)recv->data.property_set.msg_id,
                   recv->data.property_set.params_len,
                   recv->data.property_set.params);

            /* TODO: 以下代码演示如何对来自云平台的属性设置指令进行应答, 用户可取消注释查看演示效果 */
            /*
            {
                aiot_dm_msg_t msg;

                memset(&msg, 0, sizeof(aiot_dm_msg_t));
                msg.type = AIOT_DMMSG_PROPERTY_SET_REPLY;
                msg.data.property_set_reply.msg_id = recv->data.property_set.msg_id;
                msg.data.property_set_reply.code = 200;
                msg.data.property_set_reply.data = "{}";
                int32_t res = aiot_dm_send(dm_handle, &msg);
                if (res < 0) {
                    printf("aiot_dm_send failed");
                }
            }
            */
        }
        break;

        /* 异步服务调用 */
        case AIOT_DMRECV_ASYNC_SERVICE_INVOKE: {
            dy_syslog(LOG_DEBUG, "msg_id = %ld, service_id = %s, params = %.*s",
                   (unsigned long)recv->data.async_service_invoke.msg_id,
                   recv->data.async_service_invoke.service_id,
                   recv->data.async_service_invoke.params_len,
                   recv->data.async_service_invoke.params);

            /* TODO: 以下代码演示如何对来自云平台的异步服务调用进行应答, 用户可取消注释查看演示效果
             *
             * 注意: 如果用户在回调函数外进行应答, 需要自行保存msg_id, 因为回调函数入参在退出回调函数后将被SDK销毁, 不可以再访问到
             */
            {
				cJSON* root=cJSON_Parse(recv->data.async_service_invoke.params);
		        if (!root)
				{
					dy_syslog(LOG_WARNING, "cJSON_Parse payload failed");
					return ;
		        }

				cJSON_AddStringToObject(root, "sn", pmeta->alias_sn);
				cJSON_AddStringToObject(root, "identifier", recv->data.async_service_invoke.service_id);
				cJSON_AddNumberToObject(root, "mi", (unsigned long)recv->data.async_service_invoke.msg_id);

				char topic[TOPIC_MAX_LEN] = {0};
				char *payload = cJSON_PrintUnformatted(root);
				snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/%s/device/%s/data/%s", g_alink_var.sn_str, "alink", pmeta->alias_sn, TOPIC_EVT_SET_RGLT);
				dy_syslog(LOG_DEBUG, "topic:%s payload %d %s", topic,strlen(payload),payload);
				ipc_session_publish(g_alink_var.session, topic, (unsigned char*)payload, strlen(payload));

				//store the identifier
			    {
			        char key_str[64] = {0};
			        regulate_cmd_backup_t fields;

			        strcpy(fields.src_identifier, recv->data.async_service_invoke.service_id);
			        strcpy(fields.sn, pmeta->alias_sn);
			        fields.mi = recv->data.async_service_invoke.msg_id;

			        sprintf(key_str, "%s_%u", pmeta->alias_sn, (unsigned int)recv->data.async_service_invoke.msg_id);
					dy_syslog(LOG_DEBUG, "insert key_str:%s", key_str);
			        g_alink_var.identifier_backup.insert(&g_alink_var.identifier_backup, key_str, &fields, sizeof(fields));
			    }

				free(payload);
				cJSON_Delete(root);
            }
        }
        break;

        /* 同步服务调用 */
        case AIOT_DMRECV_SYNC_SERVICE_INVOKE: {
            dy_syslog(LOG_DEBUG, "msg_id = %ld, rrpc_id = %s, service_id = %s, params = %.*s",
                   (unsigned long)recv->data.sync_service_invoke.msg_id,
                   recv->data.sync_service_invoke.rrpc_id,
                   recv->data.sync_service_invoke.service_id,
                   recv->data.sync_service_invoke.params_len,
                   recv->data.sync_service_invoke.params);

            /* TODO: 以下代码演示如何对来自云平台的同步服务调用进行应答, 用户可取消注释查看演示效果
             *
             * 注意: 如果用户在回调函数外进行应答, 需要自行保存msg_id和rrpc_id字符串, 因为回调函数入参在退出回调函数后将被SDK销毁, 不可以再访问到
             */

            /*
            {
                aiot_dm_msg_t msg;

                memset(&msg, 0, sizeof(aiot_dm_msg_t));
                msg.type = AIOT_DMMSG_SYNC_SERVICE_REPLY;
                msg.data.sync_service_reply.rrpc_id = recv->data.sync_service_invoke.rrpc_id;
                msg.data.sync_service_reply.msg_id = recv->data.sync_service_invoke.msg_id;
                msg.data.sync_service_reply.code = 200;
                msg.data.sync_service_reply.service_id = "SetLightSwitchTimer";
                msg.data.sync_service_reply.data = "{}";
                int32_t res = aiot_dm_send(dm_handle, &msg);
                if (res < 0) {
                    printf("aiot_dm_send failed");
                }
            }
            */
        }
        break;

        /* 下行二进制数据 */
        case AIOT_DMRECV_RAW_DATA: {
            dy_syslog(LOG_DEBUG, "raw data len = %d", recv->data.raw_data.data_len);
            /* TODO: 以下代码演示如何发送二进制格式数据, 若使用需要有相应的数据透传脚本部署在云端 */
            /*
            {
                aiot_dm_msg_t msg;
                uint8_t raw_data[] = {0x01, 0x02};

                memset(&msg, 0, sizeof(aiot_dm_msg_t));
                msg.type = AIOT_DMMSG_RAW_DATA;
                msg.data.raw_data.data = raw_data;
                msg.data.raw_data.data_len = sizeof(raw_data);
                aiot_dm_send(dm_handle, &msg);
            }
            */
        }
        break;

        /* 二进制格式的同步服务调用, 比单纯的二进制数据消息多了个rrpc_id */
        case AIOT_DMRECV_RAW_SYNC_SERVICE_INVOKE: {
            dy_syslog(LOG_DEBUG, "raw sync service rrpc_id = %s, data_len = %d",
                   recv->data.raw_service_invoke.rrpc_id,
                   recv->data.raw_service_invoke.data_len);
        }
        break;

        default:
            break;
    }
}

/* 属性上报函数演示 */
int32_t alink_send_property_post(void *dm_handle, char *params)
{
    aiot_dm_msg_t msg;

    memset(&msg, 0, sizeof(aiot_dm_msg_t));
    msg.type = AIOT_DMMSG_PROPERTY_POST;
    msg.data.property_post.params = params;

    return aiot_dm_send(dm_handle, &msg);
}

/* 事件上报函数演示 */
int32_t alink_send_event_post(void *dm_handle, char *event_id, char *params)
{
    aiot_dm_msg_t msg;

    memset(&msg, 0, sizeof(aiot_dm_msg_t));
    msg.type = AIOT_DMMSG_EVENT_POST;
    msg.data.event_post.event_id = event_id;
    msg.data.event_post.params = params;

    return aiot_dm_send(dm_handle, &msg);
}

/* 演示了获取属性LightSwitch的期望值, 用户可将此函数加入到main函数中运行演示 */
int32_t alink_send_get_desred_requset(void *dm_handle)
{
    aiot_dm_msg_t msg;

    memset(&msg, 0, sizeof(aiot_dm_msg_t));
    msg.type = AIOT_DMMSG_GET_DESIRED;
    msg.data.get_desired.params = "[\"LightSwitch\"]";

    return aiot_dm_send(dm_handle, &msg);
}

/* 演示了删除属性LightSwitch的期望值, 用户可将此函数加入到main函数中运行演示 */
int32_t alink_send_delete_desred_requset(void *dm_handle)
{
    aiot_dm_msg_t msg;

    memset(&msg, 0, sizeof(aiot_dm_msg_t));
    msg.type = AIOT_DMMSG_DELETE_DESIRED;
    msg.data.get_desired.params = "{\"LightSwitch\":{}}";

    return aiot_dm_send(dm_handle, &msg);
}

int alink_publish(alink_var_t *var, iotx_dev_meta_info_t *pmeta, int type, tag_table_t *tag)
{
	if(type == 0)
	{
		pmeta->post_cnt++;
		dy_syslog(LOG_DEBUG, "%s type %d post_cnt %d params %s", __FUNCTION__, type, pmeta->post_cnt, tag->tag_node);
		alink_send_property_post(pmeta->dm_handle, tag->tag_node);
		if(pmeta->post_cnt > 2)
		{
			dy_syslog(LOG_WARNING, "%s post msg error post_cnt %d !!!", __FUNCTION__, pmeta->post_cnt);
			system("/etc/init.d/alink restart");
		}
	}
	else if(type == 1)
	{
		dy_syslog(LOG_DEBUG, "%s type %d event_id %s params %s ", __FUNCTION__, type, tag->identifier, tag->tag_node);
		alink_send_event_post(pmeta->dm_handle, tag->identifier, tag->tag_node);
	}
	else if(type == 2)
	{
		aiot_dm_msg_t msg;
		regulate_cmd_backup_t *fields = NULL;

		//restore the identifier
	    {
	        char key_str[64] = {0};
	        data_info_t *data_info = NULL;

	        sprintf(key_str, "%s_%u", pmeta->alias_sn, tag->mi);
	        data_info = var->identifier_backup.get(&var->identifier_backup, key_str);
	        if (data_info)
	        {
	            fields = (regulate_cmd_backup_t *)data_info->data;
				dy_syslog(LOG_DEBUG, "get key_str:%s", key_str);
	        }
	    }
		
		if(fields)
		{
	        memset(&msg, 0, sizeof(aiot_dm_msg_t));
	        msg.type = AIOT_DMMSG_ASYNC_SERVICE_REPLY;
	        msg.data.async_service_reply.msg_id = tag->mi;
	        msg.data.async_service_reply.code = 200;
	        msg.data.async_service_reply.service_id = tag->identifier;
	        msg.data.async_service_reply.data = tag->tag_node;

			dy_syslog(LOG_DEBUG, "%s type %d service_id %s params %s", __FUNCTION__, type, tag->identifier, tag->tag_node);
	        int res = aiot_dm_send(pmeta->dm_handle, &msg);
	        if (res < 0) {
	            dy_syslog(LOG_DEBUG, "aiot_dm_send service_id %s failed!!!", tag->identifier);
	        }
		}
	}
	
    return 0;
}

int alink_mqtt_connect(iotx_dev_meta_info_t *pmeta)
{
	int res = STATE_SUCCESS;
	unsigned short port = 443;
	char *url = "iot-as-mqtt.cn-shanghai.aliyuncs.com"; /* 阿里云平台上海站点的域名后缀 */
	char host[100] = {0}; /* 用这个数组拼接设备连接的云平台站点全地址, 规则是 ${productKey}.iot-as-mqtt.cn-shanghai.aliyuncs.com */
	aiot_sysdep_network_cred_t cred; /* 安全凭据结构体, 如果要用TLS, 这个结构体中配置CA证书等参数 */
	
	/* 创建SDK的安全凭据, 用于建立TLS连接 */
	memset(&cred, 0, sizeof(aiot_sysdep_network_cred_t));
	cred.option = AIOT_SYSDEP_NETWORK_CRED_SVRCERT_CA;	/* 使用RSA证书校验MQTT服务端 */
	cred.max_tls_fragment = 16384; /* 最大的分片长度为16K, 其它可选值还有4K, 2K, 1K, 0.5K */
	cred.sni_enabled = 1;								/* TLS建连时, 支持Server Name Indicator */
	cred.x509_server_cert = ali_ca_cert;				 /* 用来验证MQTT服务端的RSA根证书 */
	cred.x509_server_cert_len = strlen(ali_ca_cert);	 /* 用来验证MQTT服务端的RSA根证书长度 */
		
    /* 创建1个MQTT客户端实例并内部初始化默认参数 */
    pmeta->mqtt_handle = aiot_mqtt_init();
    if (pmeta->mqtt_handle == NULL) {
        dy_syslog(LOG_WARNING, "sn %s aiot_mqtt_init failed!!!", pmeta->device_name);
        return -1;
    }

    snprintf(host, 100, "%s.%s", pmeta->product_key, url);
    /* 配置MQTT服务器地址 */
    aiot_mqtt_setopt(pmeta->mqtt_handle, AIOT_MQTTOPT_HOST, (void *)host);
    /* 配置MQTT服务器端口 */
    aiot_mqtt_setopt(pmeta->mqtt_handle, AIOT_MQTTOPT_PORT, (void *)&port);
    /* 配置设备productKey */
    aiot_mqtt_setopt(pmeta->mqtt_handle, AIOT_MQTTOPT_PRODUCT_KEY, pmeta->product_key);
    /* 配置设备deviceName */
    aiot_mqtt_setopt(pmeta->mqtt_handle, AIOT_MQTTOPT_DEVICE_NAME, pmeta->device_name);
    /* 配置设备deviceSecret */
    aiot_mqtt_setopt(pmeta->mqtt_handle, AIOT_MQTTOPT_DEVICE_SECRET, pmeta->device_secret);
    /* 配置网络连接的安全凭据, 上面已经创建好了 */
    aiot_mqtt_setopt(pmeta->mqtt_handle, AIOT_MQTTOPT_NETWORK_CRED, (void *)&cred);
    /* 配置MQTT事件回调函数 */
    aiot_mqtt_setopt(pmeta->mqtt_handle, AIOT_MQTTOPT_EVENT_HANDLER, (void *)alink_mqtt_event_handler);

    /* 创建DATA-MODEL实例 */
    pmeta->dm_handle = aiot_dm_init();
    if (pmeta->dm_handle == NULL) {
        dy_syslog(LOG_WARNING, "sn %s aiot_dm_init failed!!!", pmeta->device_name);
        return -1;
    }
    /* 配置MQTT实例句柄 */
    aiot_dm_setopt(pmeta->dm_handle, AIOT_DMOPT_MQTT_HANDLE, pmeta->mqtt_handle);
    /* 配置消息接收处理回调函数 */
    aiot_dm_setopt(pmeta->dm_handle, AIOT_DMOPT_USERDATA, (void *)pmeta);
	aiot_dm_setopt(pmeta->dm_handle, AIOT_DMOPT_RECV_HANDLER, (void *)alink_dm_recv_handler);

    /* 与服务器建立MQTT连接 */
    res = aiot_mqtt_connect(pmeta->mqtt_handle);
    if (res < STATE_SUCCESS) {
        /* 尝试建立连接失败, 销毁MQTT实例, 回收资源 */
        aiot_mqtt_deinit(&pmeta->mqtt_handle);
        dy_syslog(LOG_WARNING, "sn %s aiot_mqtt_connect failed: -0x%04X\n", pmeta->device_name, -res);
        return -1;
    }

	pmeta->thread_running = 1;

    /* 创建一个单独的线程, 专用于执行aiot_mqtt_process, 它会自动发送心跳保活, 以及重发QoS1的未应答报文 */
    res = pthread_create(&pmeta->mqtt_process_thread, NULL, alink_mqtt_process_thread, pmeta);
    if (res < 0) {
        dy_syslog(LOG_WARNING, "pthread_create alink_mqtt_process_thread failed: %d\n", res);
        return -1;
    }

    /* 创建一个单独的线程用于执行aiot_mqtt_recv, 它会循环收取服务器下发的MQTT消息, 并在断线时自动重连 */
    res = pthread_create(&pmeta->mqtt_recv_thread, NULL, alink_mqtt_recv_thread, pmeta);
    if (res < 0) {
        dy_syslog(LOG_WARNING, "pthread_create alink_mqtt_recv_thread failed: %d\n", res);
        return -1;
    }
	
    return 0;
}

static void alink_node_connect(alink_var_t *var)
{
	int i;

	/* 配置SDK的底层依赖 */
	aiot_sysdep_set_portfile(&g_aiot_sysdep_portfile);
	/* 配置SDK的日志输出 */
	aiot_state_set_logcb(alink_state_logcb);

	var->link_node_cnt = 0;
	for(i=0; i<var->nodes_cfg_table->node_cnt; i++)
    {
    	if(strlen(var->nodes_cfg_table->node[i].device_secret))
		{
			linkkit_info_list_t* linkkit = calloc(sizeof(linkkit_info_list_t), 1);

			strncpy(linkkit->meta_info.alias_sn, var->nodes_cfg_table->node[i].sn, sizeof(linkkit->meta_info.alias_sn));
			strncpy(linkkit->meta_info.device_name, var->nodes_cfg_table->node[i].device_name, sizeof(linkkit->meta_info.device_name));
			strncpy(linkkit->meta_info.device_secret, var->nodes_cfg_table->node[i].device_secret, sizeof(linkkit->meta_info.device_secret));
			strncpy(linkkit->meta_info.product_key, var->nodes_cfg_table->node[i].product_key, sizeof(linkkit->meta_info.product_key));
			
   			if(!strlen(linkkit->meta_info.device_secret))
				continue;

			if(!strlen(linkkit->meta_info.device_name))
				strncpy(linkkit->meta_info.device_name, linkkit->meta_info.alias_sn, sizeof(linkkit->meta_info.device_name));
			
			var->link_node_cnt++;
	        dy_syslog(LOG_DEBUG,"===%s %d sn %s (%s %s %s)===\n",__FUNCTION__,__LINE__,linkkit->meta_info.alias_sn,linkkit->meta_info.device_name,linkkit->meta_info.device_secret,linkkit->meta_info.product_key);
	        alink_mqtt_connect(&linkkit->meta_info);
	        list_add_tail(&linkkit->list,&var->node_list);
		}
    }
}

static int mqtt_msg_process(alink_var_t *var, int type, ipc_msg_t *mqtt_msg)
{
	int ret;
	tag_table_t tag;

	//先检查参数
	ret = get_tag_data_from_str(&tag, mqtt_msg->payload);
    if (ret < 0)
    {
        dy_syslog(LOG_ERR, "get data structure failed");
        return -1;
    }
	
	linkkit_info_list_t *linkkit = NULL;
	list_for_each_entry(linkkit, &var->node_list, list)
	{
		if(strcmp(linkkit->meta_info.device_name, tag.sn) == 0 || strcmp(linkkit->meta_info.alias_sn, tag.sn) == 0)
		{
			if (tag.tag_node == NULL)
			{
		        dy_syslog(LOG_WARNING,"type %d tag_node is NULL!!!",type);
		        break;
		    }
			alink_publish(var, &linkkit->meta_info, type, &tag);
			break;
		}
	}

	return 0;
}

void msg_mqtt_recv(alink_var_t *var, int type, ipc_msg_t *mqtt_msg)
{
	data_list_t* data = calloc(1,sizeof(data_list_t)+mqtt_msg->payloadLen+1);	
	if(data)
	{
		strncpy(data->mqtt.topic, mqtt_msg->topic, sizeof(data->mqtt.topic));
		data->type = type;
		data->mqtt.payloadLen = mqtt_msg->payloadLen;
		memcpy(data->mqtt.payload, mqtt_msg->payload, data->mqtt.payloadLen);
		var->data_cnt++;
		//dy_syslog(LOG_DEBUG, "++%s data_cnt %d\n",__FUNCTION__,var->data_cnt);
    	list_add_tail(&data->list,&var->data_list);
	}

	var->data_flag = 1;
}

static void mqtt_msg_process_loop(alink_var_t *var)
{
	data_list_t* data = NULL;	
	list_for_each_entry(data, &var->data_list, list)
	{
		//dy_syslog(LOG_DEBUG, "var->data_cnt %d", var->data_cnt);
		var->data_cnt--;
		mqtt_msg_process(var, data->type, &data->mqtt);
		list_del(&data->list);
		return;
	}

	var->data_cnt = 0;
	var->data_flag = 0;
}

static void alink_loop(alink_var_t *var)
{	
    while (1)
    {
		if(var->data_flag)
        {
			usleep(50000);
        }
		else
		{
			usleep(2000000);
		}

		mqtt_msg_process_loop(var);
    }
}

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

    snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/+/device/+/data_filtered/property/+", var->sn_str);
    ipc_session_subscribe(session, topic);

    snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/+/device/+/data_filtered/event/+", var->sn_str);
    ipc_session_subscribe(session, topic);

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

static int alink_mqtt_handle_recv_msg(void *obj, ipc_msg_t *msg)
{
    alink_var_t *var = (alink_var_t *)obj;
    unsigned char matched;
    int ret;
    dy_syslog(LOG_DEBUG, " MQTT client: received MQTT topic:%s payload length:%d",
            msg->topic, msg->payloadLen);

    ret = mosquitto_topic_matches_sub("ipc/+/+/device/+/data_filtered/property/+", msg->topic, &matched);
    if (ret == 0 && matched)
    {
        //实时数据
        msg_mqtt_recv(var, 0, msg);
    }
    else
    {
        //事件
        ret = mosquitto_topic_matches_sub("ipc/+/+/device/+/data_filtered/event/+", msg->topic, &matched);
        if (ret == 0 && matched)
        {
            msg_mqtt_recv(var, 1, msg);
        }
        else
        {
            ret = mosquitto_topic_matches_sub("ipc/+/+/device/+/data_filtered/service/+", msg->topic, &matched);
            if (ret == 0 && matched)
            {
                msg_mqtt_recv(var, 2, msg);
            }
        }
    }

	return 0;
}

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

    snprintf(clientId, MAX_CLIENT_ID_LEN, "INT_linkkit_gw_%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, alink_mqtt_handle_recv_msg, NULL);
    alink_subscribe_all(var);
    ipc_session_start(var->session);
}

int alink_init(alink_var_t *var)
{
    const char *board_name = NULL;

    check_make_dir(NODES_CACHE);
    check_make_dir(NODES_CFG);
    get_board_sn(var->sn_str);
    get_process_name(var->proc_name);
	
	INIT_LIST_HEAD(&var->node_list);
	INIT_LIST_HEAD(&var->data_list);
	
    board_name = get_board_name();

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

    if (load_nodes_cfg(&var->nodes_cfg_table, NODES_CFG_PATH) == -1)
    {
        dy_syslog(LOG_ERR, "load nodes cfg fail");
    }
	
    if (load_templates_cfg(&var->template_table, TEMPLATES_CFG_PATH) == -1)
    {
        dy_syslog(LOG_ERR, "load template cfg fail");
    }
	
    kv_array_init(&var->identifier_backup, 32);
	
	alink_mqtt_client_init(var);
	alink_node_connect(var);
	
    return 0;
}


int main(int argc, char *argv[])
{
    alink_var_t *var = &g_alink_var;

	dy_syslog(LOG_DEBUG, "\n\n\
            |********************************************|\n\
            |           alink start  X_X         |\n\
            |********************************************|\n");
    memset(var, 0, sizeof(alink_var_t));
    alink_init(var);
    alink_loop(var);

    return 0;
}
