#include <arpa/inet.h>
#include <string.h>
#include "judger.h"
#include "tunnel_manager.h"
#include "bmser_my_log.h"
#include "bmser_msg_ring.h"
#include "dynamic_array.h"
#include "parse_cfg.h"
#include "bmser_trans.h"
#include "mqtt_msg_proc.h"
#include "id_transform.h"
#include "bmser_dev.h"
#include "bmser_project_para.h"
#include "bmser_ems_db_control_db.h"
#include "can_inf.h"

#define BMSER_DB_RETRY_TIME_LOOP_MS 3
#define BMSER_DB_RETRY_COUNT 50

typedef int (*JUDGER_TUNNEL_PROCESS)(BMSER_TRANS_JUDGER_S self, void *arg);
typedef int (*JUDGER_PUBLISH_PROCESS)(BMSER_TRANS_JUDGER_S self, void *arg);

struct TUNNEL_WITH_CFG
{
    BMSER_TRANS_COMM_TUNNEL_S tunnel;
    struct BMSER_TRANS_TUNNEL_CFG cfg;
};

static struct _BMSER_TRANS_JUDGER_S
{
    struct TUNNEL_WITH_CFG *pTunnelWithCfg_McuInBAU;  // 专门和BAU内的MCU进行通信的通道
    BMSER_DYNAMIC_ARRAY_S tunnelsWithCfg[BMSER_TRANS_STATION_ROLE_MAX];
    BMSER_DYNAMIC_ARRAY_S map_bcuid2TunnelWithCfg;
    BMSER_MSG_RING_S publishRing[BMSER_TRANS_STATION_ROLE_MAX];

    // JUDGER_PUBLISH_PROCESS publishProcess;      // 在publish MQTT前进行的操作
    JUDGER_TUNNEL_PROCESS tunnelProcess;        // 在发送给通道前进行的操作
    /* mqtt暂时不列举出来了，因为私有的直接用lib的就足够了没必要抽象 */

    // xpthread需要释放句柄资源
    BMSER_PORT_THREAD_T mqttPublishThread[BMSER_TRANS_STATION_ROLE_MAX];

    BMSER_PORT_THREAD_T dbRetryLoopThread;
    int dbRetryRunning;
    BMSER_MSG_RING_S dbRetryRing;
    pthread_mutex_t mutex; 
    int comm_state[BMSER_TRANS_MEDIUM_MAX][BMS_CFG_ROCKSNUM_MAX];       // 0: 未连接 1: 已连接
}g_judger;


int TunnelProcess(void *arg)
{
    int ret = 0;

    if (g_judger.tunnelProcess)
    {
        ret = g_judger.tunnelProcess(&g_judger, arg);
    }

    return ret;
}

/**@brief          建立从bcuid到通道的映射表
 * @param[in]      
 * @return         
 * @retval         
 * @par            修改日志:
 * <table>
 * <tr><th>Date			<th>Version	<th>Author	<th>Description
 * <tr><td>2024/03/15	<td>1.0		<td>mxg		<td>创建初始版本
 * </table>
 */
static int Bcuid2TunnelMapInit()
{
    int ret;
    int tunnelIndex;
    struct TUNNEL_WITH_CFG *pTunnelWithCfg;
    int tunnelNum = BMSER_DynamicArray_GetSize( \
        g_judger.tunnelsWithCfg[BMSER_TRANS_STATION_ROLE_MASTER]);

    g_judger.map_bcuid2TunnelWithCfg =  \
        BMSER_DynamicArray_Create(MAX_BCU_NUM_ON_BAU);
    if (g_judger.map_bcuid2TunnelWithCfg == NULL)
    {
        BMSER_MyLogError("BMSER_DynamicArray_Create failed");
        return -1;
    }

    for (tunnelIndex = 0; tunnelIndex < tunnelNum; tunnelIndex++)
    {
        int bcuIndex;
        int bcuStartID;
        int bcuNum;

        pTunnelWithCfg = BMSER_DynamicArray_Get(    \
            g_judger.tunnelsWithCfg[BMSER_TRANS_STATION_ROLE_MASTER],   \
            tunnelIndex);
        bcuStartID = pTunnelWithCfg->cfg.slaveCfg.bcuStartID;
        bcuNum = pTunnelWithCfg->cfg.slaveCfg.bcuNum;

        for (bcuIndex = bcuStartID; bcuIndex < bcuStartID + bcuNum; bcuIndex++)
        {
            ret = BMSER_DynamicArray_Put(g_judger.map_bcuid2TunnelWithCfg,  \
                                        bcuIndex,   \
                                        pTunnelWithCfg);
            if (ret < 0)
            {
                BMSER_MyLogError("BMSER_DynamicArray_Put failed");
                return -1;
            }
        }
    }

    return 0;
}

BMSER_TRANS_COMM_TUNNEL_S GetMasterTunnel(int bcuID)
{
    struct TUNNEL_WITH_CFG *pTunnelWithCfg;

    // mcu+can应用场景, trans才处理动环与堆报文
    if (bmser_jsonGetScene() == SCENE_MCU_AND_CAN)
    {
        if (bcuID == BAUDATA_BCU_ID || bcuID == POWER_ENV_BCU_ID)
        {
            return g_judger.pTunnelWithCfg_McuInBAU->tunnel;
        }
    }

    pTunnelWithCfg = BMSER_DynamicArray_Get(g_judger.map_bcuid2TunnelWithCfg, \
                                            bcuID);
    if (pTunnelWithCfg == NULL)
    {
        return NULL;
    }

    return pTunnelWithCfg->tunnel;
}

BMSER_TRANS_COMM_TUNNEL_S GetMcuInBAUTunnel()
{
    if (g_judger.pTunnelWithCfg_McuInBAU != NULL)
    {
        return g_judger.pTunnelWithCfg_McuInBAU->tunnel;
    }
    return NULL;
}

struct TUNNEL_WITH_CFG *GetMasterTunnelWithCfg(int bcuID)
{
    struct TUNNEL_WITH_CFG *pTunnelWithCfg;

    // mcu+can应用场景, trans才处理动环与堆报文
    if (bmser_jsonGetScene() == SCENE_MCU_AND_CAN)
    {
        if (bcuID == BAUDATA_BCU_ID || bcuID == POWER_ENV_BCU_ID)
        {
            return g_judger.pTunnelWithCfg_McuInBAU;
        }
    }

    pTunnelWithCfg = BMSER_DynamicArray_Get(g_judger.map_bcuid2TunnelWithCfg, \
                                            bcuID);

    return pTunnelWithCfg;
}

/**@brief          处理广播报文，转发给所有主站通道
 * @param[in]      pMsg     消息体
 * @return         操作结果
 * @retval         0：成功；    -1：失败
 * @par            修改日志:
 * <table>
 * <tr><th>Date			<th>Version	<th>Author	<th>Description
 * <tr><td>2024/03/08	<td>1.0		<td>mxg		<td>创建初始版本
 * </table>
 */
int ProcessBroadcastMsg(struct BMSER_TRANS_MSG_S *pMsg)
{
    int ret;
    int i;
    int tunnelNum;
    struct TUNNEL_WITH_CFG *pTunnelWithCfg;
    BMSER_DYNAMIC_ARRAY_S tunnelsWithCfg =  \
        g_judger.tunnelsWithCfg[BMSER_TRANS_STATION_ROLE_MASTER];

    tunnelNum = BMSER_DynamicArray_GetSize(tunnelsWithCfg);
    for (i = 0; i < tunnelNum; i++)
    {
        pTunnelWithCfg = BMSER_DynamicArray_Get(tunnelsWithCfg, i);
        if (pTunnelWithCfg->cfg.stationRole == BMSER_TRANS_STATION_ROLE_MASTER)
        {
            /*
                send to tunnel's msg ring
            */
            ret = SendMsgToTunnel(pTunnelWithCfg->tunnel, \
                                    (char *)pMsg,   \
                                    sizeof(*pMsg) + pMsg->datalen);
            if (ret)
            {
                BMSER_MyLogError("send msg to tunnel failed");

                // 这里不退出，因为广播出去可能有部分通道就是不通的. ？
                continue;
            }
        }
    }

    return 0;
}

int SendMsgToRetryThread(struct BMSER_TRANS_MSG_S* pMsg, BMSER_MQTT_CONNECTION_TYPE_E mqtt_conn_type, uint16_t bandId) {
    if (bmser_jsonGetScene() == SCENE_MCU_AND_CAN || (pMsg->info.dbInfo.devId != BAUDATA_BCU_ID && pMsg->info.dbInfo.devId != POWER_ENV_BCU_ID)) {
        struct TUNNEL_WITH_CFG *pTunnelWithCfg = GetMasterTunnelWithCfg(pMsg->info.dbInfo.devId);
        if (!pTunnelWithCfg) {
            BMSER_MyLogError("send comm state to judger failed");
            return -1;
        }
        if (!g_judger.comm_state[pTunnelWithCfg->cfg.medium][pMsg->info.dbInfo.devId]) {
            return -1;
        }
    }

    struct BMSER_MSG_RING_INFO ringInfo = {0};
    size_t msg_size = pMsg->datalen + sizeof(struct BMSER_TRANS_MSG_S);

    ringInfo.header.dataSize = msg_size + sizeof(BMSER_MQTT_CONNECTION_TYPE_E) + sizeof(bandId);
    memcpy(ringInfo.buffer, pMsg, msg_size);
    memcpy(ringInfo.buffer + msg_size, &mqtt_conn_type,
           sizeof(BMSER_MQTT_CONNECTION_TYPE_E));
    memcpy(ringInfo.buffer + msg_size + sizeof(BMSER_MQTT_CONNECTION_TYPE_E), &bandId,
           sizeof(bandId));
    int ret = BMSER_MsgRingPush(g_judger.dbRetryRing, &ringInfo);
    if (ret) {
        BMSER_MyLogError("ring push failed");

        return -1;
    }

    return 0;
}

int ProcessUnicastMsg(struct BMSER_TRANS_MSG_S *pMsg)
{
    int ret;
    int slaveid;
    BMSER_TRANS_COMM_TUNNEL_S tunnel;

    if (pMsg->type == BMSER_TRANS_MSG_DB)
    {
        slaveid = pMsg->info.dbInfo.devId;
    }
    else if (pMsg->type == BMSER_TRANS_MSG_BINARY)
    {
        slaveid = pMsg->info.binaryInfo.msgFlowBcuId;
    }
    else
    {
        BMSER_MyLogError("unknown %d msg type", pMsg->type);
        return -1;
    }

    tunnel = GetMasterTunnel(slaveid);

    /*
        send to tunnel's msg ring
    */
    ret = SendMsgToTunnel(tunnel, (char *)pMsg, \
                        sizeof(*pMsg) + pMsg->datalen);
    if (ret)
    {
        BMSER_MyLogError("send msg to tunnel failed");

        return -1;
    }

    return 0;
}

/**@brief          处理需要上报的消息时在BDU上需要做的一些策略
 * @param[in]      pMsg         消息内容
 * @param[in]      ring         消息环
 * @return         操作结果
 * @retval         0：成功；-1：失败 1：不再需要发送MQTT
 * @par            修改日志:
 * <table>
 * <tr><th>Date			<th>Version	<th>Author	<th>Description
 * <tr><td>2024/03/15	<td>1.0		<td>mxg		<td>创建初始版本
 * </table>
 */
static int ProcessPublishMsg_BDU(struct BMSER_TRANS_MSG_S *pMsg,    \
                                BMSER_MSG_RING_S ring)
{
    int ret;

    /* 
        在BDU上面

        id转换：
        1. 从站收到BAU发来的消息，上报MQTT、转发时都需要转换id
        2. 主站收到消息，转发给别的从站时，需要将舱内id转为平铺的设备id。如果直接上报mqtt，则不需要转换

        1. 如果消息的协议类型属于专门用于转发的（BMSER_TRANS_MSG_FORWARDING）
        由trans的通道a收到之后，需要转发给通道b（根据目标地址确定）、
            这个时候，需要处理广播报文，转发给所有主站通道，注意不再需要发送MQTT消息
        2. 如果不是
            a. 消息来源是主站收到的（主控响应）,需要转发给所有从站（bau）、mqtt上报，
            因为消息来源无法区分响应的是来自bdu自身还是bau的请求
            b. 从站收到的（总控请求），直接mqtt上报

        在bau上，BMSER_TRANS_MSG_FORWARDING类型的协议只会是bau主动发送，bcu回复的协议
        又按照之前定义的协议回复。另外，bau上面的从站主要是处理和EMS / PCS通讯的，
        这些通宵协议没有放到trans内。所以在bau上，收到
        1. 主站收到的消息，直接mqtt上报
        2. 其他情况不会有，不处理。也可以不进行判断是否是主站，因为其他情况不会有
     */
    if (ring == g_judger.publishRing[BMSER_TRANS_STATION_ROLE_SLAVE])
    {
        ret = IdTransform_Bcuid2Slaveid(pMsg);
        if (ret)
        {
            BMSER_MyLogError("transform bcu id to slave id failed");
            return -1;
        }
    }

    if (pMsg->type == BMSER_TRANS_MSG_FORWARDING)
    {
        pMsg->type = BMSER_TRANS_MSG_BINARY;        // 主站通道按照bin发送
        if (pMsg->info.binaryInfo.msgFlowBcuId == BROADCAST_BCUID || pMsg->info.binaryInfo.msgFlowBcuId == BROADCAST_BCUID_FF)
        {
            ret = ProcessBroadcastMsg(pMsg);
            if (ret)
            {
                BMSER_MyLogError("process broadcast message failed");
                // 这里不退出，因为广播出去可能有部分通道就是不通的. ？
            }
        }
        else
        {
            ret = ProcessUnicastMsg(pMsg);
            if (ret)
            {
                BMSER_MyLogError("process unicast message failed");
                return -1;
            }
        }

        return 1;
    }
    else
    {
        if (ring == g_judger.publishRing[BMSER_TRANS_STATION_ROLE_MASTER])
        {
            int i;
            int slaveTunnelNum = BMSER_DynamicArray_GetSize( \
                g_judger.tunnelsWithCfg[BMSER_TRANS_STATION_ROLE_SLAVE]);
            struct TUNNEL_WITH_CFG *pTunnelWithCfg;

            ret = IdTransform_Slaveid2Bcuid(pMsg);
            if (ret)
            {
                BMSER_MyLogError("transform bcu id to slave id failed");
                return -1;
            }

            for (i = 0; i < slaveTunnelNum; i++)
            {
                pTunnelWithCfg = BMSER_DynamicArray_Get(  \
                    g_judger.tunnelsWithCfg[BMSER_TRANS_STATION_ROLE_SLAVE], \
                    i);
                ret = SendMsgToTunnel(pTunnelWithCfg->tunnel,   \
                                        (char *)pMsg, \
                                        sizeof(*pMsg) + pMsg->datalen);
                if (ret < 0)
                {
                    BMSER_MyLogError("SendMsgToTunnel failed");
                }
            }
        }else{}
    }

    return 0;
}

static int PublishProcess(struct BMSER_TRANS_MSG_S *pMsg,   \
                            BMSER_MSG_RING_S ring)
{
    int ret = 0;
    enum BMSER_DEV_TYPE devType = BMSER_TRANS_GetDevType();

    if (devType == BMSER_BDU)
    {
        ret = ProcessPublishMsg_BDU(pMsg, ring);
        if (ret)
        {
            BMSER_MyLogError("process msg for bdu failed");
            return -1;
        }
    }

    return ret;
}

static int ProcessRingMsg(BMSER_MSG_RING_S ring)
{
    int ret;
    struct BMSER_MSG_RING_INFO ringInfo;

    if (NULL == ring)
    {
        return 0;
    }

    ret = BMSER_MsgRingPop(ring, &ringInfo);
    if (ret)
    {
        BMSER_MyLogError("ring pop failed");

        return -1;
    }

#if 0
    if (g_judger.publishProcess)
    {
        ret = g_judger.publishProcess(&g_judger, ringInfo.buffer);
        if (ret)
        {
            BMSER_MyLogError("publish process failed");

            return -1;
        }
    }
#else
    ret = PublishProcess(   \
            (struct BMSER_TRANS_MSG_S *)ringInfo.buffer, ring);
    if (ret < 0)
    {
        BMSER_MyLogError("publish process failed");

        return -1;
    }
    else if (ret == 1)
    {
        /*ret == 1表示不需要上报MQTT*/
        return 0;
    }
#endif

    ret = MqttMsgProc_PublishMsg((char *)ringInfo.buffer);
    if (ret)
    {
        BMSER_MyLogError("publish msg failed");

        return -1;
    }

    return 0;
}

static void *MqttPublishThread(void *arg)
{
    int ret;
    BMSER_MSG_RING_S ring = (BMSER_MSG_RING_S)arg;

    if (NULL == ring)
    {
        BMSER_MyLogWarn("oops! ring[%p] is null", ring);

        return NULL;
    }

    while (1)
    {
        ret = ProcessRingMsg(ring);
        if (ret)
        {
            BMSER_MyLogError("process ring for master msg failed");

            continue;
        }
    }

    BMSER_MyLogError("MqttPublishThread exit unpossible");

    return NULL;
}

static void* DbRetryLoopThread(void* arg) {
    int ret         = 0;
    int retry_count = 0;

    while (g_judger.dbRetryRunning) {
        int is_send_packet                  = 0;
        struct BMSER_MSG_RING_INFO ringInfo = {0};
        ret                                 = BMSER_MsgRingPeek(g_judger.dbRetryRing, &ringInfo);
        if (ret) {
            BMSER_MyLogError("ring peek failed");
            continue;
        }

        struct BMSER_TRANS_MSG_S* pMsg = (struct BMSER_TRANS_MSG_S*)ringInfo.buffer;
        BMSER_MQTT_CONNECTION_TYPE_E mqtt_conn_type =
            (BMSER_MQTT_CONNECTION_TYPE_E) *
            ((char*)ringInfo.buffer + sizeof(struct BMSER_TRANS_MSG_S) + pMsg->datalen);
        EMS_DEV_REG_DATA_S regData = {0};
        uint16_t bankId = 0, rackId = 0;
        struct DB_INFO_S* pDbInfo = &pMsg->info.dbInfo;
        bmser_dev_calcBankIdAndRackIdByDevId(pDbInfo->devId, &bankId, &rackId);
        bankId =
            (uint16_t) *
            ((char*)ringInfo.buffer + sizeof(struct BMSER_TRANS_MSG_S) + pMsg->datalen + sizeof(BMSER_MQTT_CONNECTION_TYPE_E));

        uint16_t dataBuf[TRANS_MSG_BUFFER_SIZE] = {0};
        regData.bankId                          = bankId;
        regData.rackId                          = rackId;
        regData.startAddr                       = pDbInfo->startAddr;
        regData.regCount                        = pDbInfo->regCnt;
        regData.len                             = regData.regCount * 2;
        if (regData.regCount > TRANS_MSG_BUFFER_SIZE) {
            retry_count = 0;
            /* 单次查询数量限制为TRANS_MSG_BUFFER_SIZE */
            BMSER_MyLogError("func READ: regCount over the limit(%d)!", TRANS_MSG_BUFFER_SIZE);
            ret = BMSER_MsgRingPop(g_judger.dbRetryRing, &ringInfo);
            if (ret) {
                BMSER_MyLogError("ring pop failed");
            }
            usleep(BMSER_DB_RETRY_TIME_LOOP_MS * 1000);
            continue;
        }

        int isHeapOrEnvData = bmser_jsonGetScene() != SCENE_MCU_AND_CAN &&
                (pMsg->info.dbInfo.devId == BAUDATA_BCU_ID || pMsg->info.dbInfo.devId == POWER_ENV_BCU_ID);
        if (retry_count == 0) {
            // 非MCU+CAN场景下同步堆数据时，bms_app收到Set/binaryCmd会将索引db对应的值也更新。
            // 所以该情况下不需要通过Get/binaryCmd请求下位机，从而更新索引db对应的值
            if (!(pDbInfo->func == FUNC_READ_DB && isHeapOrEnvData) &&
                (CheckDBInSyncList(pDbInfo->startAddr) == 1)) {
                BmserProtocolExt_RestRegisterDbLastRecvTime(&regData);
            }
            if (pDbInfo->func == FUNC_WRITE_DB && isHeapOrEnvData) {
                ret = PublishSetBinaryCmd(pMsg, bankId);
                if (ret) {
                    BMSER_MyLogError("PublishGetRespBinaryCmd failed, ret = %d", ret);
                    usleep(BMSER_DB_RETRY_TIME_LOOP_MS * 1000);
                    continue;
                }
            }
        }

        BMSER_CHECK_DB_INVALID_INFO_S info = {0};
        uint16_t invalid_db_addr[TRANS_MSG_BUFFER_SIZE]     = {0};
        long long last_recv[TRANS_MSG_BUFFER_SIZE] = {0};
        info.invalid_db_addr = invalid_db_addr;
        info.last_recv = last_recv;
        ret = BmserProtocolExt_ReadRegisterDb_CheckByTime(&regData, (uint8_t*)dataBuf, regData.len, 0, &info);
        if (ret == 0) {
            if (pDbInfo->func == FUNC_WRITE_DB) {
                pDbInfo->func  = FUNC_WRITE_DB_RESP;
            } else if (pDbInfo->func == FUNC_READ_DB) {
                pDbInfo->func  = FUNC_READ_DB_RESP;
            }
            is_send_packet = 1;
        } else {
            BMSER_MyLogError("read reg db fail addr = 0x%x, ret = %d.", regData.startAddr, ret);
            retry_count++;
            if (retry_count < BMSER_DB_RETRY_COUNT) {
                BMSER_TRANS_COMM_TUNNEL_S tunnel = NULL;
                tunnel                           = GetMasterTunnel(pDbInfo->devId);
                /*
                    send to tunnel's msg ring
                */
                if (NULL == tunnel) {
                    BMSER_MyLogWarn("tunnel %d not found", pDbInfo->devId);
                    usleep(BMSER_DB_RETRY_TIME_LOOP_MS * 1000);
                    continue;
                }

                if (pDbInfo->func == FUNC_WRITE_DB) {
                    if (retry_count > 1) {
                        usleep(BMSER_DB_RETRY_TIME_LOOP_MS * 1000);
                        continue;
                    }
                    pMsg->type = BMSER_TRANS_MSG_DB;
                }

                ret = SendMsgToTunnel(tunnel, (char*)pMsg, sizeof(struct BMSER_TRANS_MSG_S) + pMsg->datalen);
                if (ret) {
                    BMSER_MyLogError("send msg to tunnel failed");
                    usleep(BMSER_DB_RETRY_TIME_LOOP_MS * 1000);
                    continue;
                }
            } else {
                if (pDbInfo->func == FUNC_WRITE_DB) {
                    pDbInfo->func  = FUNC_WRITE_DB_RESP;
                } else if (pDbInfo->func == FUNC_READ_DB) {
                    pDbInfo->func  = FUNC_READ_DB_RESP;
                }
                is_send_packet = 1;
            }
        }

        if (is_send_packet) {
            retry_count = 0;
            ret         = PublishGetRespBinaryCmd(pDbInfo, (char*)dataBuf, bankId, rackId, mqtt_conn_type);
            if (ret) {
                BMSER_MyLogError("PublishGetRespBinaryCmd failed, ret = %d", ret);
            }
            ret = BMSER_MsgRingPop(g_judger.dbRetryRing, &ringInfo);
            if (ret) {
                BMSER_MyLogError("ring pop failed");
            }
        }

        usleep(BMSER_DB_RETRY_TIME_LOOP_MS * 1000);
    }

    BMSER_MyLogError("DbRetryLoopThread exit unpossible");

    return NULL;
}

static int GetTunnelNumFromCfg(enum BMSER_TRANS_STATION_ROLE role,  \
                                struct BMSER_TRANS_CFG *pCfg)
{
    int tunnelNum = 0;
    int tunnelNumMax;
    int tunnelIndex;
    struct BMSER_TRANS_TUNNEL_CFG *pTunnelCfg;

    tunnelNumMax = BMSER_DynamicArray_GetSize(pCfg->tunnelCfgs);
    for (tunnelIndex = 0; tunnelIndex < tunnelNumMax; tunnelIndex++)
    {
        pTunnelCfg = BMSER_DynamicArray_Get(pCfg->tunnelCfgs, tunnelIndex);
        if (pTunnelCfg->stationRole == role)
        {
            tunnelNum++;
        }
        else
        {
            continue;
        }
    }

    return tunnelNum;
}

int SendMsgToJudger(enum BMSER_TRANS_STATION_ROLE stationRole,  \
                    char *pMsg, size_t size)
{
    int ret;
    struct BMSER_MSG_RING_INFO ringInfo;
    BMSER_MSG_RING_S publishRing;
    int ringFreeSize;

    if (NULL == pMsg)
    {
        BMSER_MyLogError("null param");

        return -1;
    }

    publishRing = g_judger.publishRing[stationRole];

    /* 如果消息超过了容量会丢老的消息，
    这个时候简单一点处理，看下空余空间够不够，够了再放到队列 */
    ringFreeSize = MSG_RING_CAPACITY - BMSER_MsgRingGetSize(publishRing);
    if (ringFreeSize < MAX_MSG_SIZE_IN_RING)
    {
        sleep(1);

        ringFreeSize = MSG_RING_CAPACITY - BMSER_MsgRingGetSize(publishRing);
        if (ringFreeSize < MAX_MSG_SIZE_IN_RING)
        {
            BMSER_MyLogError("ring free size %d is lower than %lu after 1s", \
                            ringFreeSize, MAX_MSG_SIZE_IN_RING);
            BMSER_MyLogToFile("ring free size %d is lower than %lu after 1s", \
                            ringFreeSize, MAX_MSG_SIZE_IN_RING);
            return -1;
        }
    }

    ringInfo.header.dataSize = size;
    memcpy(ringInfo.buffer, pMsg, size);

    ret = BMSER_MsgRingPush(publishRing, &ringInfo);
    if (ret)
    {
        BMSER_MyLogError("push failed");

        return -1;
    }

    return 0;
}

int SendCommStateToJudger(int bcuID, int medium, int state) {
    struct TUNNEL_WITH_CFG *pTunnelWithCfg = GetMasterTunnelWithCfg(bcuID);
    if (!pTunnelWithCfg) {
        BMSER_MyLogDebug("send comm state to judger failed, bcuID: %d, medium: %d", bcuID, medium);
        return -1;
    }
    int is_connected = 0;
    if (medium != pTunnelWithCfg->cfg.medium) {
        BMSER_MyLogError("unknown physical medium");
        return -1;
    }

    pthread_mutex_lock(&g_judger.mutex);

    if (g_judger.comm_state[medium][bcuID] == 0 && g_judger.comm_state[medium][bcuID] != state) {
        BMSER_MsgRingClear(g_judger.dbRetryRing);
    }
    g_judger.comm_state[medium][bcuID] = state;
    for (int i = 0; i < BMSER_TRANS_MEDIUM_MAX; ++i) {
        for (int j = 0; j < BMS_CFG_ROCKSNUM_MAX; ++j) {
            if (g_judger.comm_state[i][j] == 1) {
                is_connected = 1;
                break;
            }
        }
    }

    pthread_mutex_unlock(&g_judger.mutex);

    int ret = 0, i = 0;
    int bank_num = bmser_jsonGetBankNum();
    EMS_DEV_REG_DATA_S regData = {0};

    regData.rackId = 0;
    regData.startAddr = LINUX_RESERVE_DB_ADDR;
    regData.regCount = 1;
    regData.len = regData.regCount * 2;

    // 包含堆
    for(i = 0;i < bank_num; i++)
    {
        regData.bankId = i;
        ret = BmserProtocolExt_UpdateDbBit(&regData, UI_FAILT_BIT_INDEX, !is_connected);
        if(ret) {
            dy_syslog(LOG_DEBUG, "BmserProtocolExt_UpdateDbBit err, ret %d.\n", ret);
        }
    }

    return 0;
}

static int CreateTunnels(enum BMSER_TRANS_STATION_ROLE role, \
                            struct BMSER_TRANS_CFG *pCfg)
{
    int ret;
    int tunnelNum;
    int tunnelIndex;
    BMSER_TRANS_COMM_TUNNEL_S tunnel;
    struct BMSER_TRANS_TUNNEL_CFG *pTunnelCfg;
    struct TUNNEL_WITH_CFG *pTunnelWithCfg;
    struct BMSER_MSG_RING_CFG cfg = {
        .autoLock = true,
        .capacity = MSG_RING_CAPACITY,
        .headerSize = sizeof(struct BMSER_MSG_RING_HEADER),
        .isBlocking = true,
    };

    tunnelNum = GetTunnelNumFromCfg(role, pCfg);
    if (tunnelNum <= 0)
    {
        BMSER_MyLogWarn("role[%d], tunnelNum <= 0", role);
        return 0;   // no tunnel, but is typical
    }

    g_judger.tunnelsWithCfg[role] = BMSER_DynamicArray_Create(tunnelNum);
    if (g_judger.tunnelsWithCfg[role] == NULL)
    {
        BMSER_MyLogError("create tunnelsWithCfg failed");
        return -1;
    }

    for (tunnelIndex = 0; tunnelIndex < tunnelNum; tunnelIndex++)
    {
        pTunnelCfg = BMSER_DynamicArray_Get(pCfg->tunnelCfgs, \
                                            tunnelIndex);
        if (NULL == pTunnelCfg)
        {
            BMSER_MyLogError("tunnel cfg is NULL");

            return -1;
        }

        tunnel = BMSER_TRANS_MakeTunnel(*pTunnelCfg);
        if (NULL == tunnel)
        {
            BMSER_MyLogError("make tunnel failed");

            return -1;
        }

        ret = RegisterSendMsgToJudgerCallback(tunnel, SendMsgToJudger);
        if (ret)
        {
            BMSER_MyLogError("register send msg to judger callback failed");

            return -1;
        }

        ret = RegisterSendCommStateToJudgerCallback(tunnel, SendCommStateToJudger);
        if (ret)
        {
            BMSER_MyLogError("register send comm state to judger callback failed");

            return -1;
        }

        pTunnelWithCfg = BMSER_PORT_Calloc(1, sizeof(struct TUNNEL_WITH_CFG));
        if (NULL == pTunnelWithCfg)
        {
            BMSER_MyLogError("calloc failed");

            return -1;  // 此处不进行内存回收，保证进程退出即可
        }

        pTunnelWithCfg->tunnel = tunnel;
        pTunnelWithCfg->cfg = *pTunnelCfg;

        ret = BMSER_DynamicArray_Put(g_judger.tunnelsWithCfg[role],   \
                                    tunnelIndex,    \
                                    pTunnelWithCfg);
        if (ret)
        {
            BMSER_MyLogError("tunnel put failed");

            return -1;
        }
        /* 由于BAU消息流功能不受场景限制，因此任意场景都创建BAU通道 */
        /* 需要把和mcu形式的BAU进行通信的通道找出来 */
        if (pTunnelCfg->medium == BMSER_TRANS_CAN_MEDIUM && \
            (g_judger.pTunnelWithCfg_McuInBAU == NULL ||   \
                (pTunnelCfg->slaveCfg.bcuNum <   \
                g_judger.pTunnelWithCfg_McuInBAU->cfg.slaveCfg.bcuNum)))
        {
            g_judger.pTunnelWithCfg_McuInBAU = pTunnelWithCfg;
        }
    }

    g_judger.publishRing[role] = BMSER_MsgRingCreate(cfg);
    if (g_judger.publishRing[role] == NULL)
    {
        BMSER_MyLogError("msg ring create failed");

        return -1;
    }

    return 0;
}

static int StartTunnels()
{
    int ret;
    int role;
    int tunnelIndex;
    int tunnelNum;
    struct TUNNEL_WITH_CFG *pTunnelWithCfg;

    for (role = 0; role < BMSER_TRANS_STATION_ROLE_MAX; role++)
    {
        tunnelNum = BMSER_DynamicArray_GetSize(g_judger.tunnelsWithCfg[role]);
        for (tunnelIndex = 0; tunnelIndex < tunnelNum; tunnelIndex++)
        {
            pTunnelWithCfg =    \
                BMSER_DynamicArray_Get(g_judger.tunnelsWithCfg[role],   \
                                        tunnelIndex);
            ret = BMSER_TRANS_StartTunnel(pTunnelWithCfg->tunnel);
            if (ret)
            {
                BMSER_MyLogError("start tunnel failed");

                return -1;
            }
        }
    }
    return 0;
}

static int DbRetryInit() {
    struct BMSER_MSG_RING_CFG cfg = {
        .autoLock = true,
        .capacity = MSG_RING_CAPACITY * 10,
        .headerSize = sizeof(struct BMSER_MSG_RING_HEADER),
        .isBlocking = true,
    };
    g_judger.dbRetryRing = BMSER_MsgRingCreate(cfg);
    if (g_judger.dbRetryRing == NULL) {
        BMSER_MyLogError("db retry ring create failed");
        return -1;
    }

    return 0;
}

int Judger_Init()
{
    int ret;
    int i;
    struct BMSER_TRANS_CFG *pCfg = GetCfg();

    for (i = 0; i < BMSER_TRANS_STATION_ROLE_MAX; i++)
    {
        ret = CreateTunnels(i, pCfg);
        if (ret)
        {
            BMSER_MyLogError("CreateTunnels failed");

            return -1;
        }
    }

    ret = Bcuid2TunnelMapInit();
    if (ret)
    {
        BMSER_MyLogError("Bcuid2TunnelMapInit failed");

        return -1;
    }

    ret = StartTunnels();
    if (ret)
    {
        BMSER_MyLogError("start tunnels failed");

        return -1;
    }

    ret = DbRetryInit();
    if (ret) {
        BMSER_MyLogError("db retry init failed");

        return -1;
    }

    ret = pthread_mutex_init(&g_judger.mutex, NULL);
    if (ret) {
        BMSER_MyLogError("xpthread_mutex_init failed");

        return -1;
    }

    /*
        mqtt
    */
    ret = MqttMsgProc_Init();
    if (ret)
    {
        BMSER_MyLogError("mqtt init failed");
    }

    return 0;
}

int Judger_Start()
{
    int i;
    int ret;

    /* start mqtt receiving thread*/
    ret = MqttMsgProc_Start();
    if (ret)
    {
        BMSER_MyLogError("MqttMsgProc_Start failed");

        return -1;
    }

    for (i = 0; i < BMSER_TRANS_STATION_ROLE_MAX; i++)
    {
        /* mqtt publish thread */
        ret = BMSER_PORT_CreateThread_Autodestroy(&g_judger.mqttPublishThread[i],  \
                                        MqttPublishThread,  \
                                        g_judger.publishRing[i]);
        if (ret)
        {
            BMSER_MyLogError("MqttPublishThread create failed");

            return -1;
        }
    }
    /* db retry thread */
    g_judger.dbRetryRunning = 1;
    ret = BMSER_PORT_CreateThread_Autodestroy(&g_judger.dbRetryLoopThread,  \
                                    DbRetryLoopThread,  \
                                    NULL);
    if (ret)
    {
        BMSER_MyLogError("DbRetryLoopThread create failed");

        return -1;
    }

    return 0;
}

