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

int fanuc_period_rs_triger(fanuc_thread_t *fanuc_thread, node_cfg_t *node);


gw_port_e ports[] = {TCP_CLIENT};

static int fanuc_client_init(fanuc_thread_t *fanuc_thread, tcp_node_list_t *tcp_node, char *ip_addr, int port)
{	
	int allocated = 0;
  	int ret = 0;
	char cnc_id[40];
	uint32_t cnc_ids[4];

	if(cnc_startupprocess(0, "focas.log") != EW_OK) {
	    dy_syslog(LOG_ERR, "Failed to create required log file!");
	    return -1;
  	}

	ret = cnc_allclibhndl3(ip_addr, port, 10, &tcp_node->flibHndl);
	if(ret != EW_OK)
	{
	    dy_syslog(LOG_ERR, "Failed to connect to cnc! (%s %d %d)", ip_addr, port, ret);
	    goto cleanup;
  	}

	ret = cnc_rdcncid(tcp_node->flibHndl, (unsigned long *)cnc_ids);
	if(ret != EW_OK)
	{
		dy_syslog(LOG_ERR, "Failed to read cnc id!");
		ret = 1;
		goto cleanup;
	}

	snprintf(cnc_id, 40, "%08x-%08x-%08x-%08x", cnc_ids[0], cnc_ids[1],cnc_ids[2], cnc_ids[3]);
	dy_syslog(LOG_DEBUG, "machine id: %s flibHndl %d", cnc_id, tcp_node->flibHndl);
	//return 0;
cleanup:
	if(cnc_freelibhndl(tcp_node->flibHndl) != EW_OK)
		dy_syslog(LOG_ERR, "Failed to free library handle!");
	tcp_node->flibHndl = 0;
	return 0;
}

static int fanuc_msg_tag_data_ctrl(fanuc_thread_t *fanuc_thread, ipc_msg_t *mqtt_msg)
{
    char find = 0;
    pp_regulate_signal_t regulate_data;
    fanuc_data_list_t *fanuc_data;
    int ret;
    int port;

    //先检查参数
    ret = get_rglt_data_from_json(&regulate_data, mqtt_msg->payload);
    if (ret < 0 || regulate_data.len <= 0)
    {
        dy_syslog(LOG_ERR, "parse real data structure failed");
        return -1;
    }

    fanuc_data = calloc(sizeof(fanuc_data_list_t), 1);
    if (fanuc_data)
    {
        memcpy(&fanuc_data->rs, &regulate_data, sizeof(regulate_data));

        {
            tcp_node_list_t *tcp_node = NULL;
            list_for_each_entry(tcp_node, &fanuc_thread->tcp_node_list, list)
            {
                fanuc_thread->period_triggered = 1;
                fanuc_data_list_t *cmd = NULL;
                dy_syslog(LOG_DEBUG, "tcp_port(%d %d) tcp_ip_addr(%s %s) \n", tcp_node->tcp_port, fanuc_data->rs.tcp_port, tcp_node->tcp_ip_addr, fanuc_data->rs.tcp_ip_addr);
                if (tcp_node->tcp_port == fanuc_data->rs.tcp_port && strcmp(tcp_node->tcp_ip_addr, fanuc_data->rs.tcp_ip_addr) == 0)
                {
                    find = 0;
                    list_for_each_entry(cmd, &tcp_node->data_list, list)
                    {
                        if (strcmp(cmd->rs.sn, regulate_data.sn) == 0 &&
                                cmd->rs.period == fanuc_data->rs.period &&
                                cmd->rs.len == fanuc_data->rs.len &&
                                (memcmp(cmd->rs.data, fanuc_data->rs.data, cmd->rs.len) == 0))
                        {
                            find = 1;
                            break;
                        }
                    }
                    if (!find)
                    {
                        list_add_tail(&fanuc_data->list, &tcp_node->data_list);
                        dy_syslog(LOG_DEBUG, "list_add_tail tcp_port(%d) tcp_ip_addr(%s) sn:%s\n", fanuc_data->rs.tcp_port, fanuc_data->rs.tcp_ip_addr, fanuc_data->rs.sn);
                    }
                }
            }
        }

        dy_syslog(LOG_INFO, "sn %s data[0] %d identifier %s period %d", fanuc_data->rs.sn, fanuc_data->rs.data[0], regulate_data.src_identifier, regulate_data.period);
        if (find)
        {
            dy_syslog(LOG_WARNING, "Repeated commands !!!");
            free(fanuc_data);
        }
    }

    return 0;
}

static int fanuc_tag_data_report(fanuc_thread_t *fanuc_thread, char *data, int len, unsigned char port)
{
    pp_real_time_data_t real_data = {0};
    regulate_cmd_backup_t *fields = NULL;
	time_t now = time(NULL);

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

        sprintf(key_str, "fanuc");
        data_info = fanuc_thread->identifier_backup.get(&fanuc_thread->identifier_backup, key_str);
        if (data_info)
        {
            fields = (regulate_cmd_backup_t *)data_info->data;
        }
    }
    real_data.port = port;
    real_data.len = len;
    real_data.instruction_code = data[1];
    real_data.data = malloc(real_data.len);
	if (fields != NULL && ((now - fields->ts) < 10))
    {
        real_data.mi = fields->mi;
        strcpy(real_data.sn, fields->sn);
        strncpy(real_data.src_identifier, fields->src_identifier, sizeof(real_data.src_identifier));
        strcpy(real_data.requester, fields->requester);
        real_data.down_data = malloc(fields->len);
        memcpy(real_data.down_data, fields->data, fields->len);
        real_data.down_len = fields->len;
    }

    memcpy(real_data.data, data, real_data.len);

    send_to_proto_parser(fanuc_thread->session, &real_data);
    flush_device_status_by_sn(fanuc_thread->ds, ACTION_RCV, real_data.sn);

    free(real_data.data);
	if (real_data.down_data)
	{
		free(real_data.down_data);
	}

    return 0;
}

void fanuc_rdposition(unsigned short flibHndl, cJSON *tag_data, short type)
{
	char key[32] = {0};
    ODBPOS pos[MAX_AXIS];
    short num = MAX_AXIS;
	
    short ret = cnc_rdposition(flibHndl, type, &num, pos);
    dy_syslog(LOG_DEBUG, "cnc_rdposition ret %d", ret);
    if(!ret) {
        int i;
        for(i = 0 ; i < num ; i++) {
            dy_syslog(LOG_DEBUG, "position %c = %d\n", pos[i].abs.name, pos[i].abs.data);
			switch(type)
			{
				case 0:
					sprintf(key, "absolute_%c", pos[i].abs.name);
					break;
				case 1:
					sprintf(key, "machine_%c", pos[i].abs.name);
					break;
				case 2:
					sprintf(key, "relative_%c", pos[i].abs.name);
					break;
				case 3:
					sprintf(key, "distance_%c", pos[i].abs.name);
					break;
				default:
					break;
			}
			cJSON_AddNumberToObject(tag_data, key, pos[i].abs.data);
        }
    }
}

/* 获取机台进给 */
int fanuc_cnc_actf(unsigned short flibHndl, cJSON *tag_data)
{
	ODBACT odbact;
	short res = cnc_actf(flibHndl, &odbact);
	dy_syslog(LOG_DEBUG, "获取机台进给 cnc_actf res %d\n", res);
	if (res == EW_OK)
	{
		dy_syslog(LOG_DEBUG, "获取机台进给 odbact.data %d\n", odbact.data);
		cJSON_AddNumberToObject(tag_data, "axis_feedrate", odbact.data);
	}
}

/* 获取机台转速 */
int fanuc_cnc_acts(unsigned short flibHndl, cJSON *tag_data)
{
  ODBACT spindle;
  short res = cnc_acts(flibHndl, &spindle);
  dy_syslog(LOG_DEBUG, "cnc_acts res %d\n", res);
  if (res == EW_OK) {
	dy_syslog(LOG_DEBUG, "获取机台转速 spindle.data %d\n", spindle.data);
	cJSON_AddNumberToObject(tag_data, "spindle_speed", spindle.data);
  }
}

/* 获取机台负载 */
int fanuc_cnc_rdspload(unsigned short flibHndl, cJSON *tag_data)
{
	ODBSPN odbspn[8];
	short type = -1;
	short res = cnc_rdspload(flibHndl, type, odbspn);
	dy_syslog(LOG_DEBUG, "获取机台负载 cnc_rdspload res:%d MAX_SPINDLE %d\n", res, MAX_SPINDLE);
	if (res == EW_OK) {
		//dy_syslog(LOG_DEBUG, "获取机台负载 odbspn.datano %d data %d\n", odbspn.datano, odbspn.data[0]);
		//cJSON_AddNumberToObject(tag_data, "load_information", odbspn.datano);
	}
}

/* 获取CNC 状态信息 */
int fanuc_cnc_statinfo(unsigned short flibHndl, cJSON *tag_data)
{
  ODBST odbst;
  short res = cnc_statinfo(flibHndl, &odbst);
  dy_syslog(LOG_DEBUG, "cnc_statinfo res:%d\n", res);
  if (res == EW_OK)
  {
	dy_syslog(LOG_DEBUG, "获取CNC 状态信息 hdck %d tmmode %d aut:%d run %d,edit %d,motion %d,mstb %d,emergency %d,alarm %d \n", odbst.hdck,odbst.tmmode,odbst.aut,odbst.run,odbst.motion,odbst.mstb,odbst.emergency,odbst.alarm,odbst.edit);
	cJSON_AddNumberToObject(tag_data, "status_hdck", odbst.hdck);
	cJSON_AddNumberToObject(tag_data, "status_tmmode", odbst.tmmode);
	cJSON_AddNumberToObject(tag_data, "status_aut", odbst.aut);
	cJSON_AddNumberToObject(tag_data, "status_run", odbst.run);
	cJSON_AddNumberToObject(tag_data, "status_motion", odbst.motion);
	cJSON_AddNumberToObject(tag_data, "status_mstb", odbst.mstb);
	cJSON_AddNumberToObject(tag_data, "status_emergency", odbst.emergency);
	cJSON_AddNumberToObject(tag_data, "status_alarm", odbst.alarm);
	cJSON_AddNumberToObject(tag_data, "status_edit", odbst.edit);
  }
}

/* 获取加工工件数量 */
int getProcessTimes(unsigned short flibHndl, cJSON *tag_data) {
	ODBPARANUM odbparanum;

	short res = cnc_rdparanum(flibHndl, &odbparanum);
	dy_syslog(LOG_DEBUG, "cnc_rdparanum res:%d\n", res);

	IODBPSD iodbpsd2s;

	short s_number[8];
	s_number[0] = 6712;

	short e_number[8];
	e_number[0] = 6712;

	short length[8];
	length[0] = 16;

	res = cnc_rdparar(flibHndl, s_number, (short) -1, e_number, length, &iodbpsd2s);
	dy_syslog(LOG_DEBUG, "cnc_rdparar res:%d\n", res);
}

/* 获取加工工件数量 */
int fanuc_cnc_sysinfo(unsigned short flibHndl, cJSON *tag_data)
{
	ODBSYS sysinfo;
	char var[8];
	short ret = cnc_sysinfo(flibHndl, &sysinfo);
	dy_syslog(LOG_DEBUG, "cnc_sysinfo ret %d \n", ret);
	if (ret == EW_OK)
	{
		dy_syslog(LOG_DEBUG, "  Max Axis: %d\n", sysinfo.max_axis);
		dy_syslog(LOG_DEBUG, "  CNC Type: %c%c\n", sysinfo.cnc_type[0], sysinfo.cnc_type[1]);
		dy_syslog(LOG_DEBUG, "  MT Type: %c%c\n", sysinfo.mt_type[0], sysinfo.mt_type[1]);
		dy_syslog(LOG_DEBUG, "  Series: %c%c\n", sysinfo.series[0], sysinfo.series[1], sysinfo.series[2], sysinfo.series[3]);
		dy_syslog(LOG_DEBUG, "  Version: %c%c\n", sysinfo.version[0], sysinfo.version[1], sysinfo.version[2], sysinfo.version[3]);
		dy_syslog(LOG_DEBUG, "  Axes: %c%c\n", sysinfo.axes[0], sysinfo.axes[1]);		

		cJSON_AddNumberToObject(tag_data, "max_axis", sysinfo.max_axis);
		sprintf(var, "%c%c", sysinfo.cnc_type[0], sysinfo.cnc_type[1]);
		cJSON_AddStringToObject(tag_data, "cnc_type", var);
		sprintf(var, "%c%c", sysinfo.mt_type[0], sysinfo.mt_type[1]);
		cJSON_AddStringToObject(tag_data, "mt_type", var);
		sprintf(var, "%c%c%c%c", sysinfo.series[0], sysinfo.series[1], sysinfo.series[2], sysinfo.series[3]);
		cJSON_AddStringToObject(tag_data, "series", var);
		sprintf(var, "%c%c%c%c", sysinfo.version[0], sysinfo.version[1], sysinfo.version[2], sysinfo.version[3]);
		cJSON_AddStringToObject(tag_data, "version", var);
		sprintf(var, "%c%c", sysinfo.axes[0], sysinfo.axes[1]);
		cJSON_AddStringToObject(tag_data, "axis_number", var);
	}
}

/* 运行总时间 */
int fanuc_total_machining_time(unsigned short flibHndl, cJSON *tag_data)
{
	IODBPSD iodbpsd1, iodbpsd2;
	int32_t totalTime = 0;
	short ret = cnc_rdparam(flibHndl, 6751, 0, sizeof(IODBPSD), &iodbpsd1);
	dy_syslog(LOG_DEBUG, "运行总时间 cnc_rdparam ret %d \n", ret);
	if (EW_OK == ret)
	{
		totalTime = iodbpsd1.u.ldata / 1000;//换算成秒
		ret = cnc_rdparam(flibHndl, 6752, 0, sizeof(IODBPSD), &iodbpsd2);
		if (EW_OK == ret)
		{
		  totalTime += iodbpsd2.u.ldata * 60;//换算成秒             
		  dy_syslog(LOG_DEBUG, "pSet->totalMachiningTime is %d s\n", totalTime);
		  cJSON_AddNumberToObject(tag_data, "total_machining_time", totalTime);
		}
		else
		{
		  dy_syslog(LOG_DEBUG, "cnc_rdparam 6752 error %d\n", ret);
		}
	}
}

/* 切削总时间 */
int fanuc_total_cutting_time(unsigned short flibHndl, cJSON *tag_data)
{
	IODBPSD iodbpsd1, iodbpsd2;
	int32_t totalTime = 0;
	short ret = cnc_rdparam(flibHndl, 6753, 0, sizeof(IODBPSD), &iodbpsd1);
	dy_syslog(LOG_DEBUG, "切削总时间cnc_rdparam ret %d", ret);
	if (EW_OK == ret)
	{
		totalTime = iodbpsd1.u.ldata / 1000;//换算成秒
		ret = cnc_rdparam(flibHndl, 6754, 0, sizeof(IODBPSD), &iodbpsd2);
		if (EW_OK == ret)
		{
		  totalTime += iodbpsd2.u.ldata * 60;//换算成秒
		  dy_syslog(LOG_DEBUG, "pSet->totalCuttingTime is %d s", totalTime);
		  cJSON_AddNumberToObject(tag_data, "total_cutting_time", totalTime);
		}
		else
		{
		  dy_syslog(LOG_DEBUG, "cnc_rdparam 6754 error %d", ret);
		}
	}
}

/* 循环时间 */
int fanuc_circle_time(unsigned short flibHndl, cJSON *tag_data)
{
	IODBPSD iodbpsd1, iodbpsd2;
	int32_t totalTime = 0;
	short ret = cnc_rdparam(flibHndl, 6757, 0, sizeof(IODBPSD), &iodbpsd1);
	dy_syslog(LOG_DEBUG, "切削总时间cnc_rdparam ret %d", ret);
	if (EW_OK == ret)
	{
		totalTime = iodbpsd1.u.ldata / 1000;//换算成秒
		ret = cnc_rdparam(flibHndl, 6758, 0, sizeof(IODBPSD), &iodbpsd2);
		if (EW_OK == ret)
		{
		  totalTime += iodbpsd2.u.ldata * 60;//换算成秒
		  dy_syslog(LOG_DEBUG, "circle_time is %d s", totalTime);
		  cJSON_AddNumberToObject(tag_data, "circle_time", totalTime);
		}
		else
		{
		  dy_syslog(LOG_DEBUG, "cnc_rdparam 6754 error %d", ret);
		}
	}
}

//加工件数
int fanuc_parts(unsigned short flibHndl, cJSON *tag_data)
{
	IODBPSD iodbpsd1;
	
	short ret = cnc_rdparam(flibHndl, 6711, 0, sizeof(IODBPSD), &iodbpsd1);
	dy_syslog(LOG_DEBUG, "加工件数cnc_rdparam ret %d", ret);
	if (EW_OK == ret)
	{
		dy_syslog(LOG_DEBUG, "parts:%d ", iodbpsd1.u.ldata);
		cJSON_AddNumberToObject(tag_data, "parts", iodbpsd1.u.ldata);
	}
}

//加工件总数
int fanuc_total_parts(unsigned short flibHndl, cJSON *tag_data)
{
	IODBPSD iodbpsd1;
	
	short ret = cnc_rdparam(flibHndl, 6712, 0, sizeof(IODBPSD), &iodbpsd1);
	dy_syslog(LOG_DEBUG, "加工件总数cnc_rdparam ret %d", ret);
	if (EW_OK == ret)
	{
		dy_syslog(LOG_DEBUG, "total_parts:%d ", iodbpsd1.u.ldata);
		cJSON_AddNumberToObject(tag_data, "total_parts", iodbpsd1.u.ldata);
	}
}

//开机总时间
int fanuc_power_on_time(unsigned short flibHndl, cJSON *tag_data)
{
	IODBPSD iodbpsd1;
	
	short ret = cnc_rdparam(flibHndl, 6750, 0, sizeof(IODBPSD), &iodbpsd1);
	dy_syslog(LOG_DEBUG, "加工件总数cnc_rdparam ret %d", ret);
	if (EW_OK == ret)
	{
		dy_syslog(LOG_DEBUG, "power_on_time:%d ", iodbpsd1.u.ldata);
		cJSON_AddNumberToObject(tag_data, "power_on_time", iodbpsd1.u.ldata);
	}
}

//主轴、伺服轴的负载
int fanuc_load_meter(unsigned short flibHndl, cJSON *tag_data)
{
	ODBSVLOAD sv;
	ODBSPLOAD sp;

	short a = 1;//伺服轴的数量
	short ret = cnc_rdsvmeter(flibHndl, &a, &sv);
	dy_syslog(LOG_DEBUG, "cnc_rdsvmeter ret %d a %d\n", ret,a);
	if (EW_OK == ret)
	{
		cJSON_AddNumberToObject(tag_data, "sv_load", sv.svload.data);
		dy_syslog(LOG_DEBUG, "伺服轴1负载 svload.data %d", sv.svload.data);
	}
	ret = cnc_rdspmeter(flibHndl, 1, &a, &sp);
	dy_syslog(LOG_DEBUG, "cnc_rdspmeter ret %d a %d\n", ret,a);
	if (EW_OK == ret)
	{
		cJSON_AddNumberToObject(tag_data, "spindle_load", sp.spload.data);
		cJSON_AddNumberToObject(tag_data, "spindle_speed", sp.spspeed.data);
		dy_syslog(LOG_DEBUG, "主轴的负载 spload %d spspeed %d", sp.spload.data, sp.spspeed.data);
	}
}


static int fanuc_get_data(tcp_node_list_t *tcp_node, cJSON *tag_data)
{	
	unsigned short flibHndl;
	int allocated = 0;
  	int ret = 0;
	char cnc_id[40];
	uint32_t cnc_ids[4];

	if(cnc_startupprocess(0, "focas.log") != EW_OK) {
	    dy_syslog(LOG_ERR, "Failed to create required log file!");
	    return -1;
  	}

	ret = cnc_allclibhndl3(tcp_node->tcp_ip_addr, tcp_node->tcp_port, 10, &flibHndl);
	if(ret != EW_OK)
	{
	    dy_syslog(LOG_ERR, "Failed to connect to cnc! (%s %d %d)", tcp_node->tcp_ip_addr, tcp_node->tcp_port, ret);
	    goto cleanup;
  	}

	ret = cnc_rdcncid(flibHndl, (unsigned long *)cnc_ids);
	if(ret != EW_OK)
	{
		dy_syslog(LOG_ERR, "Failed to read cnc id!");
		ret = 1;
		goto cleanup;
	}

	snprintf(cnc_id, 40, "%08x-%08x-%08x-%08x", cnc_ids[0], cnc_ids[1],cnc_ids[2], cnc_ids[3]);
	cJSON_AddStringToObject(tag_data, "cnc_id", cnc_id);
	dy_syslog(LOG_DEBUG, "machine id: %s flibHndl %d", cnc_id, flibHndl);

	fanuc_rdposition(flibHndl, tag_data, 0);
	fanuc_rdposition(flibHndl, tag_data, 1);
	fanuc_rdposition(flibHndl, tag_data, 2);
	fanuc_rdposition(flibHndl, tag_data, 3);

	fanuc_cnc_actf(flibHndl, tag_data);
	fanuc_cnc_acts(flibHndl, tag_data);

	fanuc_cnc_rdspload(flibHndl, tag_data);
	fanuc_cnc_statinfo(flibHndl, tag_data);
	getProcessTimes(flibHndl, tag_data);
	fanuc_cnc_sysinfo(flibHndl, tag_data);
	fanuc_total_machining_time(flibHndl, tag_data);
	fanuc_total_cutting_time(flibHndl, tag_data);
	fanuc_circle_time(flibHndl, tag_data);
	fanuc_parts(flibHndl, tag_data);
	fanuc_total_parts(flibHndl, tag_data);
	fanuc_power_on_time(flibHndl, tag_data);
	fanuc_load_meter(flibHndl, tag_data);
	
cleanup:
	if(cnc_freelibhndl(flibHndl) != EW_OK)
		dy_syslog(LOG_ERR, "Failed to free library handle!");
	flibHndl = 0;
	return 0;
}


static void fanuc_bus_send(fanuc_thread_t *fanuc_thread, pp_regulate_signal_t *rs, unsigned char port, unsigned int communication_timeout, tcp_node_list_t *tcp_node)
{
    int ret;
    char *data = rs->data;
    int len = rs->len;
    long long cur_ms = clock_get_ms();
    int try_cnt = 0;
    
    flush_device_status_by_sn(fanuc_thread->ds, ACTION_SND, rs->sn);
fanuc_flag:
    //if (port == MODBUS_TCP)
    {
        //tcp发送数据,阻塞等待执行数据返回
        //if (!flibHndl)
        //fanuc_client_init(fanuc_thread, tcp_node, tcp_node->tcp_ip_addr, tcp_node->tcp_port);

        //if (tcp_node->flibHndl)
        {
        	pp_real_time_data_t real_data = {0};
        	char *rsp_str = NULL;
        	cJSON *tag_data = cJSON_CreateObject();

			dy_syslog_hex(LOG_DEBUG, data, len, "sn:%s  len %d  identifier:%s    网关->PORT[%d] send_cnt[0x%x]", rs->sn, len, rs->src_identifier, port,tcp_node->send_cnt++);
			dy_syslog(LOG_DEBUG, "src_identifier %s strcmp %d rs->data %s", rs->src_identifier,strcmp(rs->src_identifier, "rdposition"),rs->data);
        	if(strcmp(rs->src_identifier, "readCncData") == 0)
        	{
        		cJSON *data = cJSON_Parse(rs->data);
				if(data)
				{
					short type;
					GET_JSON_VALUE_INT(data, "type", type);
					dy_syslog(LOG_DEBUG, "type %d", type);
					fanuc_get_data(tcp_node, tag_data);
				}
        	}

			rsp_str = cJSON_PrintUnformatted(tag_data);
			dy_syslog(LOG_DEBUG, "rsp_str %s", rsp_str);
			cJSON_Delete(tag_data);


			if (rsp_str != NULL)
		    {
		        real_data.mi = rs->mi;
		        strcpy(real_data.sn, rs->sn);
		        strncpy(real_data.src_identifier, rs->src_identifier, sizeof(real_data.src_identifier));

		        real_data.port = rs->port;
		        real_data.len = strlen(rsp_str);
		        real_data.data = rsp_str;

		        send_to_proto_parser(fanuc_thread->var->session, &real_data);
		    }

		    free(rsp_str);
        }
        //else
            //dy_syslog(LOG_WARNING, "tcp_ip_addr:%s port:%d flibHndl is 0!!!", tcp_node->tcp_ip_addr, tcp_node->tcp_port);
    }
}

static int ds_flush_timer(fanuc_thread_t *fanuc_thread)
{
    time_t now = time(NULL);

    TIMER_CONFIRM(fanuc_thread->node_status_timer);
    if (fanuc_thread->period_triggered == 0)
    {
        //fanuc_period_rs_triger(fanuc_thread, fanuc_thread->node);
    }

    if (now >= fanuc_thread->nodes_status_sync_date + 30)
    {
        static time_t lasttime = 0;
        fanuc_thread->nodes_status_sync_date = now;
        char *node_status = ds_print(fanuc_thread->ds);

        dy_syslog(LOG_INFO, "online_status_changed:%d strlen(node_status):%d, (now-lasttime):%d",
                  fanuc_thread->ds->online_status_changed, strlen(node_status), now - lasttime);
        if (node_status && strlen(node_status) > 20 &&
                (fanuc_thread->ds->online_status_changed || now > lasttime + 3600)) //在线状态变化上报，1小时也会上报
        {
            char topic[TOPIC_MAX_LEN];
            int ret = 0;

            snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/%s/%s", fanuc_thread->var->sn_str, fanuc_thread->var->proc_name, TOPIC_EVT_NODESSTATUS);

            ret = ipc_session_publish(fanuc_thread->session, topic, node_status, strlen(node_status));
            if (ret == 0)
            {
                lasttime = now;
                fanuc_thread->ds->online_status_changed = false;
            }
            node_sta_backup(fanuc_thread->ds, FANUC_NODE_STATUS_BAK_FILE);
        }
        if (node_status)
        {
            write_file_data(NODES_CACHE "/fanuc_status.json", node_status, strlen(node_status));
            free(node_status);
        }
    }
}

static int send_poll(fanuc_thread_t *fanuc_thread, struct list_head *data_list, int port, tcp_node_list_t *tcp_node)
{
    int ret, find = 0;
    unsigned int communication_timeout = 200;
    fanuc_data_list_t *send_data = NULL;
    long long now = get_time_stamp_ms();

    list_for_each_entry(send_data, data_list, list)
    {
        if (!send_data->rs.period)
        {
            find = 1;
            if (communication_timeout < send_data->rs.communication_timeout)
            {
                communication_timeout = send_data->rs.communication_timeout;
            }
            fanuc_bus_send(fanuc_thread, &send_data->rs, port, communication_timeout, tcp_node);

            list_del(&send_data->list);
            free(send_data->rs.data);
            free(send_data);
            break;
        }
        else
        {
            //第一次发送延时
            if (send_data->last_send < 5)
            {
                send_data->last_send++;
                find = 1;
                break;
            }
            if (fanuc_thread->period_flag == 0) continue;
            ret = check_timeout_msecond(now, &send_data->last_send, send_data->rs.period);
            if (ret)
            {
                find = 1;
                if (communication_timeout < send_data->rs.communication_timeout)
                {
                    communication_timeout = send_data->rs.communication_timeout;
                }
                fanuc_bus_send(fanuc_thread, &send_data->rs, port, communication_timeout, tcp_node);
                break;
            }
        }
    }

    return find;
}

static int fanuc_send_poll(fanuc_thread_t *fanuc_thread, int time_trigger)
{
	if(time_trigger)
    	TIMER_CONFIRM(fanuc_thread->send_timer);

    if (!fanuc_thread->start_flag)
    {
        return 0;
    }

    tcp_node_list_t *tcp_node = NULL;
    list_for_each_entry(tcp_node, &fanuc_thread->tcp_node_list, list)
    {
        if (send_poll(fanuc_thread, &tcp_node->data_list, MODBUS_TCP, tcp_node) == 1)
            return 1;
    }

    return 0;
}

static void *fanuc_thread_loop(void *param)
{
	fanuc_thread_t *fanuc_thread = (fanuc_thread_t *)param;
    int ret = -1, maxfd, i, need_send = 0;
    fd_set rset;
    struct timeval timeout;

    while (1)
    {
        SELECT_INIT();
        SELECT_ADD_FD(fanuc_thread->send_timer);
        SELECT_ADD_FD(fanuc_thread->node_status_timer);

		if (need_send)
        {
            timeout.tv_usec = 5000;
            timeout.tv_sec = 0;
        }
        else
        {
            timeout.tv_usec = 0;
            timeout.tv_sec = 5;
        }

        ret = select(maxfd + 1, &rset, 0, 0, &timeout);
        if (ret < 0)
        {
            dy_syslog(LOG_INFO, "errno %d\n", errno);

            if (errno == EINTR)
            {
                continue;
            }
            else
            {
                break;
            }
        }
		else if (ret == 0)
        {
            need_send = fanuc_send_poll(fanuc_thread, 0);
        }
        else if (ret > 0)
        {
            if (fanuc_thread->send_timer > 0 && FD_ISSET(fanuc_thread->send_timer, &rset))
            {
                FD_CLR(fanuc_thread->send_timer, &rset);
                need_send = fanuc_send_poll(fanuc_thread, 1);
            }
            if (fanuc_thread->node_status_timer > 0 && FD_ISSET(fanuc_thread->node_status_timer, &rset))
            {
                FD_CLR(fanuc_thread->node_status_timer, &rset);
                ds_flush_timer(fanuc_thread);
            }
        }
    }
}

static void fanuc_subscribe_all(fanuc_thread_t *fanuc_thread, node_cfg_t *node)
{
    ipc_session_t *ipc_session = fanuc_thread->session;
    char topic[TOPIC_MAX_LEN] = {0};
    int i = 0;
	
    snprintf(topic, TOPIC_MAX_LEN, "ipc/+/%s/device/%s/data/%s", FANUC_APP_KEY, node->sn, TOPIC_SEND_RGLT_SIGNAL_RAW_DATA);
    ipc_session_subscribe(ipc_session, topic);
    snprintf(topic, TOPIC_MAX_LEN, "ipc/+/%s/%s", FANUC_APP_KEY, TOPIC_NOTIFY_UPGRADE);
    ipc_session_subscribe(ipc_session, topic);
}

static int fanuc_mqtt_handle_recv_msg(void *obj, ipc_msg_t *mqtt_msg)
{
    fanuc_thread_t *fanuc_thread = (fanuc_thread_t *)obj;
    dy_syslog(LOG_DEBUG, "received MQTT topic:%s payload length:%d", mqtt_msg->topic, mqtt_msg->payloadLen);
    if (strstr(mqtt_msg->topic, TOPIC_SEND_RGLT_SIGNAL_RAW_DATA))
    {
        //发送下行数据
        fanuc_msg_tag_data_ctrl(fanuc_thread, mqtt_msg);
    }
}

// 建立与内部broker之间的MQTT连接
static int fanuc_mqtt_client_init(fanuc_thread_t *fanuc_thread, node_cfg_t *node)
{
    char clientId[MAX_CLIENT_ID_LEN] = {0};

    snprintf(clientId, MAX_CLIENT_ID_LEN, "INT_fanuc_%s_%s", fanuc_thread->var->sn_str, node->sn);
    fanuc_thread->session = ipc_session_new(clientId, (void*)fanuc_thread, IPC_DEFAULT);
    if (fanuc_thread->session == NULL) return -1;

    ipc_session_set_callbacks(fanuc_thread->session, fanuc_mqtt_handle_recv_msg, NULL);

    fanuc_subscribe_all(fanuc_thread, node);
    ipc_session_start(fanuc_thread->session);
}

int fanuc_period_rs_triger(fanuc_thread_t *fanuc_thread, node_cfg_t *node)
{
    int i, j, k;
	int find = 0;
	fanuc_var_t *var = fanuc_thread->var;

    //建立TCP连接
    tcp_node_list_t *tcp_node = NULL;
    list_for_each_entry(tcp_node, &fanuc_thread->tcp_node_list, list)
    {
        if (strcmp(tcp_node->tcp_ip_addr, node->tcp_ip_addr) == 0 &&
                tcp_node->tcp_port == node->tcp_port)
        {
            find = 1;
            break;
        }
    }

    if (find == 0)
    {
        tcp_node = calloc(sizeof(tcp_node_list_t), 1);

		INIT_LIST_HEAD(&tcp_node->data_list);
        fanuc_client_init(fanuc_thread, tcp_node, node->tcp_ip_addr, node->tcp_port);
        strncpy(tcp_node->tcp_ip_addr, node->tcp_ip_addr, SN_MAX_LEN);
        tcp_node->tcp_port = node->tcp_port;
        list_add_tail(&tcp_node->list, &fanuc_thread->tcp_node_list);
    }

    for (i = 0; i < var->template_table->template_cnt; i++)
    {
    	if(strcmp(node->template_id, var->template_table->template[i].template_id))
			continue;
		dy_syslog(LOG_DEBUG, "template_cnt====sn:%s template_cnt:%d serviceCnt:%d", node->sn,var->template_table->template_cnt,var->template_table->template[i].service_tab->serviceCnt);
        //解析周期命令
        for (j = 0; j < var->template_table->template[i].service_tab->serviceCnt; j++)
        {
            if (var->template_table->template[i].service_tab->service[j].server_period && !var->template_table->template[i].service_tab->service[j].direction)
            {
                //节点信息
                if (strcmp(node->app_key, FANUC_APP_KEY) == 0 && strcmp(node->template_id, var->template_table->template[i].template_id) == 0)
                {
                    char *data = service_rs_data(&var->template_table->template[i].service_tab->service[j], node->sn);
                    char topic[TOPIC_MAX_LEN] = {0};
                    snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/%s/device/%s/data/%s", var->sn_str, port_enum2char(node->port), node->sn, TOPIC_EVT_SET_RGLT);
                    ipc_session_publish(fanuc_thread->session, topic, data, strlen(data));
                    free(data);
                }
            }
        }
		break;
    }
	dy_syslog(LOG_DEBUG, "====sn:%s end...", node->sn);
}

int fanuc_start_thread(fanuc_var_t *var, node_cfg_t *node)
{
	int j;
	pthread_t thread_fanuc;
    fanuc_thread_t *fanuc_thread = NULL;
	char *buf = NULL;

	dy_syslog(LOG_DEBUG, "=========sn:%s=======", node->sn);

    fanuc_thread = calloc(1, sizeof(fanuc_thread_t));
    if (fanuc_thread == NULL)
    {
        dy_syslog(LOG_ERR, "Malloc failed");
        return -1;
    }

	fanuc_thread->var = var;
	fanuc_thread->node = node;

	dev_status_init(&fanuc_thread->ds, fanuc_thread->nodes_cfg_table, fanuc_thread->template_table, ports, ARRAY_SIZE(ports));
    node_sta_recovery(fanuc_thread->ds, FANUC_NODE_STATUS_BAK_FILE);

    kv_array_init(&fanuc_thread->identifier_backup, 32);

    INIT_LIST_HEAD(&fanuc_thread->tcp_node_list);

    fanuc_thread->send_timer = my_timer_create();
    if (fanuc_thread->send_timer > 0)
    {
    	dy_syslog(LOG_INFO, "send_timer %d ms", 200);
        my_timer_set(fanuc_thread->send_timer, 1, 200);
    }
    fanuc_thread->node_status_timer = my_timer_create();
    if (fanuc_thread->node_status_timer > 0)
    {
        my_timer_set(fanuc_thread->node_status_timer, 1, 5000);
    }

    fanuc_mqtt_client_init(fanuc_thread, node);
    fanuc_period_rs_triger(fanuc_thread, node);

    fanuc_thread->start_flag = 1;
    fanuc_thread->period_flag = 1;

    pthread_create(&thread_fanuc, NULL, fanuc_thread_loop, (void *)fanuc_thread);
}

