#include <stdio.h>
#include <stdlib.h>

#include "ipc_session.h"
#include "dy_utils/dy_common.h"
#include "common.h"
#include "dy_utils/dy_ipc.h"

void internal_mqtt_subscribe_all(shanghai_S5_var_t *var)
{
    ipc_session_t *ipc_session = var->session_int;
    char topic[TOPIC_MAX_LEN] = {0};
    int i;

    snprintf(topic, TOPIC_MAX_LEN, "ipc/+/%s/#", APP_KEY);
    ipc_session_subscribe(ipc_session, topic);
}

int send_data_to_front(shanghai_S5_var_t *var, node_cfg_t *node, char *buf, int len)
{
    char *data_b64 = NULL;
    char topic[TOPIC_MAX_LEN] = {0};
    cJSON *mqtt_json = cJSON_CreateObject();
    char *data_tmp = NULL;

    sprintf(topic, "GW/%s/transparent", node->sn);

    cJSON_AddStringToObject(mqtt_json, "sn", node->sn);
    data_b64 = (char *)malloc(1024);
    b64_encode(buf, len, data_b64, 1024);
    cJSON_AddStringToObject(mqtt_json, "payload", data_b64);
    cJSON_AddNumberToObject(mqtt_json, "len", strlen(data_b64));

    data_tmp = cJSON_Print(mqtt_json);

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

    free(data_b64);
    free(data_tmp);
    cJSON_Delete(mqtt_json);

    flush_device_status_by_sn(var->ds, ACTION_SND, node->sn);
}

int shanghai_s5_send_rglt(shanghai_S5_var_t *var, pp_regulate_signal_t *regulate_data)
{
    char *data_b64 = NULL;
    char topic[TOPIC_MAX_LEN] = {0};
    cJSON *mqtt_json = NULL;
    int len;
    char *data_tmp = NULL;
    int mi;
    int period;
    int buff_len = regulate_data->len * 2 + 1024;

    data_b64 = (char *)malloc(buff_len);
    if (data_b64 == NULL)
    {
        dy_syslog(LOG_ERR, "malloc fail");
        return -1;
    }

    //store the identifier
    {
        char key_str[64] = {0};
        char *identifier = NULL;
        regulate_cmd_backup_t fields;
        time_t now = time(NULL);

        strcpy(fields.src_identifier, regulate_data->src_identifier);
        strcpy(fields.sn, regulate_data->sn);
        fields.mi = regulate_data->mi;
        fields.ts = now;
        sprintf(key_str, "%s", regulate_data->sn);
        free(identifier);
        var->identifier_backup.insert(&var->identifier_backup, key_str, &fields, sizeof(fields));
    }

    //将数据加入list,由poll函数去发送
    mqtt_json = cJSON_CreateObject();
    sprintf(topic, "GW/%s/transparent", regulate_data->sn);
    b64_encode(regulate_data->data, regulate_data->len, data_b64, buff_len);
    cJSON_AddStringToObject(mqtt_json, "sn", regulate_data->sn);
    cJSON_AddStringToObject(mqtt_json, "payload", data_b64);
    cJSON_AddNumberToObject(mqtt_json, "len", strlen(data_b64));

    data_tmp = cJSON_Print(mqtt_json);

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

    free(data_b64);
    free(data_tmp);
    cJSON_Delete(mqtt_json);

    flush_device_status_by_sn(var->ds, ACTION_SND, regulate_data->sn);
}

static int forward_data_to_front(shanghai_S5_var_t *var, ipc_msg_t *mqtt_msg)
{
    pp_regulate_signal_t regulate_data;
    int ret = 0;

    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
    {
        ret = shanghai_s5_send_rglt(var, &regulate_data);
        if (ret < 0)
        {
            dy_syslog(LOG_ERR, "parse real data structure failed");
            return -1;
        }
    }

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

static int internal_mqtt_handle_recv_msg(void *obj, ipc_msg_t *mqtt_msg)
{
    shanghai_S5_var_t *var = (shanghai_S5_var_t *)obj;
    //dy_syslog(LOG_DEBUG, " MQTT client:%s received MQTT topic:%s payload length:%d",
    //          session->client.clientId, mqtt_msg->topic, mqtt_msg->payloadLen);
    if (strstr(mqtt_msg->topic, TOPIC_SEND_RGLT_SIGNAL_RAW_DATA))
    {
        //发送下行数据
        forward_data_to_front(var, mqtt_msg);
    }
}

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

    snprintf(clientId, MAX_CLIENT_ID_LEN, "SHS5_%s", var->sn_str);
    var->session_int = ipc_session_new(clientId, (void*)var, IPC_DEFAULT);
    if (var->session_int == NULL) return -1;

    ipc_session_set_callbacks(var->session_int, internal_mqtt_handle_recv_msg, NULL);
    internal_mqtt_subscribe_all(var);
    ipc_session_start(var->session_int);
}
