#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <time.h>
#include <sys/select.h>
#include "cJSON.h"
#include "my_public.h"
#include "multi_meters.h"
#include "mqtt_session.h"
#include "lnxall_ubuslog.h"

#define POWER_IDENTIFIER "sanxianggongluxinxi"
#define POWER_TAG "TAG_272F"
#define MULTI_METERS_CFG_FILE  "/app/config/ems/multi_meters.json"
#define EMS_STRUCTURE_CFG_FILE "/app/config/ems/ems_structure.json" 
#define DEFAULT_MI 12323
#define TAG_LEN 10

typedef struct mqtt_recv_msg_s {
    CHAR identifier[IDENTIFIER_LEN];
    CHAR devSn[SN_MAX_LEN];
    WORD timestamp;
    WORD mi;
} mqtt_recv_msg_t;

typedef struct service_s {
    CHAR identifier[IDENTIFIER_LEN];
    CHAR inidentifier[IDENTIFIER_LEN];
    CHAR reportMeterType[DEV_TYPE_MAX_LEN];
    CHAR tag[TAG_LEN];
} service_t;

static int handle_mqtt_conn_state_change(void *obj, int state)
{
    switch (state)
    {
    case MQTT_CONNECTED:
        dbg_syslog(LOG_NOTICE, "mqtt connected");
        break;
    case MQTT_DISCONNECTED:
        dbg_syslog(LOG_ERR, "mqtt disconnected!");
        break;
    default:
        dbg_syslog(LOG_ERR, "mqtt unknown state!\n");
    }

    return 0;
}

static void store_power_recv(meter_t *pMeter, int meterCnt, const char *devSn, double powerRecv)
{
    if (pMeter == NULL || devSn == NULL) {
        dbg_syslog(LOG_ERR, "invalid argument! pMeter: %p, devSn: %p", pMeter, devSn);
        return;
    }

    for (int i = 0; i < meterCnt; i++) {
        if (strcmp(pMeter[i].devSn, devSn) == 0)
        {
            pMeter[i].p = powerRecv;
            pMeter[i].lastUpdateTs = time(NULL);
        }
    }
}

static void sum_total_power(meter_t *pMeter, int meterCnt, double *total_power)
{
    if (pMeter == NULL || total_power == NULL) {
        dbg_syslog(LOG_ERR, "invalid argument! pMeter: %p, total_power: %p", pMeter, total_power);
        return;
    }

    for (int i = 0; i < meterCnt; i++) {
        *total_power += pMeter[i].p;
    }
}

static void pub_total_power(multi_meters_var_t *var, const char *inIdentifier, const char *powerType, double totalPower)
{
    CHAR tagNode[64] = {0};
    cJSON *tagData = NULL;
    CHAR *payload = NULL;

    if (var == NULL || inIdentifier == NULL || powerType == NULL) {
        dbg_syslog(LOG_ERR, "invalid argument! var: %p, powerType: %p", var, powerType);
        return;
    }

    tagData = cJSON_CreateObject();
    if (tagData == NULL) {
        dbg_syslog(LOG_ERR, "cjson create object failed!");
        goto cleanup;
    }

    dbg_syslog(LOG_DEBUG, "inIdentifier : %s", inIdentifier);
    dbg_syslog(LOG_DEBUG, "powerType    : %s", powerType);
    dbg_syslog(LOG_DEBUG, "totalPower   : %lf", totalPower);

    cJSON_AddStringToObject(tagData, "identifier", inIdentifier);
    cJSON_AddStringToObject(tagData, "sn", var->snStr);
    cJSON_AddNumberToObject(tagData, "timestamp", time(NULL));
    cJSON_AddNumberToObject(tagData, "mi", DEFAULT_MI);
    sprintf(tagNode, "{\"%s\":%f}", powerType, totalPower);
    cJSON_AddStringToObject(tagData, "tag_node", tagNode);

    payload = cJSON_Print(tagData);

    char topic[TOPIC_MAX_LEN] = {0};
    snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/%s/device/%s/data/service/%s", var->snStr, "aaa", var->destSn, inIdentifier);
    dbg_syslog(LOG_INFO, "send topic  : %s", topic);
    dbg_syslog(LOG_INFO, "send payload: %s", payload);
    ipc_session_publish(var->session, topic, payload, strlen(payload));

cleanup:
    if (tagData) {
        cJSON_Delete(tagData);
    }
    if (payload) {
        free(payload);
    }    
}


static void report_total_power(multi_meters_var_t *var, const char *powerType, double totalPower)
{
    char *inIdentifier = NULL;

    if (var == NULL || powerType == NULL) {
        dbg_syslog(LOG_ERR, "invalid argument! var: %p, powerType: %p", var, powerType);
        return;
    }

    if (strcmp(powerType, "p1") == 0) {
        if (var->metersState.isP1MeterOffline) {
            dbg_syslog(LOG_ERR, "there is grid meter offline! stop reporting p1!!!");
            return;
        }

        inIdentifier = "gridpower";
    }
    else if (strcmp(powerType, "p3") == 0) {
        if (var->loadMeterCnt != 0) {
            if (var->metersState.isP3MeterOffline) {
                dbg_syslog(LOG_ERR, "there is load meter offline! stop reporting p3!!!");
                return;
            }
        }
        else {
            if (var->metersState.isP1MeterOffline || var->metersState.isP3MeterOffline) {
                dbg_syslog(LOG_ERR, "there is meter offline! stop reporting p3!!! isP1MeterOffline: %d, isP3MeterOffline: %d", 
                            var->metersState.isP1MeterOffline, var->metersState.isP3MeterOffline);
                return;
            }
        }
    
        inIdentifier = "loadpower";
    }
    else {
        dbg_syslog(LOG_ERR, "invalid power type: %s", powerType);
    }

    pub_total_power(var, inIdentifier, powerType, totalPower);
}

static int analyze_msg_recv(char *payload, mqtt_recv_msg_t *mqtt_recv_msg, double *powerRecv)
{
    cJSON *root    = NULL;
    cJSON *tags    = NULL;
    CHAR  *tagNode = NULL;
    int result     = 0;

    if (payload == NULL || mqtt_recv_msg == NULL || powerRecv == NULL) {
        dbg_syslog(LOG_ERR, "invalid argument! payload: %p, mqtt_recv_msg: %p, powerRecv: %p", payload, mqtt_recv_msg, powerRecv);
        result = -1;
        goto cleanup;
    }
    
    root = cJSON_Parse(payload);
    if (root == NULL) {
        result = -1;
        goto cleanup;
    }

    GET_JSON_VALUE_STRING_LOG(root, "identifier", mqtt_recv_msg->identifier);
    GET_JSON_VALUE_STRING_LOG(root, "sn", mqtt_recv_msg->devSn);
    GET_JSON_VALUE_INT_LOG(root, "mi", mqtt_recv_msg->mi);
    GET_JSON_VALUE_INT_LOG(root, "timestamp", mqtt_recv_msg->timestamp);

    GET_JSON_VALUE_DY_STRING_LOG(root, "tag_node", tagNode);
    if (tagNode == NULL) {
        dbg_syslog(LOG_ERR, "tagNode is not exist in mqtt msg received");
        result = -1;
        goto cleanup;
    }

    tags = cJSON_Parse(tagNode);
    if (tags == NULL) {
        dbg_syslog(LOG_ERR, "cJSON_Parse tagNode failed!");
        result = -1;
        goto cleanup;
    }
    GET_JSON_VALUE_DOUBLE_LOG(tags, POWER_TAG, *powerRecv);

cleanup:
    if (tags) {
        cJSON_Delete(tags);
    }
    if (tagNode) {
        free(tagNode);
    }
    if (root) {
        cJSON_Delete(root);
    }

    return result;
}

static int handle_mqtt_msg_recv(void *obj, ipc_msg_t *mqtt_msg)
{
    double total_p1  = 0;
    double total_p2  = 0;
    double total_p3  = 0;
    double powerRecv = 0;
    mqtt_recv_msg_t mqtt_recv_msg = {0};
    multi_meters_var_t *var = NULL;

    if (obj == NULL || mqtt_msg == NULL) {
        dbg_syslog(LOG_ERR, "invalid argument! obj: %p, mqtt_msg: %p", obj, mqtt_msg);
        return -1;
    }

    dbg_syslog(LOG_INFO, "recv topic  : %s", mqtt_msg->topic);
    dbg_syslog(LOG_INFO, "recv payload: %s", mqtt_msg->payload);

    if (analyze_msg_recv(mqtt_msg->payload, &mqtt_recv_msg, &powerRecv) != 0) {
        dbg_syslog(LOG_ERR, "analyze_msg_recv() failed!");
        return -1;
    }

    var = (multi_meters_var_t *)obj;

    if (strstr(mqtt_recv_msg.devSn, "_gridmeter")) {
        store_power_recv(var->pgridMeters, var->gridMeterCnt, mqtt_recv_msg.devSn, powerRecv);
        sum_total_power(var->pgridMeters, var->gridMeterCnt, &total_p1);
        report_total_power(var, "p1", total_p1);
    }
    else if (strstr(mqtt_recv_msg.devSn, "_loadmeter")) {
        store_power_recv(var->ploadMeters, var->loadMeterCnt, mqtt_recv_msg.devSn, powerRecv);
        sum_total_power(var->ploadMeters, var->loadMeterCnt, &total_p3);
        report_total_power(var, "p3", total_p3);
    }
    else if (strstr(mqtt_recv_msg.devSn, "_bmsmeter")) {
        store_power_recv(var->pbmsMeters, var->bmsMeterCnt, mqtt_recv_msg.devSn, powerRecv);

        if (var->loadMeterCnt != 0) {
            dbg_syslog(LOG_DEBUG, "there is load meter, so bms meter is not used!");
            return 0;
        }

        sum_total_power(var->pgridMeters, var->gridMeterCnt, &total_p1);
        sum_total_power(var->pbmsMeters, var->bmsMeterCnt, &total_p2);
        total_p3 = total_p1 - total_p2;
        report_total_power(var, "p3", total_p3);
    }

    return 0;
}

static void multi_meters_subscribe_all(multi_meters_var_t *var)
{
    CHAR topic[TOPIC_MAX_LEN] = {0};
    snprintf(topic, TOPIC_MAX_LEN, "ipc/+/+/device/+/data/service/%s", POWER_IDENTIFIER);
    dbg_syslog(LOG_NOTICE, "sub_topic: %s\r\n", topic);
    ipc_session_subscribe(var->session, topic);
}

static int mqtt_client_init(multi_meters_var_t *var)
{
    char clientId[MAX_CLIENT_ID_LEN] = {0};

    if (var == NULL) {
        dbg_syslog(LOG_ERR, "invalid argument!");
        return -1;
    }
    
    snprintf(clientId, MAX_CLIENT_ID_LEN, "INT_multi_meters_%s", var->snStr);
    var->session = ipc_session_new(clientId, (void *)var, IPC_DEFAULT);
    if (var->session == NULL) {
        dbg_syslog(LOG_ERR, "session create failed!!!");
        return -1;
    }
    ipc_session_set_callbacks(var->session, handle_mqtt_msg_recv, handle_mqtt_conn_state_change);
    multi_meters_subscribe_all(var);
    ipc_session_start(var->session);

    return 0;
}

static int get_meters_num(cJSON *subStructure, const char *meterType)
{
    int meterCnt = 0;

    if (subStructure == NULL || meterType == NULL) {
        dbg_syslog(LOG_ERR, "invalid argument! subStructure: %p, meterType: %p", subStructure, meterType);
        return -1;
    }

    char type[DEV_TYPE_MAX_LEN] = {'\0'};
    int num = cJSON_GetArraySize(subStructure);
    for (int i = 0; i < num; i++) {
        cJSON *dev = cJSON_GetArrayItem(subStructure, i);
        if (dev == NULL) {
            dbg_syslog(LOG_ERR, "cJSON_GetArrayItem() failed!");
            return -1;
        }

        GET_JSON_VALUE_STRING_LOG(dev, "type", type);
        if (strcmp(type, meterType) == 0) {
            meterCnt++;
        }
    }

    return meterCnt;
}

static int get_same_type_meters_cfg(cJSON *root, const char *subStructure, const char *meterType, const char *snStr, meter_t **pMeters, int *meterCnt)
{
    CHAR default_devSn[SN_MAX_LEN] = {'\0'};
    char type[DEV_TYPE_MAX_LEN] = {'\0'};

    if (root == NULL || subStructure == NULL || meterType == NULL || snStr == NULL || pMeters == NULL || meterCnt == NULL) {
        dbg_syslog(LOG_ERR, "invalid argument! root: %p, subStructure: %p, meterType: %p, snStr: %p, pMeters: %p, meterCnt: %p", 
                    root, subStructure, meterType, snStr, pMeters, meterCnt);
        return -1;
    }

    cJSON *subStructureItem = cJSON_GetObjectItemCaseSensitive(root, subStructure);
    if (subStructureItem == NULL) {
        dbg_syslog(LOG_NOTICE, "there is no subStructure(%s)", subStructure);
        return 0;
    }

    *meterCnt = get_meters_num(subStructureItem, meterType);
    dbg_syslog(LOG_NOTICE, "subStructure(%s) cnt: %d", subStructure, *meterCnt);
    if (*meterCnt == -1) {
        dbg_syslog(LOG_ERR, "get_meters_num() failed!");
        return -1;
    }
    else if (*meterCnt == 0) {
        return 0;
    }

    *pMeters = (meter_t *)calloc(*meterCnt, sizeof(meter_t));
    if (*pMeters == NULL) {
        dbg_syslog(LOG_ERR, "calloc pMeters failed");
        return -1;
    }

    int idx = 0;
    int num = cJSON_GetArraySize(subStructureItem);
    for (int i = 0; i < num; i++) {
        cJSON *dev = cJSON_GetArrayItem(subStructureItem, i);
        if (dev == NULL) {
            dbg_syslog(LOG_ERR, "cJSON_GetArrayItem() failed!");
            return -1;
        }

        GET_JSON_VALUE_STRING_LOG(dev, "type", type);
        if (strcmp(type, meterType) != 0) {
            continue;
        }

        GET_JSON_VALUE_STRING_LOG(dev, "type", (*pMeters)[idx].type);
        GET_JSON_VALUE_STRING_LOG_WITH_DEFAULT(dev, "belong_gw", (*pMeters)[idx].belongGw, snStr);
        memset(default_devSn, 0, SN_MAX_LEN);
        snprintf(default_devSn, SN_MAX_LEN, "%s%s%s", snStr, "_", (*pMeters)[idx].type);
        GET_JSON_VALUE_STRING_LOG_WITH_DEFAULT(dev, "dev_sn", (*pMeters)[idx].devSn, default_devSn);
        (*pMeters)[idx].p = 0;
        (*pMeters)[idx].lastUpdateTs = time(NULL);

        idx++;
    }

    return 0;
}

static void print_cfg_load_result(multi_meters_var_t *var)
{
    int i = 0;
    meter_t *bms_meter = NULL;

    if (var == NULL) {
        dbg_syslog(LOG_ERR, "invalid argument!");
        return;
    }

    dbg_syslog(LOG_NOTICE, "----------------destSn-----------------");
    dbg_syslog(LOG_NOTICE, "destSn: %s", var->destSn);

    dbg_syslog(LOG_NOTICE, "----------------grid meters cfg-----------------");
    dbg_syslog(LOG_NOTICE, "gridMeterCnt: %d", var->gridMeterCnt);
    for (i = 0; i < var->gridMeterCnt; i++) {
        dbg_syslog(LOG_NOTICE, "pgridMeters[%d].type: %s", i, var->pgridMeters[i].type);
        dbg_syslog(LOG_NOTICE, "pgridMeters[%d].devSn: %s", i, var->pgridMeters[i].devSn);
        dbg_syslog(LOG_NOTICE, "pgridMeters[%d].belongGw: %s", i, var->pgridMeters[i].belongGw);
    }

    dbg_syslog(LOG_NOTICE, "----------------load meters cfg-----------------");
    dbg_syslog(LOG_NOTICE, "loadMeterCnt: %d", var->loadMeterCnt);
    for (i = 0; i < var->loadMeterCnt; i++) {
        dbg_syslog(LOG_NOTICE, "ploadMeters[%d].type: %s", i, var->ploadMeters[i].type);
        dbg_syslog(LOG_NOTICE, "ploadMeters[%d].devSn: %s", i, var->ploadMeters[i].devSn);
        dbg_syslog(LOG_NOTICE, "ploadMeters[%d].belongGw: %s", i, var->ploadMeters[i].belongGw);
    }

    dbg_syslog(LOG_NOTICE, "----------------bms meters cfg-----------------");
    dbg_syslog(LOG_NOTICE, "bmsMeterCnt: %d", var->bmsMeterCnt);
    for (i = 0; i < var->bmsMeterCnt; i++) {
        dbg_syslog(LOG_NOTICE, "pbmsMeters[%d].type: %s", i, var->pbmsMeters[i].type);
        dbg_syslog(LOG_NOTICE, "pbmsMeters[%d].devSn: %s", i, var->pbmsMeters[i].devSn);
        dbg_syslog(LOG_NOTICE, "pbmsMeters[%d].belongGw: %s", i, var->pbmsMeters[i].belongGw);
    }
    dbg_syslog(LOG_NOTICE, "------------------------------------------------");
}

static int load_meters_cfg(multi_meters_var_t *var, const char *cfg_file)
{
    int result  = 0;
    CHAR *pdata = NULL;
    cJSON *root = NULL;
    char default_destSn[SN_MAX_LEN] = {'\0'};

    if (var == NULL || cfg_file == NULL) {
        dbg_syslog(LOG_ERR, "invalid argument! var: %p cfg_file: %p", var, cfg_file);
        result = -1;
        goto cleanup;
    }

    dbg_syslog(LOG_NOTICE, "cfg_file: %s", cfg_file);
    pdata = read_file_data(cfg_file);
    if (pdata == NULL) {
        dbg_syslog(LOG_ERR, "read cfg_file: %s error", cfg_file);
        result = -1;
        goto cleanup;
    }

    root = cJSON_Parse(pdata);
    if (!root) {
        dbg_syslog(LOG_ERR, "parse cfg_file: %s error", cfg_file);
        result = -1;
        goto cleanup;
    }

    snprintf(default_destSn, SN_MAX_LEN, "%s%s", var->snStr, "_EMSCTRL");
    GET_JSON_VALUE_STRING_LOG_WITH_DEFAULT(root, "sn", var->destSn, default_destSn);

    // get grid meters cfg
    if (get_same_type_meters_cfg(root, "grid", "gridmeter", var->snStr, &(var->pgridMeters), &(var->gridMeterCnt)) != 0) {
        dbg_syslog(LOG_ERR, "get gridmeter cfg failed");
        result = -1;
        goto cleanup;
    }

    // get load meters cfg
    if (get_same_type_meters_cfg(root, "load", "loadmeter", var->snStr, &(var->ploadMeters), &(var->loadMeterCnt)) != 0) {
        dbg_syslog(LOG_ERR, "get loadmeter cfg failed");
        result = -1;
        goto cleanup;
    }

    // get bms meters cfg
    if (get_same_type_meters_cfg(root, "storage", "bmsmeter", var->snStr, &(var->pbmsMeters), &(var->bmsMeterCnt)) != 0) {
        dbg_syslog(LOG_ERR, "get bmsmeter cfg failed");
        result = -1;
        goto cleanup;
    }

    print_cfg_load_result(var);

cleanup:
    if (root) {
        cJSON_Delete(root);
    }
    if (pdata) {
        free(pdata);
    }

    return result;
}

static void multi_meters_init(multi_meters_var_t *var)
{
    if (var == NULL) {
        dbg_syslog(LOG_ERR, "multi_meters init failed! invalid argument!");
        return;
    }

    get_board_sn(var->snStr);
    dbg_syslog(LOG_NOTICE, "board SN:%s", var->snStr);

    var->metersState.isP1MeterOffline = false;
    var->metersState.isP3MeterOffline = false;

    if (access(EMS_STRUCTURE_CFG_FILE, F_OK) == 0) {
        if (load_meters_cfg(var, EMS_STRUCTURE_CFG_FILE) != 0) {
            dbg_syslog(LOG_ERR, "multi_meters init failed! load_meters_cfg failed!");
            return;
        }
    }
    else {
        if (load_meters_cfg(var, MULTI_METERS_CFG_FILE) != 0) {
            dbg_syslog(LOG_ERR, "multi_meters init failed! load_meters_cfg failed!");
            return;
        }
    }

    if (mqtt_client_init(var) != 0) {
        dbg_syslog(LOG_ERR, "multi_meters init failed! mqtt_client_init failed!");
        return;
    }

    var->period_poll_timer = my_timer_create();
    if (var->period_poll_timer <= 0) {
        dbg_syslog(LOG_ERR, "multi_meters init failed! create period_poll_timer failed!");
    }
    my_timer_set(var->period_poll_timer, 1, 3000);

    dbg_syslog(LOG_NOTICE, "multi_meters init done");
}

static void determine_meters_offline(meter_t *meter, int meterCnt, bool *isMeterOffline)
{
    bool isOffline = false;

    if (meter == NULL || isMeterOffline == NULL) {
        dbg_syslog(LOG_ERR, "invalid argument! meter: %p, isMeterOffline: %p", meter, isMeterOffline);
        return;
    }

    time_t curTime = time(NULL);

    for (int i = 0; i < meterCnt; i++) {
        if (curTime - meter[i].lastUpdateTs <= 10) {
            dbg_syslog(LOG_DEBUG, "%s is online", meter[i].devSn);
            continue;
        }

        isOffline = true;
        dbg_syslog(LOG_ERR, "%s is offline!", meter[i].devSn);
    }

    if (isOffline) {
        *isMeterOffline = true;
    }
    else {
        *isMeterOffline = false;
    }

    return;
}

static int period_service_poll(multi_meters_var_t *var)
{
    if (var == NULL) {
        dbg_syslog(LOG_ERR, "invalid argument!");
        return -1;
    }

    time_t curTime = time(NULL);
    TIMER_CONFIRM(var->period_poll_timer);

    // determine grid meters offline
    if (var->gridMeterCnt > 0) {
        determine_meters_offline(var->pgridMeters, var->gridMeterCnt, &var->metersState.isP1MeterOffline);
    }

    // determine load meters offline
    if (var->loadMeterCnt > 0) {
        determine_meters_offline(var->ploadMeters, var->loadMeterCnt, &var->metersState.isP3MeterOffline);
    }

    // determine bms meters offline
    if (var->bmsMeterCnt > 0) {
        if (var->loadMeterCnt > 0) {
            dbg_syslog(LOG_DEBUG, "there is load meter, bms meter will not be checked!");
            return 0;
        }
        determine_meters_offline(var->pbmsMeters,  var->bmsMeterCnt,  &var->metersState.isP3MeterOffline);
    }

    return 0;
}

static void multi_meters_loop(multi_meters_var_t *var)
{
    fd_set rset;
    struct timeval timeout;
    int ret = -1, maxfd, i;

    if (var == NULL) {
        dbg_syslog(LOG_ERR, "invalid argument!");
        return;
    }

    while (1) {
        SELECT_INIT();
        SELECT_ADD_FD(var->period_poll_timer);
        timeout.tv_usec = 0;
        timeout.tv_sec = 5;
        ret = select(maxfd + 1, &rset, 0, 0, &timeout);
        if (ret < 0) {
            int error = errno;
            dbg_syslog(LOG_ERR, "errno %d\n", error);

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

void main()
{
    multi_meters_var_t var = {0};
    memset(&var, 0, sizeof(var));
    lnxall_loglevel_set(LNXALL_LOGNOTICE, 1);

    multi_meters_init(&var);
    multi_meters_loop(&var);
}