#include "pp.h"
#include "signal.h"
#include "dy_utils/dy_pp.h"
#include "dy_utils/hash_intptr.h"
#include <poll.h>

lua_State *pre_load_scrite_lua(char *file_name)
{
    char file[256] = {0};
    lua_State *lua;
    int ret = 0;

    lua = luaL_newstate();
    luaL_openlibs(lua);

    if (luaL_dostring(lua, LUA_DDEFAULT_ENV))
    {
        dbg_syslog(LOG_ERR, "luaL_loadstring %s error!!!", LUA_DDEFAULT_ENV);
    }

    sprintf(file, "%s%s", PARSE_SCRIPT_DIR, file_name);
    ret = luaL_dofile(lua, file);
    if (ret)
    {
        dbg_syslog(LOG_ERR, "loading error msg:%s", lua_tostring(lua, -1));
        lua_pop(lua, 1);
    }

    return lua;
}

int pre_load_script(template_cfg_t *template)
{
    char *file_name = template->parser_script;
    if (file_name == NULL)
    {
        return 0;
    }

    if (strstr(file_name, LUA_EXTEND))
    {
        template->script.lua = pre_load_scrite_lua(file_name);
    }
    else if (strstr(file_name, PYTHON_EXTEND))
    {
    }
    else if (strstr(file_name, JAVASCRIPT_EXTEND))
    {
    }
    else
    {
        dbg_syslog(LOG_WARNING, "script file_name:%s error !!!", file_name);
    }

    return 0;
}

static void javascript_connect_start(pp_var_t *var)
{
    if (var->javascript.fd == 0)
    {
        int fd, connected;

        var->javascript.serverip[0] = 127;
        var->javascript.serverip[1] = 0;
        var->javascript.serverip[2] = 0;
        var->javascript.serverip[3] = 1;
        var->javascript.port = 7;
        var->javascript.timeout = 60;

        fd = socket(AF_INET, SOCK_STREAM, IPPROTO_TCP);
        if (fd < 0)
        {
            dbg_syslog(LOG_WARNING, "create socket fail ret %d, errno %d\n", fd, errno);
            return;
        }
        fcntl(fd, F_SETFL, fcntl(fd, F_GETFL) | O_NONBLOCK);

        dbg_syslog(LOG_INFO, "start connect server %u.%u.%u.%u\n",
                  var->javascript.serverip[0], var->javascript.serverip[1], var->javascript.serverip[2], var->javascript.serverip[3]);

        connected = tcp_connect_host(fd, var->javascript.serverip, var->javascript.port);
        if (connected == 0)
        {
            dbg_syslog(LOG_INFO, "javascript fd %d connect fail !!!", var->javascript.fd);
            if (var->javascript.connect_timer == 0)
            {
                var->javascript.connect_timer = my_timer_create();
            }
            if (var->javascript.connect_timer > 0)
            {
                my_timer_set(var->javascript.connect_timer, 1, 1 * 1000);
            }
        }
        else
        {
            dbg_syslog(LOG_INFO, "javascript fd %d connect successful", fd);
            if (var->javascript.connect_timer > 0)
            {
                my_timer_delete(var->javascript.connect_timer);
                var->javascript.connect_timer = 0;
            }
        }

        var->javascript.fd = fd;
        var->javascript.connected = connected;
    }
}

int javascript_connect_timer(pp_var_t *var)
{
    TIMER_CONFIRM(var->javascript.connect_timer);

    if (var->javascript.connected)
    {
        dbg_syslog(LOG_INFO, "javascript fd %d connect successful", var->javascript.fd);
        var->javascript.retry_cnt = 0;
        if (var->javascript.connect_timer > 0)
        {
            my_timer_delete(var->javascript.connect_timer);
            var->javascript.connect_timer = 0;
        }
    }
    else
    {
        if (var->javascript.fd)
        {
            var->javascript.connected = tcp_connect_host(var->javascript.fd, var->javascript.serverip, var->javascript.port);
        }
        if (!var->javascript.retry_cnt)
        {
            dbg_syslog(LOG_INFO, "javascript fd %d connect fail !!!", var->javascript.fd);
        }

        var->javascript.retry_cnt++;
        if (var->javascript.retry_cnt > var->javascript.timeout)
        {
            if (var->javascript.fd)
            {
                close(var->javascript.fd);
                var->javascript.fd = 0;
            }
            javascript_connect_start(var);
            var->javascript.retry_cnt = 0;
        }
    }
    return 0;
}

static int pp_poll(pp_var_t *var)
{
    TIMER_CONFIRM(var->poll_timer);
    int i = 0;
    time_t now = time(NULL);

    for (i = 0; i < var->deamon_script_cnt; i++)
    {
        if (var->feed_info[i].last_feed == 0)
        {
            continue;
        }

        if ((now - var->feed_info[i].last_feed) > 10)
        {
            dbg_syslog(LOG_ERR, "template %s watchdog timeout, reboot the PP", var->feed_info[i].template_id);
            system(PP_RESTART_CMD);
        }
    }

    return 0;
}

static void pp_loop(pp_var_t *var)
{
    int ret = -1;
    nfds_t nfds;
    struct pollfd pfds[2];

    while (1)
    {
        nfds = 0;
        if (var->javascript.connect_timer > 0) {
            pfds[nfds].fd = var->javascript.connect_timer;
            pfds[nfds].events = POLLIN;
            pfds[nfds].revents = 0;
            nfds++;
        }

        pfds[nfds].fd = var->poll_timer;
        pfds[nfds].events = POLLIN;
        pfds[nfds].revents = 0;
        nfds++;
        ret = poll(pfds, nfds, 10 * 1000);
        if (ret < 0)
        {
            int error = errno;
            printf("errno %d\n", error);
            if (error == EINTR)
            {
                continue;
            }
            else
            {
                break;
            }
        }
        else if (ret > 0)
        {
            if (nfds == 2 && pfds[0].revents)
            {
                javascript_connect_timer(var);
            }
            if (var->poll_timer > 0 && pfds[nfds - 1].revents)
            {
                pp_poll(var);
            }
        }
    }
}

static void pp_subscribe_all(pp_var_t *var)
{
    ipc_session_t *ipc_session = var->session;
    char topic[TOPIC_MAX_LEN] = {0};

    snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/+/device/+/data/%s", var->sn_str,TOPIC_UP_RAW_DATA);
    ipc_session_subscribe(ipc_session, topic);

    snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/+/device/+/data/%s", var->sn_str,TOPIC_EVT_SET_RGLT);
    ipc_session_subscribe(ipc_session, topic);

    // ipc/[GW SN]/script/[template id]/heartbeat
    snprintf(topic, TOPIC_MAX_LEN, "ipc/%s/%s/+/%s",var->sn_str, TOPIC_SCRIPT, TOPIC_SCRIPT_HB);
    ipc_session_subscribe(ipc_session, topic);

    snprintf(topic, TOPIC_MAX_LEN, "M/%s/+/device/+/service/+",var->sn_str);
    ipc_session_subscribe(ipc_session, topic);
}

static int pp_handle_mqtt_data(pp_var_t *var, ipc_msg_t *msg)
{
    char port[TOPIC_MAX_LEN] = {0};

    sscanf(msg->topic, "ipc/%*[^/]/%[^/]", port);

    if (port[0] == '\0') {
        sscanf(msg->topic, "M/%*[^/]/%[^/]/device/%*[^/]/service/%*[^/]", port);
    }

    if (strcmp(port, TOPIC_SCRIPT) == 0)
    {
        char template_id[16] = {0};
        // heartbeat for script
        // ipc/[GW SN]/script/[template id]/heartbeat
        sscanf(msg->topic, "ipc/%*[^/]/%*[^/]/%[^/]", template_id);
        if (template_id[0])
        {
            int i = 0;

            for (i = 0; i < var->deamon_script_cnt; i++)
            {
                if (strcmp(var->feed_info[i].template_id, template_id) == 0)
                {
                    time_t now = time(NULL);
                    var->feed_info[i].last_feed = now;
                    break;
                }
            }
        }
    }
    else if (strstr(msg->topic, TOPIC_UP_RAW_DATA))
    {
        return pp_parse_realtime_data(var, msg);
    }
    else if ((msg->topic[0] == 'M' && msg->topic[1] == '/' && port[0] && strcmp(port, TOPIC_GATEWAY_CONFIG)) || strstr(msg->topic, TOPIC_EVT_SET_RGLT))
    {
        return msg_set_regulate_signal_val(var, msg);
    }

    return 0;
}

static int pp_mqtt_handle_recv_msg(void *obj, ipc_msg_t *msg)
{
    pp_var_t *var = (pp_var_t *)obj;
    pp_handle_mqtt_data(var, msg);
    return 0;
}
// 建立与内部broker之间的MQTT连接
static int pp_mqtt_client_init(pp_var_t *var)
{
    char clientId[MAX_CLIENT_ID_LEN] = {0};

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

    ipc_session_set_callbacks(var->session, pp_mqtt_handle_recv_msg, NULL);
    pp_subscribe_all(var);
    ipc_session_start(var->session);
    return 0;
}

static void sig_handler(int sig)
{
    syslog(LOG_USER | LOG_NOTICE, "%s:%d receive signal:%d", __FILE__, __LINE__, sig);
    exit(0);
}

static void * build_upvalue_indices(lua_State * L)
{
    int ntop, idx, ret, ntop1;
    struct hash_intptr * hptr;
    const char * func = LNX_UPVALUES_FUNC;

    hptr = NULL;

    /* determine the top of stack */
    ntop = lua_gettop(L);

    /* fetch function `func at the top of stack */
    lua_getglobal(L, func);

    ret = lua_gettop(L);
    /* check again for the top of stack */
    if (ret != (ntop + 1) || lua_type(L, -1) != LUA_TFUNCTION) {
        dbg_syslog(LOG_ERR, "Fatal error, failed to find function '%s', type: %d",
            func, lua_type(L, -1));
        return NULL;
    }

    idx = 0;
    ntop1 = lua_gettop(L);
    for (;;) {
        size_t slen;
        const char * signm;

        idx++;
        lua_settop(L, ntop1);
        /* get the upvalue at index `idx */
        signm = lua_getupvalue(L, ntop + 1, idx);
        if (signm == NULL)
            break;
        slen = strlen(signm);
        if (slen == 0) {
            dbg_syslog(LOG_ERR, "Fatal error, invalid zero length of upvalue at %d", idx);
            continue;
        }

        ret = hash_intptr_addint(&hptr, signm, (unsigned int) slen, idx);
        if (ret < 0) {
            dbg_syslog(LOG_ERR, "Fatal error, failed to push upvalue: %s", signm);
            hash_intptr_removeall(&hptr);
            hptr = NULL;
            break;
        }
        /* fprintf(stdout, "\tupvalue[%d]: %s\n", idx, signm);
        fflush(stdout); */
    }

    lua_settop(L, ntop);
    return (void *) hptr;
}

static void build_all_luaexp(const template_table_t * ptemp, object_table_t * objtab)
{
    int idx, onum, ret;
    object_cfg_t * ocfg;

    ret = 0;
    onum = objtab ? objtab->object_cnt : 0;
    if (onum <= 0)
        return;

    /* process object table one by one */
    for (idx = 0; idx < onum; ++idx) {
        unsigned int snum, jdx;
        object_service_cfg_t * scfg;
        object_service_table_t * stab;

        ocfg = &(objtab->object[idx]);
        stab = ocfg->service_tab;
        snum = stab ? stab->serviceCnt : 0;
        if (snum == 0)
            continue;

        /* process services one by one */
        for (jdx = 0; jdx < snum; ++jdx) {
            lua_State * news;
            object_property_table_t * pout;

            scfg = &(stab->service[jdx]);
            pout = scfg->poutput;
            /* destroy old, existing Lua state machine */
            if (pout->luast != NULL) {
                lua_close(pout->luast);
                pout->luast = NULL;
            }

            /* construct the Lua script */
            ret = build_property_luaexp(ptemp, scfg, ocfg->object_model_id);
            if (ret <= 0) {
                /* no script constructed, next! */
                continue;
            }

            news = luaL_newstate();
            if (news) {
                /* load standard lua libraries */
                luaL_openlibs(news);
            }
            if (news == NULL) {
                fputs("Error, system out of memory!\n", stderr);
                fflush(stderr);
                continue;
            }

            /* load extra, external Lua modules */
            ret = luaL_dostring(news, LUA_DDEFAULT_ENV);
            if (ret != 0) {
                lua_close(news);
                dbg_syslog(LOG_ERR, "Fatal lua error, failed to load DEFAULT_ENV: %d", ret);
                continue;
            }

            /* load constructed lua script now */
            ret = (pout->lua_script != NULL) ? luaL_dostring(news, pout->lua_script) : -1;
            if (ret != 0) {
                dbg_syslog(LOG_ERR, "Fatal lua error, failed to load script: %d", ret);
                fprintf(stderr, "Error, invalid constructed lua script:\n%s\n",
                    pout->lua_script ? : "nil");
                fflush(stderr);
                lua_close(news);
                continue;
            }

            pout->luast = news;
            if (pout->sign_mark != NULL) {
                struct hash_intptr * hptr;
                hptr = (struct hash_intptr *) pout->sign_mark;
                hash_intptr_removeall(&hptr);
                pout->sign_mark = NULL;
            }
            /* so far, so good, time to construct the upvalue index */
            pout->sign_mark = build_upvalue_indices(news);
        }
    }
}

static int pp_init(pp_var_t *var)
{
    int i;
    struct sigaction sig_action;
    sig_action.sa_handler = sig_handler;
    sig_action.sa_flags = SA_RESTART;
    sigemptyset(&sig_action.sa_mask);
    sigaction(SIGINT, &sig_action, NULL);
    sigaction(SIGTERM, &sig_action, NULL);
    sigaction(SIGABRT, &sig_action, NULL);
    get_board_sn(var->sn_str);
    get_process_name(var->proc_name);
    if (load_nodes_cfg(&var->nodes_cfg_table, NODES_CFG_PATH) == -1)
    {
        dbg_syslog(LOG_ERR, "load nodes cfg fail");
    }
    if (load_templates_cfg(&var->template_table, TEMPLATES_CFG_PATH) == -1)
    {
        dbg_syslog(LOG_ERR, "load template cfg fail");
    }
    if (load_objects_cfg(&var->object_table, OBJECTS_CFG_PATH) == -1)
    {
        dbg_syslog(LOG_ERR, "load object cfg fail");
    }
    var->deamon_script_cnt = 0;
    if (var->template_table)
    {
        build_all_luaexp(var->template_table, var->object_table);
        for (i = 0; i < var->template_table->template_cnt; i++)
        {
            if (var->template_table->template[i].use_parser_url)
            {
                pp_start_script_thread(var, &var->template_table->template[i]);
            }
        }
    }

    if (var->deamon_script_cnt > 0)
    {
        int j = 0;

        var->feed_info = calloc(var->deamon_script_cnt, sizeof(feed_info_t));
        for (i = 0; i < var->template_table->template_cnt; i++)
        {
            if (var->template_table->template[i].use_parser_url == SCRIPT_DEAMON)
            {
                var->feed_info[j].template_id = var->template_table->template[i].template_id;
                var->feed_info[j].last_feed = 0;
                j++;
            }
        }
    }

    pp_mqtt_client_init(var);

    if (var->lua == NULL)
    {
        var->lua = luaL_newstate();
        luaL_openlibs(var->lua);

        if (luaL_dostring(var->lua, LUA_DDEFAULT_ENV))
        {
            dbg_syslog(LOG_ERR, "luaL_loadstring %s error!!!", LUA_DDEFAULT_ENV);
        }
    }

    if (var->python == 0)
    {
        //初始化，载入python的扩展模块
        Py_Initialize();
        //判断初始化是否成功
        if (!Py_IsInitialized())
        {
            dbg_syslog(LOG_WARNING, "Python init failed!\n");
        }
        //导入当前路径
        PyRun_SimpleString("import sys");
        PyRun_SimpleString("sys.path.append('./')");
        PyRun_SimpleString("sys.path.append('/app/script/')");

        var->python = 1;
    }

    var->poll_timer = my_timer_create();
    if (var->poll_timer > 0)
    {
        my_timer_set(var->poll_timer, 1, 5000);
    }

    // javascript_connect_start(var);

    dbg_syslog(LOG_INFO, "init done");

    return 0;
}

int main(int argc, char *argv[])
{
    pp_var_t var = {0};

    memset(&var, 0, sizeof(var));
    lnxall_loglevel_set(LNXALL_LOGNOTICE, 1);

    pp_init(&var);
    pp_loop(&var);

    return 0;
}
