#include <errno.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <sys/select.h>
#include <sys/stat.h>
#include <sys/syscall.h>
#include <sys/time.h>
#include <sys/types.h>
#include <unistd.h>
#include "main.h"
#include "energy_config.h"
#include "data_process.h"
#include "ctrl_process.h"
#include "data_collect.h"
#include "public.h"
#include <getopt.h>
#include <float.h>

#define PCS_DATA_REPORT_IDENTIFIER "Read32_40Data1"
#define PCS_VOL_CURRENT_IDENTIFIER "Read101_125Data"
#define METER_POWRE_REPORT_IDENTIFIER "sanxianggongluxinxi"
#define METER_ENERGY_REPORT_IDENTIFIER "yougongzongdianliang"
#define METER_VOL_CURRENT_IDENTIFIER "sanxiangcv"
#define BMS_READ_STATE_IDENTIFIER "Read100C_1016"
#define BMS_CHARGE_DISCHARGE_IDENTIFIER "Cluster_charge_discharge"

static energy_var_t sys_var = {0};

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

    while (1)
    {
        SELECT_INIT();
        SELECT_ADD_FD(var->calculate_timer);

        timeout.tv_usec = 0;
        timeout.tv_sec = 1;

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

            if (errno == EINTR)
            {
                continue;
            }
            else
            {
                break;
            }
        }
        else if (ret > 0)
        {
            if (var->calculate_timer > 0 && FD_ISSET(var->calculate_timer, &rset))
            {
                FD_CLR(var->calculate_timer, &rset);
                // 读数据
            }
        }
        if (var->eng_ass_grp_num == 0)
        {
            dbg_syslog(LOG_ERR, "pcs count=0, 检查数据库up_device_id设置");
        }
    }
}

static void energy_ex_subscribe_all(energy_var_t *var)
{
    ipc_session_t *session = var->exsession;
    char topic[TOPIC_MAX_LEN] = {0};
    int i = 0;
    for (i = 0; i < var->eng_ass_grp_num; i++)
    {
        if (var->grp[i].eng_ass_grp_cfg.pcs_sn)
        {
            snprintf(topic, TOPIC_MAX_LEN, "/ems/%s/service", var->grp[i].eng_ass_grp_cfg.pcs_sn);
            dbg_syslog(LOG_INFO, "topic:%s\r\n", topic);
            ipc_session_subscribe(session, topic);
            ipc_session_subscribe(session, topic);
        };
        if (var->grp[i].eng_ass_grp_cfg.meter_sn)
        {
            snprintf(topic, TOPIC_MAX_LEN, "/ems/%s/service", var->grp[i].eng_ass_grp_cfg.meter_sn);
            dbg_syslog(LOG_INFO, "topic:%s\r\n", topic);
            ipc_session_subscribe(session, topic);
        };
        if (var->grp[i].eng_ass_grp_cfg.bms_dev_num != 0)
        {
            for (int j = 0; j < var->grp[i].eng_ass_grp_cfg.bms_dev_num; j++)
            {
                snprintf(topic, TOPIC_MAX_LEN, "/ems/%s/service", var->grp[i].eng_ass_grp_cfg.bms_sn[j]);
                dbg_syslog(LOG_INFO, "topic:%s\r\n", topic);
                ipc_session_subscribe(session, topic);
            }
        }
    }
}

static void energy_subscribe_all(energy_var_t *var)
{
    ipc_session_t *session = var->session;
    char topic[TOPIC_MAX_LEN] = {0};
    int i = 0;
    // TODO 通过查询绑定关系 订阅各个设备
    // 下面SN 重新格式化,identifier 代码里内置写死
    // 订阅 data/service 上行
    // 订阅 Set_Rglt 下行
    snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/+/device/%s/data/Set_Rglt", var->sn_str, var->eng_sn_str);
    dbg_syslog(LOG_NOTICE, "topic:%s\r\n", topic);
    ipc_session_subscribe(session, topic);
    snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/+/device/%s/data/Set_Rglt", var->sn_str, var->vpcs_sn_str);
    dbg_syslog(LOG_NOTICE, "topic:%s\r\n", topic);
    ipc_session_subscribe(session, topic);
    if (var->runMode == E_ONGW_MODE)
    {
        for (i = 0; i < var->eng_ass_grp_num; i++)
        {
            if (var->grp[i].eng_ass_grp_cfg.pcs_sn)
            {
                snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/+/device/%s/data/service/%s", "+", var->grp[i].eng_ass_grp_cfg.pcs_sn, "+");
                dbg_syslog(LOG_NOTICE, "topic:%s\r\n", topic);
                ipc_session_subscribe(session, topic);
                // snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/+/device/%s/data/service/%s", "+", var->grp[i].eng_ass_grp_cfg.pcs_sn, "+");
                // dbg_syslog(LOG_INFO, "topic:%s\r\n", topic);
                // ipc_session_subscribe(session, topic);
            };
            if (var->grp[i].eng_ass_grp_cfg.meter_sn)
            {
                snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/+/device/%s/data/service/%s", "+", var->grp[i].eng_ass_grp_cfg.meter_sn, METER_POWRE_REPORT_IDENTIFIER);
                dbg_syslog(LOG_NOTICE, "topic:%s\r\n", topic);
                ipc_session_subscribe(session, topic);
                snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/+/device/%s/data/service/%s", "+", var->grp[i].eng_ass_grp_cfg.meter_sn, METER_ENERGY_REPORT_IDENTIFIER);
                dbg_syslog(LOG_NOTICE, "topic:%s\r\n", topic);
                ipc_session_subscribe(session, topic);
                snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/+/device/%s/data/service/%s", "+", var->grp[i].eng_ass_grp_cfg.meter_sn, METER_VOL_CURRENT_IDENTIFIER);
                dbg_syslog(LOG_NOTICE, "topic:%s\r\n", topic);
                ipc_session_subscribe(session, topic);
            };
            if (var->grp[i].eng_ass_grp_cfg.bms_dev_num != 0)
            {
                for (int j = 0; j < var->grp[i].eng_ass_grp_cfg.bms_dev_num; j++)
                {
                    snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/+/device/%s/data/service/%s", "+", var->grp[i].eng_ass_grp_cfg.bms_sn[j], "+");
                    dbg_syslog(LOG_NOTICE, "topic:%s\r\n", topic);
                    ipc_session_subscribe(session, topic);
                    // snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/+/device/%s/data/service/%s", "+", var->grp[i].eng_ass_grp_cfg.bms_sn[j], "+");
                    // dbg_syslog(LOG_INFO, "topic:%s\r\n", topic);
                    // ipc_session_subscribe(session, topic);
                }
            }
        }
    }
}

static int energy_ex_mqtt_handle_recv_msg(void *obj, ipc_msg_t *mqtt_msg)
{
    cJSON *node = NULL;
    cJSON *sub = NULL;
    energy_var_t *var = (energy_var_t *)obj;
    if (strstr(mqtt_msg->topic, "service"))
    {
        node = cJSON_Parse(mqtt_msg->payload);
        if (!node)
        {
            dbg_syslog(LOG_ERR, "get node: %s error", mqtt_msg->topic);
            return 0;
        }
        char dev_sn[ENG_ASS_SN_LEN_MAX] = "";
        char identifier[64] = "";
        GET_JSON_VALUE_STRING(node, "identifier", identifier);
        GET_JSON_VALUE_STRING(node, "sn", dev_sn);
        sub = cJSON_GetObjectItemCaseSensitive(node, "tags");
        into_data_col_raw_msg(var, sub, dev_sn);
        cJSON_Delete(node);
    }
    return 0;
}

static int energy_mqtt_handle_recv_msg(void *obj, ipc_msg_t *mqtt_msg)
{
    energy_var_t *var = (energy_var_t *)obj;
    // dbg_syslog(LOG_INFO, "received MQTT topic:%s payload length:%d", mqtt_msg->topic, mqtt_msg->payloadLen);
    cJSON *node = NULL;
    char *tagNode = NULL;
    if (strstr(mqtt_msg->topic, "data/service"))
    {
        // node = api_get_json_node(mqtt_msg->payload);
        // goto out;
        node = cJSON_Parse(mqtt_msg->payload);
        if (!node)
        {
            dbg_syslog(LOG_ERR, "get node: %s error", mqtt_msg->topic);
            return 0;
        }
        char dev_sn[ENG_ASS_SN_LEN_MAX] = "";
        char identifier[64] = "";
        GET_JSON_VALUE_STRING(node, "identifier", identifier);
        GET_JSON_VALUE_STRING(node, "sn", dev_sn);
        GET_JSON_VALUE_DY_STRING(node, "tag_node", tagNode);
        if (tagNode)
        {
            cJSON *sub = cJSON_Parse(tagNode);
            into_data_col_raw_msg(var, sub, dev_sn);
            cJSON_Delete(sub);
            free(tagNode);
        }
        cJSON_Delete(node);
    }
    else if (strstr(mqtt_msg->topic, "data/Set_Rglt"))
    {
        data_ctrl_raw_msg(var, mqtt_msg->payload);
    }
    return 0;
}

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

    snprintf(clientId, MAX_CLIENT_ID_LEN, "INT_energy_%s", var->sn_str);
    if (var->runMode == E_ONGW_MODE)
    {
        var->session = ipc_session_new(clientId, (void *)var, IPC_DEFAULT);
        if (var->session == NULL)
            return -1;

        ipc_session_set_callbacks(var->session, energy_mqtt_handle_recv_msg, NULL);
        energy_subscribe_all(var);
        ipc_session_start(var->session);
    }
    else
    {
        var->session = ipc_session_new(clientId, (void *)var, IPC_DEFAULT);
        if (var->session == NULL)
            return -1;
        ipc_session_set_callbacks(var->session, energy_mqtt_handle_recv_msg, NULL);
        energy_subscribe_all(var);
        ipc_session_start(var->session);  

        var->exsession = ipc_session_new(clientId, (void *)var, IPC_DEFAULT);
        dbg_syslog(LOG_INFO,"exsession:%p",var->exsession);
        dbg_syslog(LOG_INFO, "hostname:%s port:%d username:%s password:%s", var->mqttBroker.hostname, var->mqttBroker.port, var->mqttBroker.username, var->mqttBroker.password);
        ipc_session_set_address(var->exsession, var->mqttBroker.hostname, var->mqttBroker.port, var->mqttBroker.username, var->mqttBroker.password);
        ipc_session_set_opts(var->exsession, MQTT_QOS_LEVLE0, KEEP_ALIVE_MAX);
        ipc_session_set_callbacks(var->exsession, energy_ex_mqtt_handle_recv_msg, NULL);
        energy_ex_subscribe_all(var);
        ipc_session_start(var->exsession);
    }
    return 0;
}

static void energy_var_init(energy_var_t *var)
{
    int i;
    for (i = 0; i < var->eng_ass_grp_num; i++)
    {
        var->grp[i].eng_ass_grp_ctrl.connect = -1;
        var->grp[i].eng_ass_grp_ctrl.run_state = -1;
        var->grp[i].eng_ass_grp_ctrl.protect_level = -1;
        var->grp[i].eng_ass_grp_ctrl.conn_high_power = -1;
        var->grp[i].eng_ass_grp_ctrl.switch_on_state = -1;
        var->grp[i].eng_ass_grp_ctrl.relay = -1;
        var->grp[i].eng_ass_grp_data.bms_charge_current = -100;
        var->grp[i].eng_ass_grp_data.bms_discharge_current = 100;
        var->grp[i].eng_ass_grp_data.pcs_connect_state = -1;
        var->grp[i].eng_ass_grp_data.pcs_run_state = -1;
        var->grp[i].eng_ass_grp_data.bms_work_state = -1;
        var->grp[i].eng_ass_grp_data.bms_cha_discha_state = -1;
        var->grp[i].eng_ass_grp_data.bms_switch_on_state = -1;
        var->grp[i].eng_ass_grp_data.disconnect_v = -1;
        var->grp[i].eng_ass_grp_data.fault_level = -1;
        var->grp[i].eng_ass_grp_data.bms_relay_p = -1;
        var->grp[i].eng_ass_grp_data.bms_relay_n = -1;
        var->grp[i].eng_ass_grp_data.pcs_run_flag = -1;
        var->grp[i].eng_ass_grp_data.pcs_stop_flag = -1;
        var->grp[i].eng_ass_grp_data.pcs_standby_flag = -1;
        var->grp[i].eng_ass_grp_data.pcs_remote = -1;
        var->grp[i].eng_ass_grp_data.bms_avaiable_cha_curr = -FLT_MAX;
        var->grp[i].eng_ass_grp_data.bms_avaiable_disc_curr = FLT_MAX;
        var->grp[i].eng_ass_grp_data.bms_over_hold = 0;
        var->grp[i].eng_ass_grp_data.bms_under_hold = 0;

    }

    if (!var->eng_ass_data)
    {
        var->eng_ass_data = calloc(1, sizeof(eng_ass_data_t));
        if (!var->eng_ass_data)
        {
            // dbg_syslog(LOG_WARNING, "memory not sufficient!!!%d\n", strerror(errno));
            return;
        }
        var->eng_ass_data->protect_level = -1;
        var->eng_ass_data->fault_level = -1;
        var->eng_ass_data->pcs_fault = -1;
        var->eng_ass_data->pcs_comm_fault = -1;
        var->eng_ass_data->pcs_conn_fb_fault = -1;
        var->eng_ass_data->pcs_run_fb_fault = -1;
        var->eng_ass_data->meter_comm_fault = -1;
        var->eng_ass_data->bms_fault = -1;
        var->eng_ass_data->bms_comm_fault = -1;
        var->eng_ass_data->on_service_state = -1;
    }
    var->ems_run_mode = -1;
    var->supportStandy = 0;
}

static int energy_init(energy_var_t *var)
{
    get_board_sn(var->sn_str);
    get_process_name(var->proc_name);
    snprintf(var->vpcs_sn_str, sizeof(var->vpcs_sn_str), "virtual_pcs");
    snprintf(var->eng_sn_str, sizeof(var->eng_sn_str), "%s_ENGASS", var->sn_str);
    snprintf(var->bmsmeter_sn_str, sizeof(var->bmsmeter_sn_str), "%s_bmsmeter", var->sn_str);
    snprintf(var->emsctrl_sn_str, sizeof(var->emsctrl_sn_str), "%s_EMSCTRL", var->sn_str);
    if (access(EMS_STRUCTURE_CFG_FILE, F_OK) == 0) {
        energy_load_config_from_ems_structure(var, EMS_STRUCTURE_CFG_FILE);
    }
    else {
        energy_load_config(var, ENG_ASS_CFG_FILE);
    }
    if (var->dev_online_timeout == 0  || var->dev_online_timeout == 15)//之前默认15s
    {
        var->dev_online_timeout = 150; //default 150s超时
    }
    dbg_syslog(LOG_NOTICE, "dev_online_timeout: %d", var->dev_online_timeout);
    energy_var_init(var);
    energy_mqtt_client_init(var);
    energy_services_load(var,var->module);
    data_process_init(var);
     //added by yut
    if (access(ENG_SVC_STATE, F_OK) == 0) {
        energy_state_config(var,ENG_SVC_STATE);
    }
    //ended by yut
    if (access(EMS_CFG_SAVE_FILE, F_OK) == 0) {
        energy_state_config_load_ems(var,EMS_CFG_SAVE_FILE);
    }
    if (var->ems_run_mode <= 0) // 手动模式
    {
        var->eng_ass_data->SOCMAX = -1;
        var->eng_ass_data->SOCMIN = -1;
		var->eng_ass_data->SINGLE_MAX = -1;
        var->eng_ass_data->SINGLE_MIN = -1;
        dbg_syslog(LOG_NOTICE,"ctrlMode:%d SOCMAX: %f SOC min:%f ",var->ems_run_mode,var->eng_ass_data->SOCMAX,var->eng_ass_data->SOCMIN);
    }else
    {//自动模式
        if(var->eng_ass_data->pcs_connect_state != 1)//自动模式离网需切到并网
        {
            var->eng_ass_data->pcs_connect_state = 1;
            abstract_item_t *item = var->module_item;
            if(item)
            {
                if(item->set_pcs_connect)item->set_pcs_connect(var,var->eng_ass_data->pcs_connect_state);
                save_config_file(var);
            }
        }
    }
    if (access(EMS_CFG_FILE, F_OK) == 0) {
        energy_state_config_load_ems_cfg(var,EMS_CFG_FILE);
    }
    dbg_syslog(LOG_NOTICE, "supportStandy: %d", var->supportStandy);
    data_publish_init(var);
    return 0;
}


static void print_help() {
    printf("\n--------------------------------------------------------------------------------\n");
    printf("Usage: energy_assign [OPTIONS]\n");
    printf("\n");
    printf("-v, --version         show this app version\n");
    printf("-h, --help            show this help\n");
    printf("\n--------------------------------------------------------------------------------\n");
    printf("-m, --model [kuaibu]           select model\n");
}

int main(int argc, char *argv[])
{
    int c;
    bool getversion = false;
    openlog(argv[0], LOG_PID, LOG_DAEMON);
    lnxall_loglevel_set(LNXALL_LOGNOTICE,1);
    dbg_syslog(LOG_INFO, "energy assign started\n");
    strcpy(sys_var.module, ENG_ASS_DEFAULT_SERVICE);
    while ((c = getopt(argc, argv, "m:vh")) != -1) {
        switch (c) {
            case 'v':
                getversion = true;
                break;
            case 'm':
                if (optarg)
                    strncpy(sys_var.module, optarg, sizeof(sys_var.module));
                break;
            case 'h':
            case '?':
            default:
                print_help();
                exit(0);
        }
    }
    if (getversion) {
		printf("ver: %s\n",ENG_ASS_SOFT_VER);
        exit(0);
    }
    printf("mode:%s\n",sys_var.module);
    energy_init(&sys_var);
    energy_loop(&sys_var);
    return 0;
}

