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

pthread_mutex_t mutex;
huayun_var_t ghuayun_var;

static int publish_mqtt_message(huayun_var_t *var, char *data, int len)
{
	return;
}

static int iec104_connect(uint8_t *dstip, int port)
{
	struct sockaddr_in servaddr; 
    int socketfd = 0;
	
	if((socketfd = socket(AF_INET, SOCK_STREAM, 0)) < 0){  
        dy_syslog(LOG_ERR,"create socket error: serverip %s port %d %s(errno: %d)", dstip,port,strerror(errno),errno);  
        return -1;  
    } 

    memset(&servaddr, 0, sizeof(servaddr));  
    servaddr.sin_family = AF_INET;  
    servaddr.sin_port = htons(port);
    if(inet_pton(AF_INET, dstip, &servaddr.sin_addr) <= 0){  
        dy_syslog(LOG_ERR,"inet_pton error for %s port %d",dstip,port);
        close(socketfd);
        return -1;  
    }  

    if(connect(socketfd, (struct sockaddr*)&servaddr, sizeof(servaddr)) < 0){  
        dy_syslog(LOG_ERR,"connect error: serverip %s port %d %s(errno: %d)",dstip,port,strerror(errno),errno);  
        close(socketfd);
        return -1; 
    }

	fcntl(socketfd, F_SETFL, fcntl(socketfd, F_GETFL) | O_NONBLOCK);

	dy_syslog(LOG_DEBUG,"serverip %s port %d socketfd %d connect success...",dstip,port,socketfd);
	return socketfd;
}

void *iec104_task(void *arg){

    int socketfd = (*(int *)arg);
    while(1){
        //pthread_mutex_lock(&mutex2);
        Iec10x_Scheduled(socketfd, &ghuayun_var.connect); 
		if(ghuayun_var.connect == 0)
			break;
        Iec104_StateMachine();
        //pthread_mutex_lock(&mutex2);
    }
    return ((void *)0); 
}

void iec104_msg(huayun_var_t* var)
{
    int ret;
	char Iec104_RecvBuf[MAXLINE];  

    while(1)
    {
        ret = read(var->iec104_fd,Iec104_RecvBuf,MAXLINE);
        if(ret <= 0 || ret>1500)
        {
        	if(var->data_flag)
        		usleep(500000);
			else
				usleep(2000000);
            break;
        }
        
        /* mutex lock*/
        pthread_mutex_lock(&mutex);
		dy_syslog_hex(LOG_DEBUG, Iec104_RecvBuf,ret, "len %d                   网关<-HUAYUN", ret);
        Iex104_Receive(Iec104_RecvBuf, ret);
		var->iec104_alive = 1;
        pthread_mutex_unlock(&mutex);
    }
}

void iec104_pack_save(huayun_var_t* var, IEC104_PACK_T *pack, char *sn, time_t time_stamp)
{
	int i;
	tag_table_t tag = {0};
	
	if(pack->len == 0)
		return;

	strcpy(tag.sn, sn);
	tag.time = time_stamp;
	tag.tag_node = calloc(pack->len*2+1, 1);
	tag.report = 0;
	
	for(i=0; i<pack->len; i++)
	{
		char tmp[4] = {0};
		sprintf(tmp, "%02X", pack->pack[i]);
		strcat(tag.tag_node,tmp);
	}
	tag.tag_len = strlen(tag.tag_node);
	dy_syslog(LOG_DEBUG, "db insert tag sn:%s",sn);
	dy_db_session_insert_tag(var->db, &tag);
}

void iec104_pack_update(uint8_t *pack, uint16_t len)
{
	int i;
	huayun_var_t* var = &ghuayun_var;
	tag_table_t tag = {0};

	if(!var->connect)
	{
		dy_syslog(LOG_WARNING, "huayun disconnect !!!");
		return;
	}

	//数据发送失败，更新数据库标识为未发送
	tag.tag_node = calloc(len*2+1, 1);
	tag.report = 1;

	for(i=0; i<len; i++)
	{
		char tmp[4] = {0};
		sprintf(tmp, "%02X", pack[i]);
		strcat(tag.tag_node,tmp);
	}

	db_update_by_tag_content(var->db, &tag);
	free(tag.tag_node);
}

void iec104_pack_select(huayun_var_t* var)
{
	int i,cnt = 0;
	tag_list_node_t *tag_list;
	tag_list_node_t *tmp;

	if(!var->iec104_alive)
		return;
	
	//查找数据库未上报数据
	dy_db_session_select_not_repoted(var->db, 4);
	
    list_for_each_entry_safe(tag_list, tmp, &var->db->query_list,list)
    {
    	uint16_t len = strlen(tag_list->tag.tag_node);
    	uint8_t *buf = calloc(len+1, 1);
		uint8_t *p = tag_list->tag.tag_node;
		
		for(i=0; i<len/2; i++)
		{
			char tmp[4] = {0};
			sscanf(p+i*2,"%02X",&buf[i]);
			//dy_syslog(LOG_DEBUG," buf[i %d] %02X ",i,buf[i]);
		}
		dy_syslog(LOG_DEBUG,"select_not_repoted len %d tag_node %s",len/2,tag_list->tag.tag_node);
		//发送数据包
		IEC104_Send_Spont(buf, len/2);	
		list_del(&tag_list->list);
    }
}

void iec104_send_process(huayun_var_t* var, node_info_list_t *pnode, uint16_t asdu, int call, uint8_t reason, uint8_t ValueType, uint8_t Prio, uint8_t DevType)
{
	int i;	
	int start = 0, end = 0;
	uint16_t before_addr = 0;
	
	for(i=0; i<pnode->node_tag.tag_cnt; i++)
	{
		dy_syslog(LOG_DEBUG,"tag_cnt %d %d before_addr 0x%x info_addr 0x%x start %d end %d",pnode->node_tag.tag_cnt,i,before_addr,pnode->node_tag.tag_info[i].info_addr,start,end);
		IEC104_PACK_T pack = {0};
		if(pnode->node_tag.tag_info[i].info_addr == 0)
			break;
		if(i == 0)
		{
			start = 0;
		}
		else if(pnode->node_tag.tag_info[i].info_addr - before_addr == 1)
		{
			end = i;
		}
		else if(pnode->node_tag.tag_info[i].info_addr - before_addr > 1 || before_addr > pnode->node_tag.tag_info[i].info_addr)
		{
			if(before_addr > pnode->node_tag.tag_info[i].info_addr)
				end = i;
			dy_syslog(LOG_DEBUG,"Intervall info_addr 0x%x asdu 0x%x start %d end %d",pnode->node_tag.tag_info[start].info_addr, asdu, start, end+1);
			if(call)
				IEC104_BuildDetect(asdu, reason, ValueType, Prio, DevType, pnode->node_tag.tag_info[start].info_addr, start, end+1);
			else
			{
				IEC104_BuildDetectF_Spont(1, pnode->node_tag.tag_info[start].info_addr, asdu, start, end+1, &pack);
				iec104_pack_save(var, &pack, pnode->node_info.sn, pnode->node_tag.tag_info[start].time_stamp);
			}
			start = i;
			if(i == pnode->node_tag.tag_cnt-1)
			{
				end = i;
				dy_syslog(LOG_DEBUG,"Intervall 2 info_addr 0x%x asdu 0x%x start %d end %d",pnode->node_tag.tag_info[start].info_addr, asdu, start, end+1);
				if(call)
					IEC104_BuildDetect(asdu, reason, ValueType, Prio, DevType, pnode->node_tag.tag_info[start].info_addr, start, end+1);
				else
				{
					IEC104_BuildDetectF_Spont(1, pnode->node_tag.tag_info[start].info_addr, asdu, start, end+1, &pack);
					iec104_pack_save(var, &pack, pnode->node_info.sn, pnode->node_tag.tag_info[start].time_stamp);
				}
			}
			else
				end = i;
		}

		if(((i-start >= 16 || i == pnode->node_tag.tag_cnt-1) && start != i) || pnode->node_tag.tag_cnt == 1) 
		{
			end = i;
			dy_syslog(LOG_DEBUG,"info_addr 0x%x asdu 0x%x start %d end %d",pnode->node_tag.tag_info[start].info_addr, asdu, start, end+1);
			if(call)
				IEC104_BuildDetect(asdu, reason, ValueType, Prio, DevType, pnode->node_tag.tag_info[start].info_addr, start, end+1);
			else
			{
				IEC104_BuildDetectF_Spont(1, pnode->node_tag.tag_info[start].info_addr, asdu, start, end+1, &pack);
				iec104_pack_save(var, &pack, pnode->node_info.sn, pnode->node_tag.tag_info[start].time_stamp);
			}
			start = i;
		}
			
		before_addr = pnode->node_tag.tag_info[i].info_addr;
	}
}

void iec104_period(huayun_var_t* var)
{
	int i;
	node_info_list_t *pnode = NULL;
	node_info_list_t *tmp = NULL;
	
	TIMER_CONFIRM(var->period_timer);

	list_for_each_entry_safe(pnode, tmp, &var->node_list,list)
    {
    	uint16_t asdu = pnode->node_info.port << 8 | pnode->node_info.term_addr[0];
		dy_syslog(LOG_DEBUG,"--sn %s data %p addrV 0x%x port %d term_addr[0] %d",pnode->node_info.sn,pnode->node_tag.tag_info,asdu,pnode->node_info.port,pnode->node_info.term_addr[0]);
		if(!pnode->node_tag.tag_info)
			continue;
		iec104_send_process(var, pnode, asdu, 0, 0, 0, 0, 0);		
    }
}

void iec104_call_process(uint16_t addr, uint8_t reason, uint16_t ValueType, uint8_t Prio, uint8_t DevType)
{
    int ret;
	int i;
	node_info_list_t *pnode = NULL;
	node_info_list_t *tmp = NULL;

	uint8_t TimeFlag = 1;
	uint16_t infoAddr = IEC104_INFOADDR_VALUE_HXGF;
	list_for_each_entry_safe(pnode, tmp, &ghuayun_var.node_list, list)
    {
    	int disable=1;
		uint16_t asdu = pnode->node_info.port << 8 | pnode->node_info.term_addr[0];
		dy_syslog(LOG_DEBUG,"sn %s data %p addrV 0x%x port %d term_addr[0] %d",pnode->node_info.sn,pnode->node_tag.tag_info,asdu,pnode->node_info.port,pnode->node_info.term_addr[0]);
		if(!pnode->node_tag.tag_info || (addr != 0xFFFF && addr != asdu))
			continue;
		
    	for(i=0; i<pnode->node_tag.tag_cnt; i++)
    	{
    		if(pnode->node_tag.tag_info[i].quality)
    		{
				disable = 1;
				break;
    		}
    	}
		if(disable)
		{			
			uint16_t addrV = pnode->node_info.port << 8 | pnode->node_info.term_addr[0];
			dy_syslog(LOG_DEBUG,"sn %s data %p addrV 0x%x reason %d ValueType %d Prio %d DevType %d",pnode->node_info.sn,pnode->node_tag.tag_info,addrV,reason, ValueType, Prio, DevType);
			if(pnode->node_tag.tag_info)
				iec104_send_process(&ghuayun_var, pnode, addrV, 1, reason, ValueType, Prio, DevType);
		}

		if(addr == asdu)
			break;
    }
}

int iec104_get_station_cnt(uint16_t Addr, uint16_t DevType, uint8_t *num)
{
    char *paddr = &Addr;	
    node_info_list_t *pnode = NULL;
	node_info_list_t *tmp = NULL;
	
    list_for_each_entry_safe(pnode, tmp, &ghuayun_var.node_list,list)
    {
    	//dy_syslog(LOG_DEBUG,"sn %s,paddr[%d %d] port %d term_addr[0] %d Addr 0x%x",pnode->node_info.sn,paddr[0],paddr[1],pnode->node_info.port,pnode->node_info.term_addr[0],Addr);
        if((paddr[1] == pnode->node_info.port && paddr[0] == pnode->node_info.term_addr[0]) && pnode->node_tag.tag_info)
        {
			*num = pnode->node_tag.tag_cnt;
            return 0;
        }
    }

	dy_syslog(LOG_DEBUG,"can't find Addr 0x%x asdu_num",Addr);
	
    return 0;
}

int iec104_get_station_info(uint16_t Addr, uint16_t DevType, uint8_t n, tag_info_t *tag_info)
{
    char *paddr = &Addr;	
    node_info_list_t *pnode = NULL;
	node_info_list_t *tmp = NULL;
	
    list_for_each_entry_safe(pnode, tmp, &ghuayun_var.node_list,list)
    {
    	//dy_syslog(LOG_DEBUG,"sn %s Addr 0x%x paddr[%d %d] port %d term_addr[0] %d",pnode->node_info.sn,Addr,paddr[0],paddr[1],pnode->node_info.port,pnode->node_info.term_addr[0]);
        if((paddr[1] == pnode->node_info.port && paddr[0] == pnode->node_info.term_addr[0]) && pnode->node_tag.tag_info && n <= pnode->node_tag.tag_cnt)
        {
			//(*val) = pnode->node_tag.tag_info[n].val
			//(*qds) = pnode->node_tag.tag_info[n].quality;
			//(*time_stamp) = pnode->node_tag.tag_info[n].time_stamp;
			memcpy(tag_info, &pnode->node_tag.tag_info[n], sizeof(tag_info_t));
			return 0;
        }
    }

	dy_syslog(LOG_DEBUG,"can't find Addr 0x%x i %d",Addr,n);
	
    return 1;
}

static int iec104_init(huayun_var_t *var)
{
    int ret = 0;

	var->connect = 0;
	get_mac();
	Stm32f103RegisterIec10x();
	
	if(var->iec104_fd > 0)
	{
		close(var->iec104_fd);
		var->iec104_fd = 0;
	}

	var->iec104_fd = iec104_connect(var->server_ip, var->port);
	var->period_timer = -1;
	var->alive_timer = -1;
	
	ret=pthread_create(&var->iec104_task,NULL,iec104_task,&var->iec104_fd);
    if(ret!=0)
	{  
        dy_syslog(LOG_ERR,"pthread_create error:%s",strerror(ret));  
        return ret;  
    }
	Iex104_send_login();
	var->connect = 1;
	if(var->period_timer < 1)
    	var->period_timer = my_timer_create();
    if(var->period_timer > 0)
    	my_timer_set(var->period_timer, 1, var->period*1000);

	if(var->alive_timer < 1)
    	var->alive_timer = my_timer_create();
    if(var->alive_timer > 0)
    	my_timer_set(var->alive_timer, 60, 60*1000);

    return ret;
}

static void iec104_alive(huayun_var_t* var)
{
	int i;
	node_info_list_t *pnode = NULL;
	node_info_list_t *tmp = NULL;
	
	TIMER_CONFIRM(var->alive_timer);
	dy_syslog(LOG_DEBUG, "iec104_alive %d", var->iec104_alive);
	if(var->iec104_alive == 0)
	{
		dy_syslog(LOG_WARNING, "iec104_alive %d", var->iec104_alive);
		iec104_init(var);
	}
	
	var->iec104_alive = 0;
}

static int huayun_parser_tags(huayun_var_t *var, node_info_list_t *pnode)
{
    int i,j,k;
	
    for (i = 0; i < var->object_table->object_cnt; i++)
    {
        if (strcmp(pnode->node_info.product_key, var->object_table->object[i].object_model_id) == 0)
        {
        	object_cfg_t* pobject = &var->object_table->object[i];
			
        	for(j=0; j<pobject->event_tab->eventCnt; j++)
        	{
        		if(strcmp("post", pobject->event_tab->event[j].identifier) == 0)
				{
					pnode->node_tag.tag_cnt = pobject->event_tab->event[j].poutput->propertyCnt;
					pnode->node_tag.tag_info = calloc(sizeof(tag_info_t), pnode->node_tag.tag_cnt);

					for(k=0; k<pnode->node_tag.tag_cnt; k++)
					{
						pnode->node_tag.tag_info[k].quality = 0x80;
						strncpy(pnode->node_tag.tag_info[k].identifier, pobject->event_tab->event[j].poutput->property[k].identifier, sizeof(pnode->node_tag.tag_info[k].identifier));
						
						char *iec104 = strstr(pnode->node_tag.tag_info[k].identifier, IEC104_STR);
						if(iec104 && strlen(iec104) > strlen(IEC104_STR))
						{
							iec104 += strlen(IEC104_STR);
							sscanf(iec104, "%d", &pnode->node_tag.tag_info[k].info_addr);
						}
						dy_syslog(LOG_DEBUG,"tag_cnt %d %d identifier %s info_addr %d quality %d",pnode->node_tag.tag_cnt,k,pnode->node_tag.tag_info[k].identifier,pnode->node_tag.tag_info[k].info_addr,pnode->node_tag.tag_info[k].quality);
					}
					break;
        		}
        	}
			
            break;
        }
    }

    return 0;
}

void huayun_node_parser(huayun_var_t *var)
{
	int k;
	
	for (k = 0; k < var->nodes_cfg_table->node_cnt; k++)
    {
    	node_info_list_t* pnode_list = calloc(sizeof(node_info_list_t), 1);
    	memcpy(&pnode_list->node_info, &var->nodes_cfg_table->node[k], sizeof(node_cfg_t));

		dy_syslog(LOG_DEBUG, "===sn %s port %d===", var->nodes_cfg_table->node[k].sn,var->nodes_cfg_table->node[k].port);		
		huayun_parser_tags(var, pnode_list);       
        list_add_tail(&pnode_list->list,&var->node_list);
    }
}
static int mqtt_msg_process(huayun_var_t *var, int type, ipc_msg_t *mqtt_msg)
{
	int j,ret;
	int null_cnt = 0;
	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;
    }
	
	node_info_list_t *pnode_list = NULL;
	node_info_list_t *tmp = NULL;
	list_for_each_entry_safe(pnode_list, tmp, &var->node_list, list)
	{
		if(strcmp(pnode_list->node_info.sn, tag.sn) == 0)
		{
			cJSON* tags = cJSON_Parse(tag.tag_node);
			if(tags == NULL)
			{
				dy_syslog(LOG_WARNING, "tag_node is NULL");
				break;
			}
			
			if(type == 0)
		    {
		    	cJSON* child = NULL;
			    for (child = tags->child;child;child = child->next)
				{ 
					if(child->string)
			        {
			        	for(j=0; j<pnode_list->node_tag.tag_cnt; j++)
			        	{
			        		if(strcmp(pnode_list->node_tag.tag_info[j].identifier, child->string) == 0)
			        		{
			        			double double_val = -1;
			        			GET_JSON_VALUE_DOUBLE(tags,child->string,double_val);
								pnode_list->node_tag.tag_info[j].val = (float)double_val;
								pnode_list->node_tag.tag_info[j].quality = 0;
								pnode_list->node_tag.tag_info[j].time_stamp = tag.time;
								dy_syslog(LOG_DEBUG,"===sn %s string %s double_val %1.15g val %g time_stamp %d quality %d===",tag.sn,child->string,double_val,pnode_list->node_tag.tag_info[j].val,pnode_list->node_tag.tag_info[j].time_stamp,pnode_list->node_tag.tag_info[j].quality);
								break;
			        		}
			        	}
			        }
			    }
		    }
			else if(type == 1)
		    {
				uint32_t time_stamp = 0;
				int change = 0, threshold = 0;
				node_info_list_t *ptmp = NULL;

				if(strstr(tag.identifier, "event_change"))
					change = 1;
				else if(strstr(tag.identifier, "event_threshold"))
					threshold = 1;
		        else
					goto finish_0;

				uint16_t signalAddr = 0;
				double val = 0, threshold_val = 0;
				
				cJSON* child = NULL;
			    for (child = tags->child;child;child = child->next)
				{
					if(child->string)
			        {
			        	char tmpf[64] = {0};
						char tmpb[64] = {0};
						
			        	sscanf(child->string, "%[^1-9]%d%s", tmpf,&signalAddr,tmpb);
						if(change)
			        	{
			        		if(strstr(child->string, "_raw"))
								GET_JSON_VALUE_DOUBLE(tags,child->string,threshold_val);
							else
								GET_JSON_VALUE_DOUBLE(tags,child->string,val);
			        	}
						else if(threshold)
						{
							if(strstr(child->string, "_threshold"))
								GET_JSON_VALUE_DOUBLE(tags,child->string,threshold_val);
							else
								GET_JSON_VALUE_DOUBLE(tags,child->string,val);
						}
			        }
			    }	
				dy_syslog(LOG_DEBUG,"==++++==sn %s signalAddr %d val %g threshold_val %g ischange %d ==++++==",tag.sn,signalAddr,val,threshold_val,change);
				uint16_t asdu = pnode_list->node_info.port << 8 | pnode_list->node_info.term_addr[0];
				if(change)
				{
					IEC104_BuildSignal_Spon(0, (uint8_t)val, asdu, signalAddr);

					//分位告警信号，变位信号断路器位置取反
					if(signalAddr == 85)
					{
						dy_syslog(LOG_DEBUG,"==++++==sn %s signalAddr %d==++++==",tag.sn,117);
						if(val)
							IEC104_BuildSignal_Spon(0, 0, asdu, 117);
						else
							IEC104_BuildSignal_Spon(0, 1, asdu, 117);
					}
				}
				if(threshold)
				{
					if(val > threshold_val)
						IEC104_BuildSignal_Spon(0, 1, asdu, signalAddr);
					else
						IEC104_BuildSignal_Spon(0, 0, asdu, signalAddr);
					
					//分位告警信号，变位信号断路器位置取反
					if(signalAddr == 85)
					{
						dy_syslog(LOG_DEBUG,"==++++==sn %s signalAddr %d==++++==",tag.sn,117);
						if(val > threshold_val)
							IEC104_BuildSignal_Spon(0, 0, asdu, 117);
						else
							IEC104_BuildSignal_Spon(0, 1, asdu, 117);
					}
				}
		    }
		    else
		        dy_syslog(LOG_WARNING,"Unknown topic %s error!!!",mqtt_msg->topic);
		finish_0:
			cJSON_Delete(tags);
        	break;
		}
	}
	
	return 0;
}

static void msg_mqtt_recv(huayun_var_t *var, int type, ipc_msg_t *mqtt_msg)
{
    int ret,i;
	
	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, "++topic %s data_cnt %d data_flag %d",data->mqtt.topic,var->data_cnt,var->data_flag);
    	list_add_tail(&data->list,&var->data_list);
	}

	var->data_flag = 1;
}

static int mqtt_msg_process_loop(huayun_var_t *var)
{
	data_list_t* data = NULL;
	data_list_t* tmp = NULL;
	
	list_for_each_entry_safe(data, tmp, &var->data_list, list)
	{
		var->data_cnt--;
		mqtt_msg_process(var, data->type, &data->mqtt);
		list_del(&data->list);
		break;
	}

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

static void huayun_loop(huayun_var_t *var)
{
	int ret = -1, maxfd=0;
    fd_set	rset;
    struct timeval timeout;

    while(1)
    {
    	SELECT_INIT();
		SELECT_ADD_FD(var->iec104_fd);
		SELECT_ADD_FD(var->period_timer);
		SELECT_ADD_FD(var->alive_timer);
		
        if(var->data_flag)
		{
			timeout.tv_usec = 200000;
			timeout.tv_sec = 0;
		}
		else
		{
			timeout.tv_usec = 0;
			timeout.tv_sec = 2;
		}
		
        ret = select(maxfd+1, &rset, 0, 0, &timeout);
        if(ret < 0)
        {
            dy_syslog(LOG_WARNING, "%s_%d errno %d", __FUNCTION__, __LINE__, errno);

            if(errno == EINTR)
                continue;
            else
                break;
        }
        else if(ret > 0)
        {
			if(var->iec104_fd > 0 && FD_ISSET(var->iec104_fd, &rset))
            {
                FD_CLR(var->iec104_fd, &rset);

                //接收消息并处理
                iec104_msg(var);
            }
			if(var->period_timer > 0 && FD_ISSET(var->period_timer, &rset))
            {
                FD_CLR(var->period_timer, &rset);
                //接收消息并处理
                iec104_period(var);
            }
			if(var->alive_timer > 0 && FD_ISSET(var->alive_timer, &rset))
            {
                FD_CLR(var->alive_timer, &rset);
                //接收消息并处理
                iec104_alive(var);
            }
        }
		else
		{
			if(var->connect == 0)
			{
				iec104_init(var);
			}
			else if(var->data_flag == 0)
			{
				iec104_pack_select(var);
			}
		}

		mqtt_msg_process_loop(var);
		fflush(stdout);
    }
}

static void huayun_subscribe_all(huayun_var_t *var)
{
    ipc_session_t *ipc_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(ipc_session, topic);

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

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

static int huayun_mqtt_handle_recv_msg(void *obj, ipc_msg_t *mqtt_msg)
{
    huayun_var_t *var = (huayun_var_t *)obj;
    int ret;
    bool matched;
    dy_syslog(LOG_DEBUG, " MQTT client: received MQTT topic:%s payload length:%d",
              mqtt_msg->topic, mqtt_msg->payloadLen);
    ret = mosquitto_topic_matches_sub("ipc/+/+/device/+/data_filtered/property/+", mqtt_msg->topic, &matched);
    if (ret == 0 && matched)
    {
        //实时数据
        msg_mqtt_recv(var, 0, mqtt_msg);
    }
    else
    {
        //事件
        ret = mosquitto_topic_matches_sub("ipc/+/+/device/+/data_filtered/event/+", mqtt_msg->topic, &matched);
        if (ret == 0 && matched)
        {
            msg_mqtt_recv(var, 1, mqtt_msg);
        }
        else
        {
            ret = mosquitto_topic_matches_sub("ipc/+/+/device/+/data_filtered/service/+", mqtt_msg->topic, &matched);
            if (ret == 0 && matched)
            {
                msg_mqtt_recv(var, 2, mqtt_msg);
            }
        }
    }
}

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

    snprintf(clientId, MAX_CLIENT_ID_LEN, "INT_huayun_%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, huayun_mqtt_handle_recv_msg, NULL);
    huayun_subscribe_all(var);
    ipc_session_start(var->session);
}

int huayun_init(huayun_var_t *var)
{
    int i, j;
    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");
    }

	if (load_objects_cfg(&var->object_table, OBJECTS_CFG_PATH) == -1)
    {
        dy_syslog(LOG_ERR, "load object cfg fail");
    }

	huayun_node_parser(var);	
    kv_array_init(&var->identifier_backup, 32);
    huayun_load_config(var);
	
	if(strlen(var->server_ip) <= 1 || var->port == 0)
	{
		strncpy(var->server_ip, IEC104_SERVER_IP, sizeof(var->server_ip));
		var->port = IEC104_PORT;
		var->period = IEC104_REPORT_PERIOD;
	}
	dy_syslog(LOG_INFO,"==server_ip:%s port:%d period:%d==",var->server_ip,var->port,var->period);

	huayun_mqtt_client_init(var);
	dy_db_session_init(&var->db, HUAYUN_DB, (void*)var);
	
    return 0;
}


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

	dy_syslog(LOG_DEBUG, "\n\
            |********************************************|\n\
            |           huayun start  X_X         |\n\
            |********************************************|\n");
    memset(var, 0, sizeof(huayun_var_t));
    huayun_init(var);
    huayun_loop(var);

    return 0;
}
