/*
 * Created by jiaqiang.ye@lnxall.com
 *
 * MHMP protocol implmentation
 *
 * 2022/03/30
 */

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

#include "mhmp.h"
#include <dy_utils/cJSON.h>

static int json_find_number(cJSON * json, const char * item,
    int * valint, double * valdb)
{
    cJSON * jnum;

    jnum = cJSON_GetObjectItem(json, item);
    if (jnum == NULL)
        return -1;
    if (cJSON_IsNumber(jnum)) {
        if (valint != NULL)
            *valint = jnum->valueint;
        if (valdb != NULL)
            *valdb = jnum->valuedouble;
        return 0;
    }

    if (cJSON_IsString(jnum) &&
        jnum->valuestring != NULL) {
        long vallong;
        double valdouble;
        char * jstr, * strp;

        jstr = strp = jnum->valuestring;
        if (strchr(jstr, '.') != NULL ||
            strchr(jstr, 'E') != NULL ||
            strchr(jstr, 'e') != NULL) {
            valdouble = strtod(jstr, &strp);
            vallong = (long) valdouble;
        } else {
            vallong = strtol(jstr, &strp, 10);
            valdouble = (double) vallong;
        }
        if (strp && strp == jstr)
            return -2;

        if (valint != NULL)
            *valint = (int) vallong;
        if (valdb != NULL)
            *valdb = valdouble;
        return 0;
    }
    return -3;
}

static int json_find_string(cJSON * json, const char * item,
    char * output, unsigned int outlen)
{
    cJSON * jstr;
    const char * strp;
    unsigned int jlen = 0;

    jstr = cJSON_GetObjectItem(json, item);
    if (jstr == NULL || cJSON_IsString(jstr) == 0)
        return -1;

    strp = jstr->valuestring;
    if (strp != NULL)
        jlen = (unsigned int) strlen(strp);
    if (jlen == 0) {
        output[0] = '\0';
        return 0;
    }
    if (jlen >= outlen) {
        outlen--;
        memcpy(output, strp, (size_t) outlen);
        output[outlen] = '\0';
        return (int) outlen;
    }
    strncpy(output, strp, (size_t) outlen);
    return (int) jlen;
}

static int mhmp_check_reply(cJSON * payload)
{
    int ret;
    int replyFlag;
    char replyReason[64];

    replyFlag = -1;
    ret = json_find_number(payload, "ReplyFlag", &replyFlag, NULL);
    if (ret < 0) {
        fputs("Error, ReplyFlag not found.\n", stderr);
        fflush(stderr);
        return -1;
    }
    if (replyFlag == 0)
        return 0;

    ret = json_find_string(payload,
        "ReplyReason", replyReason, sizeof(replyReason));
    if (ret < 0) {
        fputs("ERror, ReplyReason not found.\n", stderr);
        fflush(stderr);
        return -2;
    }
    fprintf(stderr, "Error, data unavailable: %d, %s\n",
        replyFlag, replyReason);
    fflush(stderr);
    return -3;
}

static int mhmp_ipc_msg(const char * psn,
    const char * msg, time_t nowt, char * * msgptr)
{
    cJSON * root, * tmp;
    char * retval;

    root = cJSON_CreateObject();
    if (root == NULL)
        return -1;

    tmp = cJSON_CreateString(psn);
    if (tmp != NULL)
        cJSON_AddItemToObject(root, "sn", tmp);

    tmp = cJSON_CreateString(msg);
    if (tmp != NULL)
        cJSON_AddItemToObject(root, "tag_node", tmp);

    tmp = cJSON_CreateNumber(0);
    if (tmp != NULL)
        cJSON_AddItemToObject(root, "mi", tmp);

    tmp = cJSON_CreateNumber((double) nowt);
    if (tmp != NULL)
        cJSON_AddItemToObject(root, "time", tmp);

    retval = cJSON_Print(root);
    cJSON_Delete(root);
    *msgptr = retval;
    return retval ? (int) strlen(retval) : -1;
}

static int mhmp_add_xaf_xvf(const struct egwIP_vibration * pinfo,
    char * jbuf, int maxlen, int curlen)
{
    int ret;
    if (maxlen <= (curlen + 1024))
        return curlen;

    if (pinfo->ev_xaf) {
        ret = snprintf(&jbuf[curlen], (size_t) (maxlen - curlen),
            ",\"Xaf\":\"%s\"", pinfo->ev_xaf);
        curlen += ret;
        if (maxlen <= (curlen + 1024))
            return curlen;
    }

    if (pinfo->ev_xvf) {
        ret = snprintf(&jbuf[curlen], (size_t) (maxlen - curlen),
            ",\"Xvf\":\"%s\"", pinfo->ev_xvf);
        curlen += ret;
        if (maxlen <= (curlen + 1024))
            return curlen;
    }

    if (pinfo->ev_yaf) {
        ret = snprintf(&jbuf[curlen], (size_t) (maxlen - curlen),
            ",\"Yaf\":\"%s\"", pinfo->ev_yaf);
        curlen += ret;
        if (maxlen <= (curlen + 1024))
            return curlen;
    }

    if (pinfo->ev_yvf) {
        ret = snprintf(&jbuf[curlen], (size_t) (maxlen - curlen),
            ",\"Yvf\":\"%s\"", pinfo->ev_yvf);
        curlen += ret;
        if (maxlen <= (curlen + 1024))
            return curlen;
    }

    if (pinfo->ev_zaf) {
        ret = snprintf(&jbuf[curlen], (size_t) (maxlen - curlen),
            ",\"Zaf\":\"%s\"", pinfo->ev_zaf);
        curlen += ret;
        if (maxlen <= (curlen + 1024))
            return curlen;
    }

    if (pinfo->ev_zvf) {
        ret = snprintf(&jbuf[curlen], (size_t) (maxlen - curlen),
            ",\"Zvf\":\"%s\"", pinfo->ev_zvf);
        curlen += ret;
    }

    return curlen;
}

#define MHMP_REALTIME_BUFSIZ 0x100000 /* 1MB */
#define MHMP_TOPIC_BUFSIZ    256
static void mhmp_publish_data(
    struct mhmp_var * mvar, const struct egwIP_info * pinfo, time_t tval)
{
    char newsn[128];
    int plen, idx, ret;
    ipc_session_t * sess;
    char * topic = NULL;
    char * payload = NULL;
    char * realpay = NULL;

    sess = mvar->ipcsess;
    payload = (char *) malloc(MHMP_REALTIME_BUFSIZ);
    if (payload == NULL) {
err0:
        fputs("Error, system out of memory!\n", stderr);
        fflush(stderr);
        return;
    }
    topic = (char *) malloc(MHMP_TOPIC_BUFSIZ);
    if (topic == NULL) {
        free(payload);
        goto err0;
    }

    plen = snprintf(payload, MHMP_REALTIME_BUFSIZ,
        "{"
        "\"DeviceID\":\"%s\","
        "\"DeviceName\":\"%s\","
        "\"EgwIP\":\"%s\","
        "\"startTime\":\"%s\","
        "\"Temperature\":%.01lf,"
        "\"X\":%.04lf,"
        "\"Y\":%.04lf,"
        "\"Z\":%.04lf,"
        "\"Xa\":%.04lf,"
        "\"Xd\":%.04lf,"
        "\"Ya\":%.04lf,"
        "\"Yd\":%.04lf,"
        "\"Za\":%.04lf,"
        "\"Zd\":%.04lf",
        pinfo->deviceID, pinfo->deviceName, pinfo->egwip_id,
        pinfo->startTime, pinfo->temperature,
        pinfo->vibration.ev_x,
        pinfo->vibration.ev_y,
        pinfo->vibration.ev_z,
        pinfo->vibration.ev_xa,
        pinfo->vibration.ev_xd,
        pinfo->vibration.ev_ya,
        pinfo->vibration.ev_yd,
        pinfo->vibration.ev_za,
        pinfo->vibration.ev_zd);
    plen = mhmp_add_xaf_xvf(&pinfo->vibration, payload, MHMP_REALTIME_BUFSIZ, plen);
    payload[plen++] = '}';
    payload[plen] = '\0';

    idx = snprintf(topic, MHMP_TOPIC_BUFSIZ,
        "ipc/%s/mhmp-realtime/device/%s/data/property/post",
        mvar->board_sn, pinfo->egwip_id);
    if (idx <= 0) {
        fputs("Error, fatal internal error!\n", stderr);
        fflush(stderr);
        goto impossible;
    }

    snprintf(newsn, sizeof(newsn), "%s_%d", pinfo->deviceID, pinfo->index_no);
    idx = mhmp_ipc_msg(newsn, payload, tval, &realpay);
    if (realpay && idx > 0) {
        ret = ipc_session_publish(sess, topic, (unsigned char *) realpay, idx);
        if (ret != 0) {
            fprintf(stderr, "Error, failed to publish message: %d\n", ret);
            fflush(stderr);
        }
    }

impossible:
    if (realpay != NULL)
        free(realpay);
    free(topic);
    free(payload);
}

static void mhmp_publish_kafka(
    struct mhmp_var * mvar, const struct egwIP_info * pinfo, time_t nowt)
{
    int plen, idx, ret;
    ipc_session_t * sess;
    char * topic = NULL;
    char * payload = NULL;

    sess = mvar->ipcsess;
    payload = (char *) malloc(MHMP_REALTIME_BUFSIZ);
    if (payload == NULL) {
err0:
        fputs("Error, system out of memory!\n", stderr);
        fflush(stderr);
        return;
    }
    topic = (char *) malloc(MHMP_TOPIC_BUFSIZ);
    if (topic == NULL) {
        free(payload);
        goto err0;
    }

    plen = snprintf(payload, MHMP_REALTIME_BUFSIZ,
        "{"
        "\"sn\":\"%s_%d\","
        "\"from_mhmp\":true,"
        "\"time\":%ld,"
        "\"tags\":{"
        "\"DeviceID\":\"%s\","
        "\"DeviceName\":\"%s\","
        "\"EgwIP\":\"%s\","
        "\"startTime\":\"%s\","
        "\"Temperature\":%.01lf,"
        "\"X\":%.04lf,"
        "\"Y\":%.04lf,"
        "\"Z\":%.04lf,"
        "\"Xa\":%.04lf,"
        "\"Xd\":%.04lf,"
        "\"Ya\":%.04lf,"
        "\"Yd\":%.04lf,"
        "\"Za\":%.04lf,"
        "\"Zd\":%.04lf",
        pinfo->deviceID, pinfo->index_no, (long) nowt,
        pinfo->deviceID, pinfo->deviceName, pinfo->egwip_id,
        pinfo->startTime, pinfo->temperature,
        pinfo->vibration.ev_x,
        pinfo->vibration.ev_y,
        pinfo->vibration.ev_z,
        pinfo->vibration.ev_xa,
        pinfo->vibration.ev_xd,
        pinfo->vibration.ev_ya,
        pinfo->vibration.ev_yd,
        pinfo->vibration.ev_za,
        pinfo->vibration.ev_zd);
    plen = mhmp_add_xaf_xvf(&pinfo->vibration, payload, MHMP_REALTIME_BUFSIZ, plen);
    payload[plen++] = '}';
    payload[plen++] = '}';
    payload[plen] = '\0';

    idx = snprintf(topic, MHMP_TOPIC_BUFSIZ,
        "ipc/%s/mhmp-realtime/device/%s/data_filtered/service/post",
        mvar->board_sn, pinfo->egwip_id);
    if (idx <= 0) {
        fputs("Error, fatal internal error!\n", stderr);
        fflush(stderr);
        goto impossible;
    }

    ret = ipc_session_publish(sess, topic, (unsigned char *) payload, plen);
    if (ret != 0) {
        fprintf(stderr, "Error, failed to publish message: %d\n", ret);
        fflush(stderr);
    }

impossible:
    free(topic);
    free(payload);
}

static void mhmp_publish_warn_info(struct mhmp_var * mvar,
    const struct egwIP_warninfo * winfo, time_t now)
{
    char newsn[128];
    int plen, ret;
    ipc_session_t * sess;
    char * topic = NULL;
    char * payload = NULL;
    char * realpay = NULL;
    const char * empty = "";

    sess = mvar->ipcsess;
    payload = (char *) malloc(MHMP_REALTIME_BUFSIZ);
    if (payload == NULL) {
err0:
        fputs("Error, system out of memory!\n", stderr);
        fflush(stderr);
        return;
    }
    topic = (char *) malloc(MHMP_TOPIC_BUFSIZ);
    if (topic == NULL) {
        free(payload);
        goto err0;
    }

    plen = snprintf(payload, MHMP_REALTIME_BUFSIZ,
        "{"
        "\"DeviceID\":\"%s\","
        "\"DeviceName\":\"%s\","
        "\"EgwIP\":\"%s\","
        "\"WarnID\":\"%s\","
        "\"WarnType\":\"%s\","
        "\"PushType\":\"%s\","
        "\"WarnValue\":\"%s\",\n"
        "\"MonitoringSite\":\"%s\","
        "\"MonitoringSiteName\":\"%s\","
        "\"StartTime\":\"%s\","
        "\"EndTime\":\"%s\","
        "\"DiagCode\":\"%s\"}",
        winfo->deviceID, winfo->deviceName, winfo->egwip_id,
        winfo->warnID ? : empty,
        winfo->warnType ? : empty,
        winfo->pushType ? : empty,
        winfo->warnValue ? : empty,
        winfo->monitoringSite ? : empty,
        winfo->monitoringSiteName ? : empty,
        winfo->startTime ? : empty,
        winfo->endTime ? : empty,
        winfo->diagCode ? : empty);

    snprintf(topic, MHMP_TOPIC_BUFSIZ,
        "ipc/%s/mhmp-warninfo/device/%s/data/property/post",
        mvar->board_sn, winfo->egwip_id);

    snprintf(newsn, sizeof(newsn), "%s_%d", winfo->deviceID, winfo->index_no);
    plen = mhmp_ipc_msg(newsn, payload, now, &realpay);

    if (realpay && plen > 0) {
        ret = ipc_session_publish(sess, topic, (unsigned char *) realpay, plen);
        if (ret != 0) {
            fprintf(stderr, "Error, failed to publish message: %d\n", ret);
            fflush(stderr);
        }
    }

    free(topic);
    free(payload);
    if (realpay != NULL)
        free(realpay);
}

static void mhmp_publish_warninfo_kafka(struct mhmp_var * mvar,
    const struct egwIP_warninfo * winfo, time_t now, const char * which)
{
    int plen, ret;
    ipc_session_t * sess;
    char * topic = NULL;
    char * payload = NULL;
    const char * empty = "";

    sess = mvar->ipcsess;
    payload = (char *) malloc(MHMP_REALTIME_BUFSIZ);
    if (payload == NULL) {
err0:
        fputs("Error, system out of memory!\n", stderr);
        fflush(stderr);
        return;
    }
    topic = (char *) malloc(MHMP_TOPIC_BUFSIZ);
    if (topic == NULL) {
        free(payload);
        goto err0;
    }

    plen = snprintf(payload, MHMP_REALTIME_BUFSIZ,
        "{"
        "\"sn\":\"%s_%d\","
        "\"from_mhmp\":true,"
        "\"time\":%ld,"
        "\"%s\":{"
        "\"DeviceID\":\"%s\","
        "\"DeviceName\":\"%s\","
        "\"EgwIP\":\"%s\","
        "\"WarnID\":\"%s\","
        "\"WarnType\":\"%s\","
        "\"PushType\":\"%s\","
        "\"WarnValue\":\"%s\",\n"
        "\"MonitoringSite\":\"%s\","
        "\"MonitoringSiteName\":\"%s\","
        "\"StartTime\":\"%s\","
        "\"EndTime\":\"%s\","
        "\"DiagCode\":\"%s\"}}",
        winfo->deviceID, winfo->index_no, (long) now, which,
        winfo->deviceID, winfo->deviceName, winfo->egwip_id,
        winfo->warnID ? : empty,
        winfo->warnType ? : empty,
        winfo->pushType ? : empty,
        winfo->warnValue ? : empty,
        winfo->monitoringSite ? : empty,
        winfo->monitoringSiteName ? : empty,
        winfo->startTime ? : empty,
        winfo->endTime ? : empty,
        winfo->diagCode ? : empty);

    snprintf(topic, MHMP_TOPIC_BUFSIZ,
        "ipc/%s/mhmp-warninfo/device/%s/data_filtered/service/post",
        mvar->board_sn, winfo->egwip_id);

    ret = ipc_session_publish(sess, topic, (unsigned char *) payload, plen);
    if (ret != 0) {
        fprintf(stderr, "Error, failed to publish message: %d\n", ret);
        fflush(stderr);
    }

    free(topic);
    free(payload);
}

static void mhmp_handle_realtime_datel(
    struct mhmp_var * mvar, cJSON * datel)
{
    time_t nowt;
    cJSON * tlist;
    cJSON * viblist;
    int idx, count, ret;
    struct egwIP_info * info = NULL;
    double temps[EGWIP_TEMP_NUM];

    info = (struct egwIP_info *) calloc(0x1, sizeof(*info));
    if (info == NULL) {
        fputs("Error, system out of memory!\n", stderr);
        fflush(stderr);
        return;
    }

    ret = json_find_string(datel, "deviceID", info->deviceID, EGWIP_NAME_MAX);
    if (ret > 0)
        ret = json_find_string(datel, "deviceName", info->deviceName, EGWIP_NAME_MAX);
    if (ret > 0)
        ret = json_find_string(datel, "egwIP", info->egwip_id, EGWIP_NAME_MAX);
    if (ret > 0)
        ret = json_find_string(datel, "startTime", info->startTime, EGWIP_TIME_LEN);
    if (ret <= 0) {
        free(info);
        fputs("Error, invalid deviceID/deviceName/egwIP found\n", stderr);
        fflush(stderr);
        return;
    }

    nowt = time(NULL);
    /* mhmp_get_time(info->timestamp, EGWIP_TIME_LEN, &nowt, 0); */
    for (idx = 0; idx < EGWIP_TEMP_NUM; ++idx)
        temps[idx] = -1;

    tlist = cJSON_GetObjectItem(datel, "temperatureList");
    if (tlist && cJSON_IsArray(tlist)) {
        count = cJSON_GetArraySize(tlist);
        if (count > EGWIP_TEMP_NUM)
            count = EGWIP_TEMP_NUM;
        for (idx = 0; idx < count; ++idx) {
            cJSON * temp;
            temp = cJSON_GetArrayItem(tlist, idx);
            if (temp != NULL) {
                if (cJSON_IsString(temp))
                    temps[idx] = strtod(temp->valuestring, NULL);
                else
                    temps[idx] = temp->valuedouble;
            }
        }
    }

    viblist = cJSON_GetObjectItem(datel, "vibrationList");
    if (viblist && cJSON_IsArray(viblist)) {
        count = cJSON_GetArraySize(viblist);
        if (count > EGWIP_TEMP_NUM)
            count = EGWIP_TEMP_NUM;
        for (idx = 0; idx < count; ++idx) {
            cJSON * temp;
            temp = cJSON_GetArrayItem(viblist, idx);
            if (temp && cJSON_IsObject(temp)) {
				cJSON * avf;
                struct egwIP_vibration * vibp;
                info->temperature = temps[idx];
                info->index_no = idx;
                vibp = &info->vibration;
                if (json_find_number(temp, "x", NULL, &vibp->ev_x) < 0)
                    continue;
                if (json_find_number(temp, "y", NULL, &vibp->ev_y) < 0)
                    continue;
                if (json_find_number(temp, "z", NULL, &vibp->ev_z) < 0)
                    continue;

                if (json_find_number(temp, "xa", NULL, &vibp->ev_xa) < 0)
                    continue;
                if (json_find_number(temp, "xd", NULL, &vibp->ev_xd) < 0)
                    continue;

                if (json_find_number(temp, "ya", NULL, &vibp->ev_ya) < 0)
                    continue;
                if (json_find_number(temp, "yd", NULL, &vibp->ev_yd) < 0)
                    continue;

                if (json_find_number(temp, "za", NULL, &vibp->ev_za) < 0)
                    continue;
                if (json_find_number(temp, "zd", NULL, &vibp->ev_zd) < 0)
                    continue;

                avf = cJSON_GetObjectItem(temp, "xaf");
                if (cJSON_IsString(avf))
                    vibp->ev_xaf = avf->valuestring;
                else
                    vibp->ev_xaf = NULL;
                avf = cJSON_GetObjectItem(temp, "xvf");
                if (cJSON_IsString(avf))
                    vibp->ev_xvf = avf->valuestring;
                else
                    vibp->ev_xvf = NULL;

                avf = cJSON_GetObjectItem(temp, "yaf");
                if (cJSON_IsString(avf))
                    vibp->ev_yaf = avf->valuestring;
                else
                    vibp->ev_yaf = NULL;
                avf = cJSON_GetObjectItem(temp, "yvf");
                if (cJSON_IsString(avf))
                    vibp->ev_yvf = avf->valuestring;
                else
                    vibp->ev_yvf = NULL;

                avf = cJSON_GetObjectItem(temp, "zaf");
                if (cJSON_IsString(avf))
                    vibp->ev_zaf = avf->valuestring;
                else
                    vibp->ev_zaf = NULL;
                avf = cJSON_GetObjectItem(temp, "zvf");
                if (cJSON_IsString(avf))
                    vibp->ev_zvf = avf->valuestring;
                else
                    vibp->ev_zvf = NULL;

                mhmp_publish_data(mvar, info, nowt);
                mhmp_publish_kafka(mvar, info, nowt);
            }
        }
    }

    free(info);
}

static void mhmp_handle_realtime_data(
    struct mhmp_var * mvar, cJSON * payload)
{
    int datnum, idx;
    cJSON * jdat = NULL;

    if (mhmp_check_reply(payload) < 0)
        return;
    jdat = cJSON_GetObjectItem(payload, "data");
    if (jdat == NULL || cJSON_IsArray(jdat) == 0) {
        fputs("Error, no data array found\n", stderr);
        fflush(stderr);
        return;
    }
    datnum = cJSON_GetArraySize(jdat);
    if (datnum <= 0)
        return;

    for (idx = 0; idx < datnum; ++idx) {
        cJSON * datel;
        datel = cJSON_GetArrayItem(jdat, idx);
        if (datel == NULL || cJSON_IsObject(datel) == 0)
            continue;
        mhmp_handle_realtime_datel(mvar, datel);
    }
}

static int mhmp_get_timetag_value(const cJSON * jdat,
    int idx, char * timetag, size_t taglen, double * valuep)
{
    double value;
    const cJSON * jel;

    jel = cJSON_GetArrayItem(jdat, idx);
    if (jel == NULL || cJSON_IsObject(jel) == 0)
        return -1;
    jel = jel->child;
    if (cJSON_IsNumber(jel))
        value = jel->valuedouble;
    else if (cJSON_IsString(jel) && jel->valuestring != NULL) {
        if (strchr(jel->valuestring, '.') != NULL)
            value = strtod(jel->valuestring, NULL);
        else
            value = (double) strtol(jel->valuestring, NULL, 0);
    } else
        return -2;
    if (timetag != NULL) {
        if (jel->string == NULL || jel->string[0] == '\0')
            return -3;
        strncpy(timetag, jel->string, taglen);
    }
    *valuep = value;
    return 0;
}

static void mhmp_handle_history_data(
    struct mhmp_var * mvar, cJSON * payload)
{
    time_t nowt;
    cJSON * viblist;
    int datnum, dix;
    int vibnum, idx, ret;
    struct egwIP_info * info;

    info = NULL;
    if (mhmp_check_reply(payload) < 0)
        return;
    payload = cJSON_GetObjectItem(payload, "data");
    if (payload == NULL || cJSON_IsArray(payload) == 0) {
        fputs("Error, data not found in history data\n", stderr);
        fflush(stderr);
        return;
    }

    datnum = cJSON_GetArraySize(payload);
    if (datnum <= 0) {
        fprintf(stderr, "Error, invalid number of history data: %d\n", datnum);
        fflush(stderr);
        return;
    }

    info = (struct egwIP_info *) malloc(sizeof(*info));
    if (info == NULL) {
        fputs("Error, system out of memory!\n", stderr);
        fflush(stderr);
        return;
    }

    nowt = 0;
    time(&nowt);

    for (dix = 0; dix < datnum; ++dix) {
        cJSON * jdat;

        jdat = cJSON_GetArrayItem(payload, dix);
        if (jdat == NULL || cJSON_IsObject(jdat) == 0)
            continue;

        viblist = cJSON_GetObjectItem(jdat, "vibrationList");
        if (viblist == NULL || cJSON_IsArray(viblist) == 0)
            continue;

        vibnum = cJSON_GetArraySize(viblist);
        if (vibnum <= 0)
            continue;

        memset(info, 0, sizeof(*info));
        ret = json_find_string(jdat, "deviceID", info->deviceID, EGWIP_NAME_MAX);
        if (ret >= 0)
            ret = json_find_string(jdat, "deviceName", info->deviceName, EGWIP_NAME_MAX);
        if (ret >= 0)
            ret = json_find_string(jdat, "egwIP", info->egwip_id, EGWIP_NAME_MAX);
        if (ret < 0)
            continue;

        for (idx = 0; idx < vibnum; ++idx) {
            cJSON * vib;
            int num, jdx;
            cJSON * j_xa, * j_xd;
            cJSON * j_ya, * j_yd;
            cJSON * j_za, * j_zd;
            cJSON * j_x, * j_y, * j_z;

            vib = cJSON_GetArrayItem(viblist, idx);
            if (vib == NULL || cJSON_IsObject(vib) == 0)
                continue;

            j_x  = cJSON_GetObjectItem(vib, "x");
            if (j_x == NULL || cJSON_IsArray(j_x) == 0)
                continue;
            num = cJSON_GetArraySize(j_x);
            if (num <= 0)
                continue;

            j_y  = cJSON_GetObjectItem(vib, "y");
            if (j_y == NULL || cJSON_IsArray(j_y) == 0)
                continue;

            j_z  = cJSON_GetObjectItem(vib, "z");
            if (j_z == NULL || cJSON_IsArray(j_z) == 0)
                continue;

            j_xa = cJSON_GetObjectItem(vib, "xa");
            if (j_xa == NULL || cJSON_IsArray(j_xa) == 0)
                continue;

            j_xd = cJSON_GetObjectItem(vib, "xd");
            if (j_xd == NULL || cJSON_IsArray(j_xd) == 0)
                continue;

            j_ya = cJSON_GetObjectItem(vib, "ya");
            if (j_ya == NULL || cJSON_IsArray(j_ya) == 0)
                continue;

            j_yd = cJSON_GetObjectItem(vib, "yd");
            if (j_yd == NULL || cJSON_IsArray(j_yd) == 0)
                continue;

            j_za = cJSON_GetObjectItem(vib, "za");
            if (j_za == NULL || cJSON_IsArray(j_za) == 0)
                continue;

            j_zd = cJSON_GetObjectItem(vib, "zd");
            if (j_zd == NULL || cJSON_IsArray(j_zd) == 0)
                continue;

            for (jdx = 0; jdx < num; ++jdx) {
                struct egwIP_vibration * vibr;

                info->temperature = -1;
                vibr = &info->vibration;

                memset(info->startTime, 0, EGWIP_TIME_LEN);
                ret = mhmp_get_timetag_value(j_x,
                    jdx, info->startTime, EGWIP_TIME_LEN, &vibr->ev_x);
                if (ret < 0)
                    continue;

                ret = mhmp_get_timetag_value(j_y,
                    jdx, NULL, 0, &vibr->ev_y);
                if (ret < 0)
                    continue;

                ret = mhmp_get_timetag_value(j_z,
                    jdx, NULL, 0, &vibr->ev_z);
                if (ret < 0)
                    continue;

                ret = mhmp_get_timetag_value(j_xa,
                    jdx, NULL, 0, &vibr->ev_xa);
                if (ret < 0)
                    continue;

                ret = mhmp_get_timetag_value(j_xd,
                    jdx, NULL, 0, &vibr->ev_xd);
                if (ret < 0)
                    continue;

                ret = mhmp_get_timetag_value(j_ya,
                    jdx, NULL, 0, &vibr->ev_ya);
                if (ret < 0)
                    continue;

                ret = mhmp_get_timetag_value(j_za,
                    jdx, NULL, 0, &vibr->ev_za);
                if (ret < 0)
                    continue;

                ret = mhmp_get_timetag_value(j_zd,
                    jdx, NULL, 0, &vibr->ev_zd);
                if (ret < 0)
                    continue;

                mhmp_publish_data(mvar, info, nowt);
                mhmp_publish_kafka(mvar, info, nowt);
            }
        }
    }

    free(info);
}

static void mhmp_handle_warninfo(
    struct mhmp_var * mvar, cJSON * payload, const char * which)
{
    time_t nowt;
    int warnnum, idx;
    cJSON * jdat = NULL;
    struct egwIP_warninfo * winfo;

    winfo = NULL;
    jdat = cJSON_GetObjectItem(payload, "data");
    if (jdat == NULL || cJSON_IsArray(jdat) == 0) {
        jdat = payload;
        warnnum = 1;
    } else
        warnnum = cJSON_GetArraySize(jdat);
    if (warnnum <= 0)
        return;

    winfo = (struct egwIP_warninfo *) malloc(sizeof(*winfo));
    if (winfo == NULL) {
        fputs("Error, system out of memory!\n", stderr);
        fflush(stderr);
        return;
    }

    nowt = 0;
    time(&nowt);

    for (idx = 0; idx < warnnum; ++idx) {
        int ret;
        cJSON * temp, * field;
        temp = cJSON_GetArrayItem(jdat, idx);
        if (temp == NULL || cJSON_IsObject(temp) == 0) {
            if (warnnum == 1)
                temp = jdat;
            else
                continue;
        }

        memset(winfo, 0, sizeof(*winfo));
        ret = json_find_string(temp, "deviceID", winfo->deviceID, EGWIP_NAME_MAX);
        if (ret < 0) {
            fputs("Error, cannot find deviceID.\n", stderr);
            fflush(stderr);
            continue;
        }
        ret = json_find_string(temp, "deviceName", winfo->deviceName, EGWIP_NAME_MAX);
        if (ret < 0) {
            fputs("Error, cannot find deviceName.\n", stderr);
            fflush(stderr);
            continue;
        }
        ret = json_find_string(temp, "egwIP", winfo->egwip_id, EGWIP_NAME_MAX);
        if (ret < 0) {
            fputs("Error, cannot find egwIP\n", stderr);
            fflush(stderr);
            continue;
        }

        field = cJSON_GetObjectItem(temp, "warnID");
        if (field && cJSON_IsString(field)) {
            int jdx = 0;
            char num[16];
            const char * fdot;

            num[0] = '\0';
            winfo->warnID = field->valuestring;
            fdot = strchr(winfo->warnID, '-');
            if (fdot)
                fdot = strchr(fdot, '.');
            if (fdot != NULL) {
                fdot++;
                while (*fdot != '.') {
                    char cha = *fdot++;
                    num[jdx++] = cha;
                    if (jdx >= (sizeof(num) - 1) || cha == '\0')
                        break;
                }
                num[jdx] = '\0';
                jdx = strtol(num, NULL, 0);
            } else
                jdx = 1;
            winfo->index_no = jdx - 1;
        }

        field = cJSON_GetObjectItem(temp, "warnType");
        if (field && cJSON_IsString(field))
            winfo->warnType = field->valuestring;

        field = cJSON_GetObjectItem(temp, "pushType");
        if (field && cJSON_IsString(field))
            winfo->pushType = field->valuestring;

        field = cJSON_GetObjectItem(temp, "monitoringSite");
        if (field && cJSON_IsString(field))
            winfo->monitoringSite = field->valuestring;

        field = cJSON_GetObjectItem(temp, "monitoringSiteName");
        if (field && cJSON_IsString(field))
            winfo->monitoringSiteName = field->valuestring;

        field = cJSON_GetObjectItem(temp, "diagCode");
        if (field && cJSON_IsString(field))
            winfo->diagCode = field->valuestring;

        field = cJSON_GetObjectItem(temp, "startTime");
        if (field && cJSON_IsString(field))
            winfo->startTime = field->valuestring;

        field = cJSON_GetObjectItem(temp, "endTime");
        if (field && cJSON_IsString(field))
            winfo->endTime = field->valuestring;

        field = cJSON_GetObjectItem(temp, "warnValue");
        if (field && cJSON_IsString(field))
            winfo->warnValue = field->valuestring;

        mhmp_publish_warn_info(mvar, winfo, nowt);
        mhmp_publish_warninfo_kafka(mvar, winfo, nowt, which);
    }

    free(winfo);
}

void mhmp_process_msg(struct mhmp_var * mvar,
    const char * msgbuf, int msglen)
{
    int ret;
    char cmdName[32];
    cJSON * payload = NULL;

    if (mhmp_verbose) {
        fprintf(stdout, "Received MHMP message:\n%s\n", msgbuf);
        fflush(stdout);
    }

    payload = cJSON_Parse(msgbuf);
    if (payload == NULL) {
        fprintf(stderr, "Error, corrupted MHMP payload, length: %d\n", msglen);
        fflush(stderr);
        return;
    }

    ret = json_find_string(payload, "cmd", cmdName, sizeof(cmdName));
    if (ret < 0) {
        fprintf(stderr, "Error, MHMP cmd not found: %d\n", ret);
        fflush(stderr);
        goto exit_process;
    }

    if (strcmp(cmdName, MHMP_KEEPALIVE) == 0) {
        /* update hb_recv timestamp */
        mvar->last_hb_recv = mhmp_uptime(NULL, NULL);
    } else if (strcmp(cmdName, MHMP_GETREALTIMEDATA) == 0) {
        mhmp_handle_realtime_data(mvar, payload);
    } else if (strcmp(cmdName, MHMP_GETHISTORYDATA) == 0) {
        mhmp_handle_history_data(mvar, payload);
    } else if (strcmp(cmdName, MHMP_GETWARNINGINFO) == 0) {
        mhmp_handle_warninfo(mvar, payload, "abnorm");
    } else if (strcmp(cmdName, MHMP_GETHISTORYWARNDATA) == 0) {
        mhmp_handle_warninfo(mvar, payload, "abnorm");
    } else if (strcmp(cmdName, MHMP_PUSHWARNINGINFO) == 0) {
        mhmp_handle_warninfo(mvar, payload, "alarms");
    } else {
        fprintf(stderr, "Error, unknown MHMP command: %s\n", cmdName);
        fflush(stderr);
    }

exit_process:
    if (payload != NULL)
        cJSON_Delete(payload);
}

static int mhmp_handle_dummy(void * obj, ipc_msg_t * msg)
{
    (void) obj; /* ignored */
    (void) msg;
    return 0;
}

int mhmp_ipc_init(struct mhmp_var * mvar)
{
    int ret;

    ret = load_nodes_cfg(&mvar->nodecfg, NODES_CFG_PATH);
    if (ret == -1)
        return -1;
    ret = load_templates_cfg(&mvar->tempcfg, TEMPLATES_CFG_PATH);
    if (ret == -1)
        return -2;
    ret = load_objects_cfg(&mvar->objcfg, OBJECTS_CFG_PATH);
    if (ret == -1)
        return -3;
    ret = get_board_sn(mvar->board_sn);
    if (ret < 0)
        return -4;

    mvar->ipcsess = ipc_session_new(NULL, mvar, IPC_MQTT);
    if (mvar->ipcsess == NULL) {
        fputs("Error, failed to create ipc session!\n", stderr);
        fflush(stderr);
        return -5;
    }

#if 0
    ret = ipc_session_subscribe(mvar->ipcsess, "ipc/+/+/device/+/data/property/+");
    if (ret < 0) {
        fputs("Error, failed to register topic!\n", stderr);
        fflush(stderr);
        goto err0;
    }
#endif

    ret = ipc_session_set_callbacks(mvar->ipcsess, mhmp_handle_dummy, NULL);
    if (ret < 0) {
        fputs("Error, failed to set IPC callback function!\n", stderr);
        fflush(stderr);
        return -6;
    }

    ret = ipc_session_start(mvar->ipcsess);
    if (ret < 0) {
        fputs("Error, failed to start IPC session!\n", stderr);
        fflush(stderr);
        return -7;
    }
    return 0;
}
