#include "vint_common.h"
#include "dy_utils/dy_ipc.h"
#include "Encryption_AES.h"

/**
*
* @author: cnscn@163.com
* @reference: lovesnow1314@http://community.csdn.net/Expert/TopicView3.asp?id=5198221 
*
* 用新子串newstr替换源字符串src中的前len个字符内所包含的oldstr子串
*
* @param char* dest 目标串，也就是替换后的新串
* @param const char* src 源字符串，被替换的字符串
* @param const char* oldstr 旧的子串，将被替换的子串
* @param const char* newstr 新的子串
* @param int len 将要被替换的前len个字符
*
* @return char* dest 返回新串的地址
*
*/
static char *strreplace(char *dest, char *src, const char *oldstr, const char *newstr, size_t len)
{
	//如果串相等，则直接返回
	if(strcmp(oldstr, newstr)==0)
		return src; 

	//子串位置指针
	char *needle;
	 

	//临时内存区
	char *tmp; 

	//把源串地址赋给指针dest，即让dest和src都指向src的内存区域
	dest = src; 

	//如果找到子串, 并且子串位置在前len个子串范围内, 则进行替换, 否则直接返回
	while((needle = strstr(dest, oldstr)) && (needle -dest <= len))
	{
		//分配新的空间: +1 是为了添加串尾的'\0'结束符
		tmp=(char*)malloc(strlen(dest)+(strlen(newstr)-strlen(oldstr))+1); 

		//把src内的前needle-dest个内存空间的数据，拷贝到arr
		strncpy(tmp, dest, needle-dest); 

		//标识串结束
		tmp[needle-dest]='\0'; 

		//连接arr和newstr, 即把newstr附在arr尾部, 从而组成新串(或说字符数组)arr
		strcat(tmp, newstr); 

		//把src中 从oldstr子串位置后的部分和arr连接在一起，组成新串arr
		strcat(tmp, needle+strlen(oldstr)); 

		//把用malloc分配的内存，复制给指针retv
		dest = strdup(tmp); 

		//释放malloc分配的内存空间
		free(tmp);
	} 

	return dest;
}

static int send_mqtt_soft_info(vint_var_t *var, char *sn, mqtt_message_t *mqtt_msg)
{
    mqtt_payload_t *mqtt_payload_p = NULL;
    char *data = NULL;
    register_info_t psoft = {0};
    node_message_action_t nm_action = {0};

    mqtt_payload_p = (mqtt_payload_t *)mqtt_msg->payload;
    data = mqtt_payload_p->payload;
    psoft.soft_ver = ntohs(*(uint16_t *)data);
    data = data + 2;
    psoft.protocol_ver = *data;
    nm_action.action = ACTION_SOFT_INFO;
    nm_action.port = VINTF_MQTT;
    strcpy(nm_action.sn, sn);
    memcpy(nm_action.res, &psoft, sizeof(psoft));
    ds_flush_status_action(var->ds, &nm_action);

    return 0;
}

// topic: G/[SN_G]/Rsp_transparent
static int vint_data_report_transparent(vint_var_t *var, mqtt_message_t *mqtt_msg)
{
    char s[4][32] = {0};
    char *topic = mqtt_msg->topic;
    int count = sscanf(topic, "%[^/]/%[^/]/%[^/]/%s", s[0], s[1], s[2], s[3]);
    int data_len = mqtt_msg->payloadLen - sizeof(mqtt_payload_t); // remove MI and UTC
    char *data = NULL;
    pp_real_time_data_t real_data = {0};
    node_cfg_t *p_node = NULL;
    char *sn = NULL;
    char *dtu_sn = NULL;
    int i = 0;
    int ret = 0;
    mqtt_payload_t *mqtt_payload = (mqtt_payload_t *)mqtt_msg->payload;
    uint32_t mi, ts;
    regulate_cmd_backup_t *fields = NULL;

    if (data_len <= 0)
    {
        dy_syslog(LOG_WARNING, "empty data");
        return 0;
    }

    if (count != 4)
    {
        count = sscanf(topic, "%[^/]/%[^/]/%s", s[0], s[1], s[2]);
        if (count == 3)
        {
            dtu_sn = s[1];
        }
        else
        {
            dy_syslog(LOG_ERR, "sscanf error");
            return -1;
        }
    }
    else
    {
        sn = s[2];
        dtu_sn = s[1];
    }

    mi = mqtt_payload->mi;
    ts = mqtt_payload->UTC;

    // store the downlink info
    for (i = 0; i < var->nodes_cfg_table->node_cnt; i++)
    {
        if (strcmp(dtu_sn, var->nodes_cfg_table->node[i].sn) == 0 && var->nodes_cfg_table->node[i].port == VINTF_MQTT)
        {
            p_node = &var->nodes_cfg_table->node[i];
        }
    }
    if (p_node == NULL)
    {
        ret = -1;
        goto out;
    }

    data = malloc(data_len);
    if (data == NULL)
    {
        dy_syslog(LOG_ERR, "malloc fail");
        return -1;
    }

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

        sprintf(key_str, "%s_%d", dtu_sn, mi);

        data_info = var->identifier_backup.get(&var->identifier_backup, key_str);
        if (data_info)
        {
            fields = (regulate_cmd_backup_t *)data_info->data;
        }
    }

    if (fields && fields->src_identifier && strlen(fields->src_identifier))
    {
        // 找对对应的下行控制数据
        strncpy(real_data.src_identifier, fields->src_identifier, sizeof(real_data.src_identifier));
        real_data.mi = fields->mi;
        strcpy(real_data.sn, fields->sn);
    }

    memcpy(data, mqtt_payload->payload, data_len);
    real_data.port = MQTT_DTU_RS485;
    real_data.len = data_len;
    real_data.data = data;
    ///<支持控制消息mi返回
    strcpy(real_data.dtu_sn, dtu_sn);

    flush_device_status_by_sn(var->ds, ACTION_RCV, dtu_sn);
    flush_device_status_by_sn(var->ds, ACTION_RCV, real_data.sn);

    p_node = NULL;
    for (i = 0; i < var->nodes_cfg_table->node_cnt; i++)
    {
        if (strcmp(real_data.sn, var->nodes_cfg_table->node[i].sn) == 0 && var->nodes_cfg_table->node[i].port == MQTT_DTU_RS485)
        {
            p_node = &var->nodes_cfg_table->node[i];
        }
    }
    if (p_node != NULL && strlen(p_node->sub_template_id))
    {
        strncpy(real_data.sub_template_id, p_node->sub_template_id, strlen(p_node->sub_template_id));
    }

    send_to_proto_parser(var->session_int, &real_data);
out:

    return ret;
}

static int reply_common_ack(vint_var_t *var, int err, char *dtu_sn, int mi)
{
    char topic[64];
    char payload[64];
    char *p = NULL;
    time_t ts = time(NULL);

    snprintf(topic, sizeof(topic), "M/%s/ACK", dtu_sn);
    p = payload;
    *(uint16_t *)p = mi;
    p = p + 2;
    *(uint32_t *)p = ts;
    p = p + 4;
    *p = err;
    p = p + 1;
    mqtt_session_publish(var->session_ext, topic, payload, p - payload);

    return 0;
}

static int reply_login_ack(vint_var_t *var, int err, char *dtu_sn, int mi)
{
    char topic[64];
    char payload[64];
    char *p = NULL;
    time_t ts = time(NULL);

    snprintf(topic, sizeof(topic), "M/%s/Rsp_LogIn", dtu_sn);
    p = payload;
    *(uint16_t *)p = mi;
    p = p + 2;
    *(uint32_t *)p = ts;
    p = p + 4;
    *p = err;
    p = p + 1;
    mqtt_session_publish(var->session_ext, topic, payload, p - payload);

    return 0;
}

// 有下行命令的返回数据
// 将根据返回数据中的MI信息查找下行时的cmd info结构，并使用该cmd info结构中的identifier
static int __vint_handle_data_no_transparent(vint_var_t *var, char *sn, uint32_t mi, uint32_t ts,
        mqtt_message_t *mqtt_msg, char *cmd, node_cfg_t *p_node, regulate_cmd_backup_t *fields)
{
    int data_len = mqtt_msg->payloadLen;
    char *data = NULL;
    pp_real_time_data_t real_data = {0};
    int i = 0;
    int ret = 0;
    mqtt_payload_t *mqtt_payload = NULL;

    if (strcmp(cmd, TOPIC_EVT_LOGIN) == 0)
    {
        flush_device_status_by_sn(var->ds, ACTION_LOGIN, sn);
        send_mqtt_soft_info(var, sn, mqtt_msg);
        goto out;
    }

	if(data_len < sizeof(mqtt_payload_t))
		goto out;

    if (strlen(p_node->sub_template_id))
    {
        strncpy(real_data.sub_template_id, p_node->sub_template_id, strlen(p_node->sub_template_id));
    }
    mqtt_payload = (mqtt_payload_t *)mqtt_msg->payload;

	strncpy(real_data.src_identifier, cmd, sizeof(real_data.src_identifier));

    strcpy(real_data.sn, sn);
    if (fields && fields->src_identifier && strlen(fields->src_identifier))
    {
        strncpy(real_data.src_identifier, fields->src_identifier, sizeof(real_data.src_identifier));
        real_data.mi = fields->mi;
    }

    data = malloc(data_len);
    if (data == NULL)
    {
        dy_syslog(LOG_ERR, "malloc fail");
        return -1;
    }
    memcpy(data, mqtt_payload->payload, data_len - sizeof(mqtt_payload_t));
    real_data.port = VINTF_MQTT;
    real_data.len = data_len - sizeof(mqtt_payload_t);
    real_data.data = data;
    real_data.ts = ts;
    dy_syslog_hex(LOG_INFO, real_data.data, real_data.len, "	len %d			网关<-VINT		src_identifier:%s", real_data.len, real_data.src_identifier);
    send_to_proto_parser(var->session_int, &real_data);
out:
	if(data)
    	free(data);
    if (strcmp(cmd, TOPIC_EVT_LOGIN) == 0)
    {
        reply_login_ack(var, 0, sn, mi);
    }
    else if (strcmp(cmd, TOPIC_EVT_LOGIN) == 0 ||
             strcmp(cmd, TOPIC_EVT_PRD_DATA) == 0 ||
             strcmp(cmd, TOPIC_DATA_REPORT) == 0 ||
             strcmp(cmd, TOPIC_DATA_HISTORY) == 0)
    {
        reply_common_ack(var, 0, sn, mi);
    }
    return ret;
}

static int vint_handle_nb_data_no_transparent(vint_var_t *var, mqtt_message_t *mqtt_msg)
{
	cJSON *root = NULL;
    char *client_id = NULL;
    char *value = NULL;
	int i = 0;
	int nb = 0;
	int data_len = 0;
	char *data = NULL;
	char device_id[64] = {0};
	char payload_T[192] = {0};
	node_cfg_t *p_node = NULL;
    pp_real_time_data_t real_data = {0};
    regulate_cmd_backup_t *fields = NULL;
	
	if (sscanf(mqtt_msg->topic, "F/%[^/]/transparent", device_id) != 1)
    {
        dy_syslog(LOG_ERR, "topic:%s can't find device_id!!!", mqtt_msg->topic);
        return -1;
    }
	
    root = cJSON_Parse(mqtt_msg->payload);
    if (root == NULL)
    {
        dy_syslog(LOG_ERR, "payload json parse error!!!");
        return -1;
    }

	value = calloc(1024, 1);
    GET_JSON_VALUE_STRING(root, "value", value);

	//payload中value数据先b64_decode
	data = calloc(MAXBUF, 1);
	data_len = b64_decode(value, data, MAXBUF);	
	dy_syslog(LOG_DEBUG, "===b64_decode data:%s len:%d", data, data_len);

	//AES解密,得到协议报文数据
	data_len = AES_Decrypt(data, data_len);
	dy_syslog(LOG_DEBUG, "===AES_Decrypt value:%s len %d", value,data_len);

	//解析协议报文中的identifier
	cJSON *jval = cJSON_Parse(data);
	if (jval == NULL)
	{
		dy_syslog(LOG_ERR, "payload data json parse error!!!");
		goto out;
	}

	GET_JSON_VALUE_STRING(jval, "T", payload_T);

	if(strlen(payload_T))
	{
		char t1[3][64] = {0};
		if (sscanf(payload_T, "%[^/]/%[^/]/%[^/]", t1[0],t1[1],t1[2]) != 3)
		{
			dy_syslog(LOG_ERR, "payload_T:%s format error!!!", payload_T);
			goto out;
		}
		if (t1[2][0] == 'B' && t1[2][1] == '_')
		{
			char *dest = NULL;
			dest = strreplace(dest, t1[2], "B_", "G_", 200);
			strncpy(real_data.src_identifier, dest, sizeof(real_data.src_identifier));
			free(dest);
		}
		else
			strncpy(real_data.src_identifier, t1[2], sizeof(real_data.src_identifier));
	}
	
	/*GET_JSON_VALUE_STRING(jval, "X", value);
	dy_syslog(LOG_DEBUG, "===value:%s data_len %d", value,data_len);
	data_len = b64_decode(value, data, MAXBUF);
	dy_syslog(LOG_DEBUG, "===data:%s data_len %d", data,data_len);*/

	//根据device_id找到SN
    for (i = 0; i < var->nodes_cfg_table->node_cnt; i++)
    {
        if (strcmp(device_id, var->nodes_cfg_table->node[i].device_id) == 0)
        {
            p_node = &var->nodes_cfg_table->node[i];
			strcpy(real_data.sn, p_node->sn);
			break;
        }
    }

	if(p_node == NULL)
	{
		dy_syslog(LOG_WARNING, "can't find SN by device_id:%s", device_id);
		goto out;
	}
	
    //restore the identifier
    {
        char key_str[64] = {0};
        data_info_t *data_info = NULL;

		sprintf(key_str, "%s", real_data.sn);

        data_info = var->identifier_backup.get(&var->identifier_backup, key_str);
        if (data_info)
        {
            fields = (regulate_cmd_backup_t *)data_info->data;
        }
    }

    if (fields && fields->src_identifier && strlen(fields->src_identifier))
    {
        // 找对对应的下行控制数据
        strncpy(real_data.src_identifier, fields->src_identifier, sizeof(real_data.src_identifier));
        real_data.mi = fields->mi;
        strcpy(real_data.sn, fields->sn);
    }

    real_data.port = VINTF_NB;
    real_data.len = data_len;
    real_data.data = data;

    flush_device_status_by_sn(var->ds, ACTION_RCV, real_data.sn);
    send_to_proto_parser(var->session_int, &real_data);
out:

	free(data);
	free(value);
	if(jval)
		cJSON_Delete(jval);
	cJSON_Delete(root);
    return 0;
}


static int vint_handle_data_no_transparent(vint_var_t *var, mqtt_message_t *mqtt_msg)
{
    char s[3][32] = {0};
    char *sn;
    uint32_t mi, ts;
    int i;
    mqtt_payload_t *mqtt_payload = NULL;
    node_cfg_t *p_node = NULL;
    int ret = 0;
    char *topic = mqtt_msg->topic;
    regulate_cmd_backup_t *fields = NULL;

    if (sscanf(topic, "%[^/]/%[^/]/%s", s[0], s[1], s[2]) != 3)
    {
        dy_syslog(LOG_ERR, "sscanf error, topic:%s", topic);
        ret = -1;
        goto out;
    }
    sn = s[1];
    mqtt_payload = (mqtt_payload_t *)mqtt_msg->payload;
    mi = mqtt_payload->mi;
    ts = mqtt_payload->UTC;
    dy_syslog(LOG_INFO, "receive topic:%s sn:%s mi:%d ts:%u", topic, sn, mi, ts);

    dy_syslog_hex(LOG_INFO, mqtt_msg->payload, mqtt_msg->payloadLen, "	len %d			网关<-VINT		", mqtt_msg->payloadLen);

    flush_device_status_by_sn(var->ds, ACTION_RCV, sn);
    for (i = 0; i < var->nodes_cfg_table->node_cnt; i++)
    {
        if (strcmp(sn, var->nodes_cfg_table->node[i].sn) == 0 && var->nodes_cfg_table->node[i].port == VINTF_MQTT)
        {
            p_node = &var->nodes_cfg_table->node[i];
        }
    }
    if (p_node == NULL)
    {
        ret = -1;
        dy_syslog(LOG_ERR, "SN:%s not belong to us", sn);
        goto out;
    }

    {
        char key_str[64] = {0};
        data_info_t *data_info = NULL;

        sprintf(key_str, "%s_%d", sn, mi);

        data_info = var->identifier_backup.get(&var->identifier_backup, key_str);
        if (data_info)
        {
            fields = (regulate_cmd_backup_t *)data_info->data;
        }
    }

    __vint_handle_data_no_transparent(var, sn, mi, ts, mqtt_msg, s[2], p_node, fields);

out:

    return ret;
}

// mqtt send  M/[SN_G]/transparent with payload
static int vint_send_mqtt_transparent_data(vint_var_t *var, pp_regulate_signal_t *regulate_data)
{
    mqtt_payload_t *mqtt_payload_p = NULL;
    uint16_t payloadLen = regulate_data->len;
    char *payload = NULL; //MI and UTC
    char *sn = regulate_data->sn;
    char *dtu_sn = regulate_data->dtu_sn;
    int i;
    char topic[128];
    int ret = 0;

    payload = malloc(payloadLen + sizeof(mqtt_payload_t));
    if (payload == NULL)
    {
        dy_syslog(LOG_ERR, "malloc failed");
        ret = -1;
        goto out;
    }
    snprintf(topic, sizeof(topic), "M/%s/transparent", dtu_sn);
    mqtt_payload_p = (mqtt_payload_t *)payload;
    mqtt_payload_p->mi = (unsigned short)regulate_data->mi;
    mqtt_payload_p->UTC = time(NULL);
    dy_syslog_hex(LOG_DEBUG, regulate_data->data, payloadLen, "payload data");
    memcpy(mqtt_payload_p->payload, regulate_data->data, payloadLen);
    dy_syslog_hex(LOG_DEBUG, payload, payloadLen + sizeof(mqtt_payload_t), "	len %d			VINT->SN:%s		topic:%s", payloadLen + sizeof(mqtt_payload_t), sn, topic);

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

        strcpy(fields.src_identifier, regulate_data->src_identifier);
        strcpy(fields.sn, sn);
        fields.mi = regulate_data->mi;

        sprintf(key_str, "%s_%d", dtu_sn, mqtt_payload_p->mi);
        var->identifier_backup.insert(&var->identifier_backup, key_str, &fields, sizeof(fields));
    }

    mqtt_session_publish(var->session_ext, topic, payload, payloadLen + sizeof(mqtt_payload_t));
    free(payload);
out:
    return ret;
}

// mqtt send  M/[SN_G]/cmd with payload
static int vint_send_mqtt_data_normal(vint_var_t *var, pp_regulate_signal_t *regulate_data)
{
    mqtt_payload_t *mqtt_payload_p = NULL;
    uint16_t payloadLen = regulate_data->len;
    char *payload = NULL;
    char *sn = regulate_data->sn;
    int i;
    char topic[128];
    int ret = 0;
    time_t ts = time(NULL);

    payload = malloc(payloadLen + sizeof(mqtt_payload_t));
    if (payload == NULL)
    {
        dy_syslog(LOG_ERR, "malloc failed");
        ret = -1;
        goto out;
    }
    snprintf(topic, sizeof(topic), "M/%s/%s", sn, regulate_data->src_identifier);
    mqtt_payload_p = (mqtt_payload_t *)payload;
    mqtt_payload_p->mi = (unsigned short)regulate_data->mi;
    mqtt_payload_p->UTC = ts;
    dy_syslog_hex(LOG_DEBUG, regulate_data->data, payloadLen, "payload data");
    memcpy(mqtt_payload_p->payload, regulate_data->data, payloadLen);
    dy_syslog_hex(LOG_DEBUG, payload, payloadLen + sizeof(mqtt_payload_t), "	len %d			VINT->SN:%s		topic:%s", payloadLen + sizeof(mqtt_payload_t), sn, topic);

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

        strcpy(fields.src_identifier, regulate_data->src_identifier);
        strcpy(fields.sn, sn);
        fields.mi = regulate_data->mi;

        sprintf(key_str, "%s_%d", sn, mqtt_payload_p->mi);
        var->identifier_backup.insert(&var->identifier_backup, key_str, &fields, sizeof(fields));
    }
    mqtt_session_publish(var->session_ext, topic, payload, payloadLen + sizeof(mqtt_payload_t));
    free(payload);
out:
    return ret;
}

int vint_send_rglt_data(vint_var_t *var, pp_regulate_signal_t *regulate_data)
{
	int i;
	node_cfg_t *p_node = NULL;
	
    if (regulate_data == NULL)
        return -1;

    flush_device_status_by_sn(var->ds, ACTION_SND, regulate_data->sn);
    for (i = 0; i < var->nodes_cfg_table->node_cnt; i++)
    {
        if (strcmp(regulate_data->sn, var->nodes_cfg_table->node[i].sn) == 0 && var->nodes_cfg_table->node[i].port == VINTF_MQTT)
        {
            p_node = &var->nodes_cfg_table->node[i];
            break;
        }
    }
    if (p_node != NULL && p_node->template->protocol == GW_PROTC_WOOLINK_DTU)
    {
        vint_send_mqtt_data_lnxall_dtu(var, regulate_data);
        goto out;
    }
	if (regulate_data->port == VINTF_NB)
    {
        // 发送到前置机
        vint_send_mqtt_data_to_front(var, regulate_data);
        goto out;
    }

    if (regulate_data->port == VINTF_TCP)
    {
        // 发送到前置机
        vint_send_tcp_data_to_front(var, regulate_data);
        goto out;
    }

    if (regulate_data->protocol == GW_PROTC_DYWL_RTU_V66)
    {
        // 新RTU协议
        vint_send_data_to_rtu(var, regulate_data);
        goto out;
    }

    if (regulate_data->port == MQTT_DTU_RS485 && regulate_data->dtu_sn != 0)
    {
        // the connect port is MQTT_DTU_RS485 and the dtu_sn0
        // this dev is a dev under DTU, should use transparent mode to send the data
        vint_send_mqtt_transparent_data(var, regulate_data);
    }
    else
    {
        vint_send_mqtt_data_normal(var, regulate_data);
    }

out:
    return 0;
}

int vint_send_mqtt_data(vint_var_t *var, mqtt_message_t *mqtt_msg)
{
    pp_regulate_signal_t regulate_data = {0};
    int ret = 0;
    time_t ts = time(NULL);

    ret = get_rglt_data_from_json(&regulate_data, mqtt_msg->payload);
    if (ret < 0)
    {
        dy_syslog(LOG_ERR, "parse real data structure failed");
        return -1;
    }
    if (regulate_data.period > 0)
    {
        period_msg_t once;

        dy_syslog(LOG_INFO, "Recevied interval ctrl cmd periad %d", regulate_data.period);
        snprintf(once.key, 128, "%s_%s", regulate_data.sn, regulate_data.src_identifier);
        once.interval = regulate_data.period;
        once.last_poll = 0;
        once.rs = calloc(1, sizeof(pp_regulate_signal_t));
        memcpy(once.rs, &regulate_data, sizeof(pp_regulate_signal_t));
        once.started = 1;
        once.rs->data = calloc(1, regulate_data.len);
        memcpy(once.rs->data, regulate_data.data, regulate_data.len);

        period_msg_add(var->period, &once);
    }
    else
    {
        vint_send_rglt_data(var, &regulate_data);
    }

out:
    free(regulate_data.data);
    return 0;
}

int vint_external_mqtt_handle_recv_msg(void *obj, mqtt_message_t *mqtt_msg)
{
    vint_var_t *var = (vint_var_t *)obj;
    char *topic = mqtt_msg->topic;
    char sn[SN_MAX_LEN] = {0};
    node_cfg_t *p_node = NULL;

    dy_syslog(LOG_DEBUG, " MQTT client: received MQTT topic:%s payload length:%d",
              mqtt_msg->topic, mqtt_msg->payloadLen);

    if (mqtt_msg->topic == NULL)
    {
        dy_syslog(LOG_ERR, "topic is NULL");
        goto out;
    }

    if (mqtt_msg->topic[0] == 'F' && topic[1] == '/')
    {
        // 前置机的topic
        int nb = vint_handle_mqtt_from_front_msg(var, mqtt_msg);
		if(nb == 1)
		{
			vint_handle_nb_data_no_transparent(var, mqtt_msg);
		}
        goto out;
    }
    else if (strncmp(mqtt_msg->topic, "RTU/", 4) == 0)
    {
        handle_new_rtu_data(var, mqtt_msg);
        goto out;
    }

    sscanf(mqtt_msg->topic, "%*[^/]/%[^/]/%*s", sn);
    if (strlen(sn) > 0)
    {
        int i = 0;
        for (i = 0; i < var->nodes_cfg_table->node_cnt; i++)
        {
            if (strcmp(sn, var->nodes_cfg_table->node[i].sn) == 0 && var->nodes_cfg_table->node[i].port == VINTF_MQTT)
            {
                p_node = &var->nodes_cfg_table->node[i];
                break;
            }
        }
        if (p_node != NULL && p_node->template->protocol == GW_PROTC_WOOLINK_DTU)
        {
            handle_lnxall_dtu(var, mqtt_msg, p_node);
            goto out;
        }
        else
        {
            if (strstr(mqtt_msg->topic, "transparent"))
            {
                //透传数据
                //DTU 透传模式时采用该处理方式网关直接与该DTU下面的设备通信
                if (topic[0] == 'G' && topic[1] == '/')
                {
                    //消息源为DTU G/[SN_G]/[SN_D]/Rsp_transparent
                    vint_data_report_transparent(var, mqtt_msg);
                }
            }
            else
            {
                //非透传数据
                //DTU非透传或RTU数据
                if (topic[0] == 'G' && topic[1] == '/')
                {
                    // 消息源为DTU G/[SN_G]/cmd
                    vint_handle_data_no_transparent(var, mqtt_msg);
                }
            }
        }
    }

out:
    return 0;
}
