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

#include "dy_utils/cJSON.h"
#include "dy_utils/dy_common.h"
#include "ipc_session.h"
#include "modbus/modbus.h"
#include "rs485_slave_common.h"

static void sort_object_cfg_services_outputData(object_cfg_t *p_object_cfg)
{
    int i, j, k;
    object_service_table_t *p_service_table = p_object_cfg->service_tab;
    object_service_cfg_t *p_tmp_service;
    object_service_cfg_t tmp_service;
    object_property_table_t *p_property_table;
    object_property_cfg_t tmp_property;

    for (i = 0; i < p_service_table->serviceCnt - 1; i++)
    {
        for (j = i + 1; j < p_service_table->serviceCnt; j++)
        {
            if (strcmp(p_service_table->service[i].identifier, p_service_table->service[j].identifier) > 0)
            {
                memcpy(&tmp_service, &p_service_table->service[i], sizeof(object_service_cfg_t));
                memcpy(&p_service_table->service[i], &p_service_table->service[j], sizeof(object_service_cfg_t));
                memcpy(&p_service_table->service[j], &tmp_service, sizeof(object_service_cfg_t));
            }
        }
    }

    for (i = 0; i < p_service_table->serviceCnt; i++)
    {
        p_tmp_service = &p_service_table->service[i];
        p_property_table = p_tmp_service->poutput;
        for (j = 0; j < p_property_table->propertyCnt - 1; j++)
        {
            for (k = j + 1; k < p_property_table->propertyCnt; k++)
            {
                if (strcmp(p_property_table->property[j].identifier, p_property_table->property[k].identifier) > 0)
                {
                    memcpy(&tmp_property, &p_property_table->property[j], sizeof(object_property_cfg_t));
                    memcpy(&p_property_table->property[j], &p_property_table->property[k], sizeof(object_property_cfg_t));
                    memcpy(&p_property_table->property[k], &tmp_property, sizeof(object_property_cfg_t));
                }
            }
        }
    }
}

static int rs485_slave_modbus_rtu_init(rs485_slave_thread_t *rs485_slave_thread)
{
    rs485_slave_port_cfg_t *rs485_slave_port_cfg = rs485_slave_thread->rs485_slave_port_cfg;
    const char *dev = rs485_slave_port_cfg->dev;
    int speed = rs485_slave_port_cfg->speed; 
    int bits = rs485_slave_port_cfg->bits;
    char cparity = 'N';
    int stop = rs485_slave_port_cfg->stop;
    modbus_t* ctx = NULL;
    modbus_mapping_t* map = NULL;
    int ret;
    int object_cnt = rs485_slave_thread->var->object_table->object_cnt;
    object_cfg_t *p_object_cfg = NULL;
    int service_cnt;
    int property_cnt;
    int i, j;
    rs485_slave_service_property_t *service_property;
    object_service_table_t *p_service_table = NULL;
    object_service_cfg_t *p_service;
    object_property_table_t *p_property_table;
    int register_cnt;

    // TODO: only support one dev now.
    if (object_cnt != 1)
    {
        dy_syslog(LOG_ERR, "more than one dev, exit!\n");
        exit(0);
    }
    p_object_cfg = &rs485_slave_thread->var->object_table->object[0];
    sort_object_cfg_services_outputData(p_object_cfg);
    p_service_table = p_object_cfg->service_tab;
    service_cnt = p_object_cfg->service_tab->serviceCnt;
    dy_syslog(LOG_DEBUG, "service count %d\n", service_cnt);
    service_property = calloc(service_cnt, sizeof(rs485_slave_service_property_t));
    for (i = 0; i < service_cnt; i++)
    {
        p_service = &p_service_table->service[i];
        property_cnt = p_service->poutput->propertyCnt;
        dy_syslog(LOG_DEBUG, "service index %d, identifier: %s, property count %d\n", i, p_service->identifier, property_cnt);
        service_property[i].cur_cnt = property_cnt;
        service_property[i].p_byte = calloc(property_cnt, sizeof(int));
        service_property[i].p_offset = calloc(property_cnt, sizeof(int));
        p_property_table = p_service->poutput;
        for (j = 0; j < property_cnt; j++)
        {
            if (j == 0) {
                service_property[i].p_offset[j] = 0;
            } else {
                service_property[i].p_offset[j] = service_property[i].p_offset[j - 1] + service_property[i].p_byte[j - 1];
            }
            service_property[i].p_byte[j] = 4; //data_type_register_bytes(p_property_table->property[j].raw_data_type); //only float now
            service_property[i].cur_len += service_property[i].p_byte[j];
            dy_syslog(LOG_DEBUG, "property index %d, identifier: %s, cur_offset: %d, cur_bytes: %d\n", j, p_property_table->property[j].identifier, service_property[i].p_offset[j], service_property[i].p_byte[j]);
        }
        
        if (i == 0) {
            service_property[i].pre_len = 0;
        } else {
            service_property[i].pre_len = service_property[i - 1].pre_len + service_property[i - 1].cur_len;
        }

        dy_syslog(LOG_DEBUG, "service index %d, pre_len %d, cur_len %d\n", i, service_property[i].pre_len, service_property[i].cur_len);
        if (i == service_cnt - 1)
        {
            rs485_slave_thread->services_property.total_len = service_property[i].pre_len + service_property[i].cur_len;
            dy_syslog(LOG_DEBUG, "service total_len %d\n", rs485_slave_thread->services_property.total_len);
        }
    }
    rs485_slave_thread->services_property.services_count = service_cnt;
    rs485_slave_thread->services_property.service_property = service_property;

    register_cnt = rs485_slave_thread->services_property.total_len / 2;
    dy_syslog(LOG_DEBUG, "register count %d\n", register_cnt);

    if (rs485_slave_port_cfg->parity == 1)
    {
        cparity = 'O';
    }
    else if (rs485_slave_port_cfg->parity == 2)
    {
        cparity = 'E';
    }

    ctx = modbus_new_rtu(dev, speed, cparity, bits, stop);
    if (ctx == NULL)
    {
        dy_syslog(LOG_ERR, "modbus_new_rtu failed!\n");
        ret = -1;
        goto end;
    }

    if (modbus_set_slave(ctx, 1) < 0)
    {
        dy_syslog(LOG_ERR, "modbus_set_slave failed!\n");
        ret = -1;
        goto end;
    }

    if (modbus_set_debug(ctx, TRUE) < 0)
    {
        dy_syslog(LOG_ERR, "modbus_set_debug failed!\n");
        ret = -1;
        goto end;
    }

    map = modbus_mapping_new(0, 0, 0, register_cnt);
    if (map == NULL)
    {
        dy_syslog(LOG_ERR, "modbus_mapping_new failed!\n");
        ret = -1;
        goto end;
    }

    rs485_slave_thread->ctx = ctx;
    rs485_slave_thread->map = map;

    return 0;

end:
    if (ctx != NULL)
    {
        modbus_free(ctx);
    }
    if (map != NULL)
    {
        modbus_mapping_free(map);
    }

    return ret;
}

static int rs485_slave_handle_register_data(rs485_slave_thread_t *rs485_slave_thread, ipc_msg_t *mqtt_msg)
{
    int ret = 0;
    cJSON *root = cJSON_Parse(mqtt_msg->payload);
    char identifier[64];
    cJSON *tags = NULL;
    cJSON *tag = NULL;
    float tmp;
    int service_index = 0;
    int i, j;
    object_service_table_t *service_tab = rs485_slave_thread->var->object_table->object[0].service_tab;
    object_property_table_t *property_table;
    rs485_slave_service_property_t *service_property;
    int register_offset;
    int size;

    if (root)
    {
        GET_JSON_VALUE_STRING(root, "identifier", identifier);
        for (i = 0; i < service_tab->serviceCnt; i++)
        {
            if (strcmp(service_tab->service[i].identifier, identifier) == 0)
            {
                service_index = i;
                break;
            }
        }
        service_property = &rs485_slave_thread->services_property.service_property[service_index];
        property_table = service_tab->service[service_index].poutput;
        tags = cJSON_GetObjectItem(root, "tags");
        size = cJSON_GetArraySize(tags);
        for (i = 0; i < size; i++)
        {
            tag = cJSON_GetArrayItem(tags, i);
            if (strcmp(tag->string, "identifier") == 0)
            {
                continue;
            }
            for (j = 0; j < property_table->propertyCnt; j++)
            {
                if (strcmp(tag->string, property_table->property[j].identifier) == 0)
                {
                    register_offset = (service_property->pre_len + service_property->p_offset[j]) / 2;
                    tmp = tag->valuedouble; //TODO: only float
                    modbus_set_float_dcba(tmp, &rs485_slave_thread->map->tab_input_registers[register_offset]);
                    break;
                }
            }
        }
    }
    else
    {
        dy_syslog(LOG_ERR, "json parse error");
        ret = -1;
    }

    cJSON_Delete(root);
    return ret;
}

static int rs485_slave_mqtt_handle_recv_msg(void *obj, ipc_msg_t *mqtt_msg)
{
    rs485_slave_thread_t *rs485_slave_thread = (rs485_slave_thread_t *)obj;

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

    return 0;
}

static void rs485_slave_subscribe_all(rs485_slave_thread_t *rs485_slave_thread)
{
    ipc_session_t *ipc_session = rs485_slave_thread->session;
    char topic[TOPIC_MAX_LEN] = {0};

    if (strlen(rs485_slave_thread->var->mqtt_server_topic.topic) > 0)
    {
        snprintf(topic, TOPIC_MAX_LEN, "%s", rs485_slave_thread->var->mqtt_server_topic.topic);
    }
    else
    {
        snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/+/device/+/data_filtered/service/+", rs485_slave_thread->var->sn_str);
    }
    ipc_session_subscribe(ipc_session, topic);
}

static int rs485_slave_mqtt_client_init(rs485_slave_thread_t *rs485_slave_thread)
{
    char clientId[128] = {0};

    snprintf(clientId, sizeof(clientId), "INT_rs485_slave_%d_%s", rs485_slave_thread->rs485_slave_port, rs485_slave_thread->var->sn_str);
    rs485_slave_thread->session = ipc_session_new(clientId, (void*)rs485_slave_thread, IPC_DEFAULT);
    if (rs485_slave_thread->session == NULL)
        return -1;

    ipc_session_set_callbacks(rs485_slave_thread->session, rs485_slave_mqtt_handle_recv_msg, NULL);
    rs485_slave_subscribe_all(rs485_slave_thread);
    ipc_session_start(rs485_slave_thread->session);
    return 0;
}

static void *rs485_slave_thread_loop(void *param)
{
    rs485_slave_thread_t *rs485_slave_thread = (rs485_slave_thread_t *)param;
    modbus_t *ctx = rs485_slave_thread->ctx;
    modbus_mapping_t *map = rs485_slave_thread->map;
    char query[MODBUS_TCP_MAX_ADU_LENGTH];
    int ret;

    if (modbus_connect(ctx) < 0)
    {
        dy_syslog(LOG_ERR, "modbus_connect failed!\n");
        goto end;
    }

    while (1)
    {
        memset(query, 0, sizeof(query));

        ret = modbus_receive(ctx, (unsigned char *) query);
        if (ret >= 0)
        {
            pthread_mutex_lock(&rs485_slave_thread->regs_lock);
            modbus_reply(ctx, (const unsigned char *) query, ret, map);
            pthread_mutex_unlock(&rs485_slave_thread->regs_lock);
        }
        else
        {
            dy_syslog(LOG_ERR, "Connection close\n\n");
        }
    }

end:
    if (ctx != NULL)
    {
        modbus_free(ctx);
        rs485_slave_thread->ctx = NULL;
    }
    if (map != NULL)
    {
        modbus_mapping_free(map);
        rs485_slave_thread->map = NULL;
    }

    return NULL;
}

int rs485_slave_start_thread(rs485_slave_var_t *var, rs485_slave_port_cfg_t *rs485_slave_port_cfg, int rs485_slave_port)
{
    pthread_t thread_rs485_slave;
    rs485_slave_thread_t *rs485_slave_thread = NULL;

    dy_syslog(LOG_DEBUG, "=========rs485_slave_port:%d=======", rs485_slave_port);

    rs485_slave_thread = calloc(1, sizeof(rs485_slave_thread_t));
    if (rs485_slave_thread == NULL)
    {
        dy_syslog(LOG_ERR, "Malloc failed");
        return -1;
    }

    pthread_mutex_init(&rs485_slave_thread->regs_lock, NULL);
    rs485_slave_port_cfg->dev = get_tty_from_port(rs485_slave_port_cfg->port);
    rs485_slave_thread->var = var;
    rs485_slave_thread->rs485_slave_port_cfg = rs485_slave_port_cfg;
    rs485_slave_thread->rs485_slave_port = rs485_slave_port;

    rs485_slave_modbus_rtu_init(rs485_slave_thread);
    rs485_slave_mqtt_client_init(rs485_slave_thread);

    pthread_create(&thread_rs485_slave, NULL, rs485_slave_thread_loop, (void *)rs485_slave_thread);

    return 0;
}
