/**
* @file         mqtt_for_front.c
* @brief        handle the MQTT data from/to front.
* @details      handle the MQTT data from/to front.
* @author       zhaoxianpeng
* @date     2019-10-28
* @version  A001
* @par Copyright (c):
*       DYIOT
* @par History:
*   version: zhaoxianpeng, 2019-10-28, initial version\n
*/
#include "vint_common.h"
#include "Encryption_AES.h"
/**
 * 发送MQTT数据到指定的SN
 * @param[in]   var 模块变量部分.
 * @param[in]   sn  目的设备SN.
 * @param[in]   payload  要发送的数据，将会经过base64_encode之后加入json然后通过MQTT发送.
 * @param[in]   payload_len  要发送的数据长度.
 * @retval    成功
 * @par 标识符
 *      保留
 * @par 其它
 *      无
 * @par 修改日志
 *      zhaoxianpeng于2019-10-28创建
 */
static int __vint_send_tcp_data_to_front(vint_var_t *var, char *sn, char *payload, int payload_len)
{
    char *encode_str = NULL;
    char topic[32] = {0};
    cJSON *mqtt_json = cJSON_CreateObject();
    int len;
    char *data_tmp = NULL;

    //将数据加入list,由poll函数去发送
    sprintf(topic, "GW/%s/transparent", sn);

    // 将payload进行base64 encode
    encode_str = (char *)malloc(MAXBUF);
    len = b64_encode(payload, payload_len, encode_str, MAXBUF);

    cJSON_AddStringToObject(mqtt_json, "sn", sn);
    cJSON_AddStringToObject(mqtt_json, "payload", encode_str);
    cJSON_AddNumberToObject(mqtt_json, "len", len);

    data_tmp = cJSON_Print(mqtt_json);

    mqtt_session_publish(var->session_ext, topic, data_tmp, strlen(data_tmp));

    free(encode_str);
    free(data_tmp);
    cJSON_Delete(mqtt_json);

    return 0;
}

/**
 * 发送MQTT数据到指定的SN
 * @param[in]   var 模块变量部分.
 * @param[in]   sn  目的设备SN.
 * @param[in]   payload  要发送的数据，将会经过base64_encode之后加入json然后通过MQTT发送.
 * @param[in]   payload_len  要发送的数据长度.
 * @retval    成功
 * @par 标识符
 *      保留
 * @par 其它
 *      无
 * @par 修改日志
 *      longwenjian于2020-10-28创建
 */
static int __vint_send_mqtt_data_to_front(vint_var_t *var, char *sn, char *payload, int payload_len)
{
	int i;
	char *encode_str = NULL;
    char topic[256] = {0};
    int len = payload_len;
	char *data = NULL;
    char *data_tmp = NULL;
	cJSON *mqtt_json = cJSON_CreateObject();

    //payload pkcs5padding
    encode_str = (char *)malloc(MAXBUF);
	data = calloc(payload_len+AES_BLOCK_SIZE*2, 1);
	memcpy(data, payload, payload_len);
	
	if(payload_len % AES_BLOCK_SIZE)
	{
		uint8_t padding  = AES_BLOCK_SIZE - len % AES_BLOCK_SIZE;
		memset(data + len, padding, padding);
		len +=padding;
	}
	else 
	{
		memset(data + len, AES_BLOCK_SIZE, AES_BLOCK_SIZE);
		len += AES_BLOCK_SIZE;
	}

	dy_syslog(LOG_DEBUG, "payload_len:%d len:%d payload:%s", payload_len, len, payload);
	dy_syslog(LOG_DEBUG, "data:%s", data);

	//payload AES_Encrypt
	AES_Encrypt(data, len, data);
	//payload b64_encode
    len = b64_encode(data, len, encode_str, MAXBUF);
	//组包，并发送
	for(i=0; i<var->nodes_cfg_table->node_cnt; i++)
	{
		if(strcmp(var->nodes_cfg_table->node[i].sn, sn) == 0)
		{
			//将数据加入list,由poll函数去发送
			sprintf(topic, "GW/%s/transparent", var->nodes_cfg_table->node[i].device_id);
			cJSON_AddStringToObject(mqtt_json, "appId", var->nodes_cfg_table->node[i].huawei_app_id);
			cJSON_AddStringToObject(mqtt_json, "deviceId", var->nodes_cfg_table->node[i].device_id);
			cJSON_AddStringToObject(mqtt_json, "value", encode_str);
			
			data_tmp = cJSON_Print(mqtt_json);
			dy_syslog(LOG_DEBUG, "====topic:%s data_tmp:%s", topic, data_tmp);
		    mqtt_session_publish(var->session_ext, topic, data_tmp, strlen(data_tmp));

		    free(data_tmp);
			break;
		}
	}

	free(data);
	free(encode_str);
    cJSON_Delete(mqtt_json);

    return 0;
}

/**
 * 发送MQTT数据到指定的SN
 * @param[in]   var 模块变量部分.
 * @param[in]   regulate_data  发送的参数信息，会记录在每个设备的cmd_info里面.
 * @retval    成功
 * @par 标识符
 *      保留
 * @par 其它
 *      无
 * @par 修改日志
 *      longwenjian于2020-06-29创建
 */
int vint_send_mqtt_data_to_front(vint_var_t *var, pp_regulate_signal_t *regulate_data)
{
    int ret = 0;
    char *sn = regulate_data->sn;
    int i;

    //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", sn);
        var->identifier_backup.insert(&var->identifier_backup, key_str, &fields, sizeof(fields));
    }
    __vint_send_mqtt_data_to_front(var, regulate_data->sn, regulate_data->data, regulate_data->len);

out:
    return ret;
}

/**
 * 发送MQTT数据到指定的SN
 * @param[in]   var 模块变量部分.
 * @param[in]   regulate_data  发送的参数信息，会记录在每个设备的cmd_info里面.
 * @retval    成功
 * @par 标识符
 *      保留
 * @par 其它
 *      无
 * @par 修改日志
 *      zhaoxianpeng于2019-10-28创建
 */
int vint_send_tcp_data_to_front(vint_var_t *var, pp_regulate_signal_t *regulate_data)
{
    int ret = 0;
    char *sn = regulate_data->sn;
    int i;

    //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, regulate_data->mi);
        var->identifier_backup.insert(&var->identifier_backup, key_str, &fields, sizeof(fields));
    }
    __vint_send_tcp_data_to_front(var, regulate_data->sn, regulate_data->data, regulate_data->len);

out:
    return ret;
}

/**
 * 接收来自前置机的MQTT message
 * 并按照接收的数据内容进行处理，当需要返回数据时，将待返回的数据加入到list里面，由poll函数来发送
 * 当收到需要通过SN的请求时，立即发送SN数据
 * @param[in]   var 模块变量部分.
 * @param[in]   mqtt_msg  收到的MQTT消息内容，包含topic和payload.
 * @retval    成功
 * @par 标识符
 *      保留
 * @par 其它
 *      无
 * @par 修改日志
 *      zhaoxianpeng于2019-10-28创建
 */
int vint_handle_mqtt_from_front_msg(vint_var_t *var, mqtt_message_t *mqtt_msg)
{
    cJSON *root = NULL;
    uint8_t *buf = NULL;
    int port, socket;
    char *ip = NULL;
    char *client_id = NULL;
    char *payload = NULL;
	char *type = NULL;
    int len = 0;
    int ack_len = 0;
    char topic[64] = {0};
    char *encode_str = NULL;
    char *data_tmp = NULL;
	int nb = 0;

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

    dy_syslog(LOG_DEBUG, "receive topic:%s payload:%s", mqtt_msg->topic, mqtt_msg->payload);

    root = cJSON_Parse(mqtt_msg->payload);
    if (root == NULL)
    {
        dy_syslog(LOG_ERR, "json parse error");
        return -1;
    }
    GET_JSON_VALUE_DY_STRING(root, "payload", payload);
	GET_JSON_VALUE_DY_STRING(root, "type", type);

	if(type && strcmp(type, "nb") == 0)
	{
		nb = 1;
		goto out;
	}

    buf = (char *)malloc(MAXBUF);
    len = b64_decode(payload, buf, MAXBUF);
    if (buf[0] == 0x60 && buf[len - 1] == 0x27)
    {
        payload_data_hdr_t hdr = {0};
        int ret;
        ret = handle_dyiot_tcp_data(var, buf, len, &hdr);
        if (strcmp(mqtt_msg->topic, "F/SearchSN") == 0)
        {
            if (ret != -1 && ret != NOT_BELONG_TO_US && strlen(hdr.sn) > 0)
            {
                //consider sync SN topic as control frame, handle it immediately
                cJSON *rsp_sn = cJSON_CreateObject();

                GET_JSON_VALUE_INT(root, "port", port);
                GET_JSON_VALUE_INT(root, "socket", socket);
                GET_JSON_VALUE_DY_STRING(root, "ip", ip);
                GET_JSON_VALUE_DY_STRING(root, "clientId", client_id);

                sprintf(topic, "GW/%s/RspSearchSN", client_id);
                cJSON_AddStringToObject(rsp_sn, "sn", hdr.sn);
                cJSON_AddNumberToObject(rsp_sn, "socket", socket);
                cJSON_AddStringToObject(rsp_sn, "ip", ip);
                cJSON_AddNumberToObject(rsp_sn, "port", port);
                // 以字符串的方式发送数据
                data_tmp = cJSON_Print(rsp_sn);

                mqtt_session_publish(var->session_ext, topic, data_tmp, strlen(data_tmp));

                cJSON_Delete(rsp_sn);

                free(ip);
                free(client_id);
                free(data_tmp);
            }
        }

        if (ret != NORETURN && ret != -1)
        {
            vint_tcp_gen_common_ack(var, &hdr, buf, ret);
            __vint_send_tcp_data_to_front(var, hdr.sn, buf, strlen(buf));
        }
    }
	else if (len > 4 && buf[0] == 0xEB && buf[1] == 0x90)
	{
		payload_data_hdr_t hdr = {0};
        int ret;
		dy_syslog(LOG_DEBUG, "topic:%s len %d 0x%x 0x%x", mqtt_msg->topic, len, buf[0], buf[1]);
		
		if(len > 61)
		{
			ret = handle_unmanned_ship_tcp_data(var, buf, 61, &hdr);
        	ret += handle_unmanned_ship_tcp_data(var, buf+61, len-61, &hdr);
		}
		else
			ret = handle_unmanned_ship_tcp_data(var, buf, len, &hdr);
        if (strcmp(mqtt_msg->topic, "F/SearchSN") == 0)
        {
            if (ret != -1 && ret != NOT_BELONG_TO_US && strlen(hdr.sn) > 0)
            {
                //consider sync SN topic as control frame, handle it immediately
                cJSON *rsp_sn = cJSON_CreateObject();

                GET_JSON_VALUE_INT(root, "port", port);
                GET_JSON_VALUE_INT(root, "socket", socket);
                GET_JSON_VALUE_DY_STRING(root, "ip", ip);
                GET_JSON_VALUE_DY_STRING(root, "clientId", client_id);

                sprintf(topic, "GW/%s/RspSearchSN", client_id);
                cJSON_AddStringToObject(rsp_sn, "sn", hdr.sn);
                cJSON_AddNumberToObject(rsp_sn, "socket", socket);
                cJSON_AddStringToObject(rsp_sn, "ip", ip);
                cJSON_AddNumberToObject(rsp_sn, "port", port);
                // 以字符串的方式发送数据
                data_tmp = cJSON_Print(rsp_sn);
				dy_syslog(LOG_INFO, "unmanned_ship topic:%s payload:%s", topic, data_tmp);
                mqtt_session_publish(var->session_ext, topic, data_tmp, strlen(data_tmp));

                cJSON_Delete(rsp_sn);

                free(ip);
                free(client_id);
                free(data_tmp);
            }
        }
	}
	else if (len > 32 && buf[0] == 0x23 && buf[1] == 0x23 && buf[len - 7] == 0x26 && buf[len - 8] == 0x26)
	{
		dy_syslog(LOG_DEBUG, "topic:%s len %d 0x%x 0x%x puyu", mqtt_msg->topic, len, buf[0], buf[1]);
		payload_data_hdr_t hdr = {0};
		int ret = handle_puyu_tcp_data(var, buf, len, &hdr);
		if (strcmp(mqtt_msg->topic, "F/SearchSN") == 0)
        {
            if (ret != -1 && ret != NOT_BELONG_TO_US && strlen(hdr.sn) > 0)
            {
                //consider sync SN topic as control frame, handle it immediately
                cJSON *rsp_sn = cJSON_CreateObject();

                GET_JSON_VALUE_INT(root, "port", port);
                GET_JSON_VALUE_INT(root, "socket", socket);
                GET_JSON_VALUE_DY_STRING(root, "ip", ip);
                GET_JSON_VALUE_DY_STRING(root, "clientId", client_id);

                sprintf(topic, "GW/%s/RspSearchSN", client_id);
                cJSON_AddStringToObject(rsp_sn, "sn", hdr.sn);
                cJSON_AddNumberToObject(rsp_sn, "socket", socket);
                cJSON_AddStringToObject(rsp_sn, "ip", ip);
                cJSON_AddNumberToObject(rsp_sn, "port", port);
                // 以字符串的方式发送数据
                data_tmp = cJSON_Print(rsp_sn);
				dy_syslog(LOG_INFO, "puyu topic:%s payload:%s", topic, data_tmp);
                mqtt_session_publish(var->session_ext, topic, data_tmp, strlen(data_tmp));

                cJSON_Delete(rsp_sn);

                free(ip);
                free(client_id);
                free(data_tmp);
            }
        }
	}
	else if (len > 35 && buf[0] == 0x54 && buf[len-1] == 0x3B && strncmp(buf, "TimeStamp", strlen("TimeStamp")) == 0)
	{
		dy_syslog(LOG_DEBUG, "topic:%s CSI-DT len %d %s", mqtt_msg->topic, len, buf);
		payload_data_hdr_t hdr = {0};
		int ret = handle_csi_dt_tcp_data(var, buf, len, &hdr);
		if (strcmp(mqtt_msg->topic, "F/SearchSN") == 0)
        {
            if (ret != -1 && ret != NOT_BELONG_TO_US && strlen(hdr.sn) > 0)
            {
                //consider sync SN topic as control frame, handle it immediately
                cJSON *rsp_sn = cJSON_CreateObject();

                GET_JSON_VALUE_INT(root, "port", port);
                GET_JSON_VALUE_INT(root, "socket", socket);
                GET_JSON_VALUE_DY_STRING(root, "ip", ip);
                GET_JSON_VALUE_DY_STRING(root, "clientId", client_id);

                sprintf(topic, "GW/%s/RspSearchSN", client_id);
                cJSON_AddStringToObject(rsp_sn, "sn", hdr.sn);
                cJSON_AddNumberToObject(rsp_sn, "socket", socket);
                cJSON_AddStringToObject(rsp_sn, "ip", ip);
                cJSON_AddNumberToObject(rsp_sn, "port", port);
                // 以字符串的方式发送数据
                data_tmp = cJSON_Print(rsp_sn);
				dy_syslog(LOG_INFO, "CSI-DT topic:%s payload:%s", topic, data_tmp);
                mqtt_session_publish(var->session_ext, topic, data_tmp, strlen(data_tmp));

                cJSON_Delete(rsp_sn);

                free(ip);
                free(client_id);
                free(data_tmp);
            }
        }
	}
out:
	if(type)
		free(type);
    free(payload);
    free(buf);
    cJSON_Delete(root);

    return nb;
}
