/*
 * @Author: ghzhang gogosanmao@126.com
 * @Date: 2023-06-19 18:33:59
 * @LastEditors: ybzhou yibo.zhou@lnxall.com
 * @LastEditTime: 2026-03-24 09:44:33
 * @FilePath: \proto_forward\src\port_transf.c
 * @Description: 这是默认设置,请设置`customMade`, 打开koroFileHeader查看配置 进行设置: https://github.com/OBKoro1/koro1FileHeader/wiki/%E9%85%8D%E7%BD%AE
 */
#include "opcua/ua-server.h"
#include "prot_transf.h"
#include "tag.h"
#include "common.h"
#include <math.h>
#include <errno.h>
#include "collector-api.h"
#include <sys/syslog.h>
#include <syslog.h>
#include "mqtt_session.h"
#include "lnxall_list.h"
#include "dy_utils/dy_common.h"
#include "dy_utils/dy_ipc.h"
#include "dy_utils/hash_intptr.h"
#include "dy_utils/modbus_utils.h"
#include "dy_utils/period_service.h"
#include "dy_utils/bitmap.h"
#include "sys/param.h"
#include "assert.h"
#include <lnxall_ubuslog.h>
#include "prot_transf.h"
#include "protocol.h"
#include "pthread.h"
#include "mac_mul/raw_sock.h"
#include "mac_mul/mac_protocol.h"
#include "cmdclient.h"
#include "opcua/ua_client.h"
#include "cc_lc_com/define.h"
#include "cs104_connection.h"
#include "abnormal_log/abnormal_log.h"
#include "bms_map.h"
#include "custom.h"
#include "ydt1363.h"
#ifdef EN_EMS
#include "../ems/jsonct.h"
#endif

#define COMMUNIT_RETRY_MAX 5

#define GET_FOR_A(sour) ((sour)[0])
#define GET_FOR_AB(sour) ((sour)[0] | ((sour)[1] << 8))
#define GET_FOR_BA(sour) ((sour)[1] | ((sour)[0] << 8))
#define GET_FOR_ABCD(sour) ((sour)[2] | ((sour)[3] << 8) | ((sour)[0] << 16) | ((sour)[1] << 24))
#define GET_FOR_DCBA(sour) ((sour)[1] | ((sour)[0] << 8) | ((sour)[3] << 16) | ((sour)[2] << 24))
#define GET_FOR_BADC(sour) ((sour)[3] | ((sour)[2] << 8) | ((sour)[1] << 16) | ((sour)[0] << 24))
#define GET_FOR_CDAB(sour) ((sour)[0] | ((sour)[1] << 8) | ((sour)[2] << 16) | ((sour)[3] << 24))
#define SET_FOR_A(value, sour) \
    do                         \
    {                          \
        (sour)[0] = (value);   \
    } while (0)
#define SET_FOR_AB(value, sour)   \
    do                            \
    {                             \
        (sour)[0] = (value);      \
        (sour)[1] = (value) >> 8; \
    } while (0)
#define SET_FOR_BA(value, sour)   \
    do                            \
    {                             \
        (sour)[1] = (value);      \
        (sour)[0] = (value) >> 8; \
    } while (0)
#define SET_FOR_ABCD(value, sour)  \
    do                             \
    {                              \
        (sour)[2] = (value);       \
        (sour)[3] = (value) >> 8;  \
        (sour)[0] = (value) >> 16; \
        (sour)[1] = (value) >> 24; \
    } while (0)
#define SET_FOR_DCBA(value, sour)  \
    do                             \
    {                              \
        (sour)[1] = (value);       \
        (sour)[0] = (value) >> 8;  \
        (sour)[3] = (value) >> 16; \
        (sour)[2] = (value) >> 24; \
    } while (0)
#define SET_FOR_BADC(value, sour)  \
    do                             \
    {                              \
        (sour)[3] = (value);       \
        (sour)[2] = (value) >> 8;  \
        (sour)[1] = (value) >> 16; \
        (sour)[0] = (value) >> 24; \
    } while (0)
#define SET_FOR_CDAB(value, sour)  \
    do                             \
    {                              \
        (sour)[0] = (value);       \
        (sour)[1] = (value) >> 8;  \
        (sour)[2] = (value) >> 16; \
        (sour)[3] = (value) >> 24; \
    } while (0)

#define GET_INT(tag_v, sour, order)                                     \
    do                                                                  \
    {                                                                   \
        tag_v.type = 0;                                                 \
        switch (order)                                                  \
        {                                                               \
        case ORDER_A:                                                   \
            tag_v.value.to_int = (int8_t)GET_FOR_A(sour);               \
            break;                                                      \
        case ORDER_AB:                                                  \
            tag_v.value.to_int = (int16_t)GET_FOR_AB(sour);             \
            break;                                                      \
        case ORDER_BA:                                                  \
            tag_v.value.to_int = (int16_t)GET_FOR_BA(sour);             \
            break;                                                      \
        case ORDER_ABCD:                                                \
            tag_v.value.to_int = (int32_t)GET_FOR_ABCD(sour);           \
            break;                                                      \
        case ORDER_DCBA:                                                \
            tag_v.value.to_int = (int32_t)GET_FOR_DCBA(sour);           \
            break;                                                      \
        case ORDER_BADC:                                                \
            tag_v.value.to_int = (int32_t)GET_FOR_BADC(sour);           \
            break;                                                      \
        case ORDER_CDAB:                                                \
            tag_v.value.to_int = (int32_t)GET_FOR_CDAB(sour);           \
            break;                                                      \
        case ORDER_ABCDEFGH:                                             \
            tag_v.type = 1; \
            tag_v.value.to_float = (int64_t)GET_FOR_ABCD(sour+4) | ((int64_t)GET_FOR_ABCD(sour)<<32);       \
            break;                                                       \
        case ORDER_GHEFCDAB:                                             \
            tag_v.type = 1; \
            tag_v.value.to_float = (int64_t)GET_FOR_CDAB(sour) | ((int64_t)GET_FOR_CDAB(sour+4)<<32);       \
            break;                                                       \
        case ORDER_BADCFEHG:                                             \
            tag_v.type = 1; \
            tag_v.value.to_float = (int64_t)GET_FOR_BADC(sour+4) | ((int64_t)GET_FOR_BADC(sour)<<32);       \
            break;                                                       \
        case ORDER_HGFEDCBA:                                             \
            tag_v.type = 1; \
            tag_v.value.to_float = (int64_t)GET_FOR_DCBA(sour) | ((int64_t)GET_FOR_DCBA(sour+4)<<32);       \
            break;                                                       \
        default:                                                        \
            proto_syslog(LOG_ERR, "GET_INT order error, order=%d", order); \
            break;                                                      \
        }                                                               \
    } while (0)

#define SET_INT(tag, sour, re_len, order)                                                                                                                              \
    do                                                                                                                                                                 \
    {                                                                                                                                                                  \
        int32_t value = 0;                                                                                                                                             \
        union int_int                                                                                                                                                 \
        {                                                                                                                                                              \
            long long_value;                                                                                                                                           \
            int32_t u32_arr[2];                                                                                                                                        \
        } long_int_v = {0};                                                                                                                                                  \
        if (check_float_is_1(tag->scale))                                                                                                                              \
        {                                                                                                                                                              \
            value = long_int_v.long_value = (tag->data_type == 0) ? (float)(write_tag->write_cache.to_int - (int)tag->offset) : write_tag->write_cache.to_float - tag->offset;                 \
        }                                                                                                                                                              \
        else                                                                                                                                                           \
        {                                                                                                                                                              \
            value = long_int_v.long_value = (tag->data_type == 0) ? (write_tag->write_cache.to_int - tag->offset) / tag->scale : (write_tag->write_cache.to_float - tag->offset) / tag->scale; \
        }                                                                                                                                                              \
        switch (order)                                                                                                                                                 \
        {                                                                                                                                                              \
        case ORDER_A:                                                                                                                                                 \
            SET_FOR_A((int8_t)value, sour);                                                                                                                          \
            re_len = 1;                                                                                                                                                \
            break;                                                                                                                                                     \
        case ORDER_AB:                                                                                                                                                 \
            SET_FOR_AB((int16_t)value, sour);                                                                                                                          \
            re_len = 2;                                                                                                                                                \
            break;                                                                                                                                                     \
        case ORDER_BA:                                                                                                                                                 \
            SET_FOR_BA((int16_t)value, sour);                                                                                                                          \
            re_len = 2;                                                                                                                                                \
            break;                                                                                                                                                     \
        case ORDER_ABCD:                                                                                                                                               \
            SET_FOR_ABCD((int32_t)value, sour);                                                                                                                        \
            re_len = 4;                                                                                                                                                \
            break;                                                                                                                                                     \
        case ORDER_DCBA:                                                                                                                                               \
            SET_FOR_DCBA((int32_t)value, sour);                                                                                                                        \
            re_len = 4;                                                                                                                                                \
            break;                                                                                                                                                     \
        case ORDER_BADC:                                                                                                                                               \
            SET_FOR_BADC((int32_t)value, sour);                                                                                                                        \
            re_len = 4;                                                                                                                                                \
            break;                                                                                                                                                     \
        case ORDER_CDAB:                                                                                                                                               \
            SET_FOR_CDAB((int32_t)value, sour);                                                                                                                        \
            re_len = 4;                                                                                                                                                \
            break;                                                                                                                                                     \
        case ORDER_ABCDEFGH:                                                                                                                                           \
            tag->data_type = 1;                                                                                                                                         \
            SET_FOR_ABCD(long_int_v.u32_arr[0], sour + 4);                                                                                                             \
            SET_FOR_ABCD(long_int_v.u32_arr[1], sour);                                                                                                                 \
            re_len = 8;                                                                                                                                                \
            break;                                                                                                                                                     \
        case ORDER_GHEFCDAB:                                                                                                                                           \
            tag->data_type = 1;                                                                                                                                         \
            SET_FOR_CDAB(long_int_v.u32_arr[0], sour);                                                                                                                 \
            SET_FOR_CDAB(long_int_v.u32_arr[1], sour + 4);                                                                                                             \
            re_len = 8;                                                                                                                                                \
            break;                                                                                                                                                     \
        case ORDER_BADCFEHG:                                                                                                                                           \
            tag->data_type = 1;                                                                                                                                         \
            SET_FOR_BADC(long_int_v.u32_arr[0], sour + 4);                                                                                                             \
            SET_FOR_BADC(long_int_v.u32_arr[1], sour);                                                                                                                 \
            re_len = 8;                                                                                                                                                \
            break;                                                                                                                                                     \
        case ORDER_HGFEDCBA:                                                                                                                                           \
            tag->data_type = 1;                                                                                                                                         \
            SET_FOR_DCBA(long_int_v.u32_arr[0], sour);                                                                                                                 \
            SET_FOR_DCBA(long_int_v.u32_arr[1], sour + 4);                                                                                                             \
            re_len = 8;                                                                                                                                                \
            break;                                                                                                                                                     \
        default:                                                                                                                                                       \
            proto_syslog(LOG_ERR, "SET INT order error, order=%d", order);                                                                                             \
            break;                                                                                                                                                     \
        }                                                                                                                                                              \
    } while (0)

#define GET_UINT(tag_v, sour, order)                                     \
    do                                                                   \
    {                                                                    \
        tag_v.type = 0;                                                  \
        switch (order)                                                   \
        {                                                                \
        case ORDER_A:                                                    \
            tag_v.value.to_int = (uint8_t)GET_FOR_A(sour);               \
            break;                                                       \
        case ORDER_AB:                                                   \
            tag_v.value.to_int = (uint16_t)GET_FOR_AB(sour);             \
            break;                                                       \
        case ORDER_BA:                                                   \
            tag_v.value.to_int = (uint16_t)GET_FOR_BA(sour);             \
            break;                                                       \
        case ORDER_ABCD:                                                 \
            tag_v.value.to_int = (uint32_t)GET_FOR_ABCD(sour);           \
            break;                                                       \
        case ORDER_DCBA:                                                 \
            tag_v.value.to_int = (uint32_t)GET_FOR_DCBA(sour);           \
            break;                                                       \
        case ORDER_BADC:                                                 \
            tag_v.value.to_int = (uint32_t)GET_FOR_BADC(sour);           \
            break;                                                       \
        case ORDER_CDAB:                                                 \
            tag_v.value.to_int = (uint32_t)GET_FOR_CDAB(sour);           \
            break;                                                       \
        case ORDER_ABCDEFGH:                                             \
            tag_v.type = 1; \
            tag_v.value.to_float = (uint64_t)GET_FOR_ABCD(sour+4) | ((uint64_t)GET_FOR_ABCD(sour)<<32);       \
            break;                                                       \
        case ORDER_GHEFCDAB:                                             \
            tag_v.type = 1; \
            tag_v.value.to_float = (uint64_t)GET_FOR_CDAB(sour) | ((uint64_t)GET_FOR_CDAB(sour+4)<<32);       \
            break;                                                       \
        case ORDER_BADCFEHG:                                             \
            tag_v.type = 1; \
            tag_v.value.to_float = (uint64_t)GET_FOR_BADC(sour+4) | ((uint64_t)GET_FOR_BADC(sour)<<32);       \
            break;                                                       \
        case ORDER_HGFEDCBA:                                             \
            tag_v.type = 1; \
            tag_v.value.to_float = (uint64_t)GET_FOR_DCBA(sour) | ((uint64_t)GET_FOR_DCBA(sour+4)<<32);       \
            break;                                                       \
        default:                                                         \
            proto_syslog(LOG_ERR, "GET_UINT order error, order=%d", order); \
            break;                                                       \
        }                                                                \
    } while (0)

#define SET_UINT(tag, sour, re_len, order)                                                                                                                             \
    do                                                                                                                                                                 \
    {                                                                                                                                                                  \
        uint32_t value = 0;                                                                                                                                            \
        union uint_int                                                                                                                                                 \
        {                                                                                                                                                              \
            long long_value;                                                                                                                                           \
            uint32_t u32_arr[2];                                                                                                                                       \
        } long_uint_v = {0};                                                                                                                                                  \
        if (check_float_is_1(tag->scale))                                                                                                                              \
        {                                                                                                                                                              \
            value = long_uint_v.long_value = (tag->data_type == 0) ? (float)(write_tag->write_cache.to_int - (int)tag->offset) : write_tag->write_cache.to_float - tag->offset;                 \
        }                                                                                                                                                              \
        else                                                                                                                                                           \
        {                                                                                                                                                              \
            value = long_uint_v.long_value = (tag->data_type == 0) ? (write_tag->write_cache.to_int - tag->offset) / tag->scale : (write_tag->write_cache.to_float - tag->offset) / tag->scale; \
        }                                                                                                                                                              \
        value += 0.5;                                                                                                                                                  \
        switch (order)                                                                                                                                                 \
        {                                                                                                                                                              \
        case ORDER_A:                                                                                                                                                  \
            SET_FOR_A((uint8_t)value, sour);                                                                                                                          \
            re_len = 1;                                                                                                                                                \
            break;                                                                                                                                                     \
        case ORDER_AB:                                                                                                                                                 \
            SET_FOR_AB((uint16_t)value, sour);                                                                                                                         \
            re_len = 2;                                                                                                                                                \
            break;                                                                                                                                                     \
        case ORDER_BA:                                                                                                                                                 \
            SET_FOR_BA((uint16_t)value, sour);                                                                                                                         \
            re_len = 2;                                                                                                                                                \
            break;                                                                                                                                                     \
        case ORDER_ABCD:                                                                                                                                               \
            SET_FOR_ABCD((uint32_t)value, sour);                                                                                                                       \
            re_len = 4;                                                                                                                                                \
            break;                                                                                                                                                     \
        case ORDER_DCBA:                                                                                                                                               \
            SET_FOR_DCBA((uint32_t)value, sour);                                                                                                                       \
            re_len = 4;                                                                                                                                                \
            break;                                                                                                                                                     \
        case ORDER_BADC:                                                                                                                                               \
            SET_FOR_BADC((uint32_t)value, sour);                                                                                                                       \
            re_len = 4;                                                                                                                                                \
            break;                                                                                                                                                     \
        case ORDER_CDAB:                                                                                                                                               \
            SET_FOR_CDAB((uint32_t)value, sour);                                                                                                                       \
            re_len = 4;                                                                                                                                                \
            break;                                                                                                                                                     \
        case ORDER_ABCDEFGH:                                                                                                                                           \
            tag->data_type = 1;                                                                                                                                         \
            SET_FOR_ABCD(long_uint_v.u32_arr[0], sour + 4);                                                                                                             \
            SET_FOR_ABCD(long_uint_v.u32_arr[1], sour);                                                                                                                 \
            re_len = 8;                                                                                                                                                \
            break;                                                                                                                                                     \
        case ORDER_GHEFCDAB:                                                                                                                                           \
            tag->data_type = 1;                                                                                                                                         \
            SET_FOR_CDAB(long_uint_v.u32_arr[0], sour);                                                                                                                 \
            SET_FOR_CDAB(long_uint_v.u32_arr[1], sour + 4);                                                                                                             \
            re_len = 8;                                                                                                                                                \
            break;                                                                                                                                                     \
        case ORDER_BADCFEHG:                                                                                                                                           \
            tag->data_type = 1;                                                                                                                                         \
            SET_FOR_BADC(long_uint_v.u32_arr[0], sour + 4);                                                                                                             \
            SET_FOR_BADC(long_uint_v.u32_arr[1], sour);                                                                                                                 \
            re_len = 8;                                                                                                                                                \
            break;                                                                                                                                                     \
        case ORDER_HGFEDCBA:                                                                                                                                           \
            tag->data_type = 1;                                                                                                                                         \
            SET_FOR_DCBA(long_uint_v.u32_arr[0], sour);                                                                                                                 \
            SET_FOR_DCBA(long_uint_v.u32_arr[1], sour + 4);                                                                                                             \
            re_len = 8;                                                                                                                                                \
            break;                                                                                                                                                     \
        default:                                                                                                                                                       \
            proto_syslog(LOG_ERR, "SET UINT order error, order=%d", order);                                                                                            \
            break;                                                                                                                                                     \
        }                                                                                                                                                              \
    } while (0)

#define GET_FLOAT(tag_v, sour, order)                                        \
    do                                                                       \
    {                                                                        \
        int tmp_int = 0;                                                     \
        float *p_float = (float *)&tmp_int;                                  \
        int(*arr)[2] = (int(*)[2]) & tag_v.value.to_float;                   \
        tag_v.type = 1;                                                      \
        switch (order)                                                       \
        {                                                                    \
        case ORDER_ABCD:                                                     \
            tmp_int = GET_FOR_ABCD(sour);                                    \
            tag_v.value.to_float = *p_float;                                 \
            break;                                                           \
        case ORDER_DCBA:                                                     \
            tmp_int = GET_FOR_DCBA(sour);                                    \
            tag_v.value.to_float = *p_float;                                 \
            break;                                                           \
        case ORDER_BADC:                                                     \
            tmp_int = GET_FOR_BADC(sour);                                    \
            tag_v.value.to_float = *p_float;                                 \
            break;                                                           \
        case ORDER_CDAB:                                                     \
            tmp_int = GET_FOR_CDAB(sour);                                    \
            tag_v.value.to_float = *p_float;                                 \
            break;                                                           \
        case ORDER_ABCDEFGH:                                                 \
            (*arr)[0] = GET_FOR_ABCD(sour + 4);                              \
            (*arr)[1] = GET_FOR_ABCD(sour);                                  \
            break;                                                           \
        case ORDER_GHEFCDAB:                                                 \
            (*arr)[0] = GET_FOR_CDAB(sour);                                  \
            (*arr)[1] = GET_FOR_CDAB(sour + 4);                              \
            break;                                                           \
        case ORDER_BADCFEHG:                                                 \
            (*arr)[0] = GET_FOR_BADC(sour + 4);                              \
            (*arr)[1] = GET_FOR_BADC(sour);                                  \
            break;                                                           \
        case ORDER_HGFEDCBA:                                                 \
            (*arr)[0] = GET_FOR_DCBA(sour);                                  \
            (*arr)[1] = GET_FOR_DCBA(sour + 4);                              \
            break;                                                           \
        default:                                                             \
            proto_syslog(LOG_ERR, "GET_FLOAT order error, order=%d", order); \
            break;                                                           \
        }                                                                    \
    } while (0)

#define SET_FLOAT(tag, sour, re_len, order)                                                                                                                                                          \
    do                                                                                                                                                                                               \
    {                                                                                                                                                                                                \
        float f_value = 0;                                                                                                                                                                           \
        union double_int                                                                                                                                                                             \
        {                                                                                                                                                                                            \
            double double_value;                                                                                                                                                                     \
            uint32_t u32_arr[2];                                                                                                                                                                     \
        } double_int_v = {0};                                                                                                                                                                              \
        if (check_float_is_1(tag->scale))                                                                                                                                                            \
        {                                                                                                                                                                                            \
            f_value = double_int_v.double_value = (tag->data_type == 0) ? (float)(write_tag->write_cache.to_int - (int)tag->offset) : (write_tag->write_cache.to_float - tag->offset);               \
        }                                                                                                                                                                                            \
        else                                                                                                                                                                                         \
        {                                                                                                                                                                                            \
            f_value = double_int_v.double_value = (tag->data_type == 0) ? (write_tag->write_cache.to_int - tag->offset) / tag->scale : (write_tag->write_cache.to_float - tag->offset) / tag->scale; \
        }                                                                                                                                                                                            \
        switch (order)                                                                                                                                                                               \
        {                                                                                                                                                                                            \
        case ORDER_ABCD:                                                                                                                                                                             \
            SET_FOR_ABCD(*(uint32_t *)&f_value, sour);                                                                                                                                               \
            re_len = 4;                                                                                                                                                                              \
            break;                                                                                                                                                                                   \
        case ORDER_DCBA:                                                                                                                                                                             \
            SET_FOR_DCBA(*(uint32_t *)&f_value, sour);                                                                                                                                               \
            re_len = 4;                                                                                                                                                                              \
            break;                                                                                                                                                                                   \
        case ORDER_BADC:                                                                                                                                                                             \
            SET_FOR_BADC(*(uint32_t *)&f_value, sour);                                                                                                                                               \
            re_len = 4;                                                                                                                                                                              \
            break;                                                                                                                                                                                   \
        case ORDER_CDAB:                                                                                                                                                                             \
            SET_FOR_CDAB(*(uint32_t *)&f_value, sour);                                                                                                                                               \
            re_len = 4;                                                                                                                                                                              \
            break;                                                                                                                                                                                   \
        case ORDER_ABCDEFGH:                                                                                                                                                                         \
            SET_FOR_ABCD(double_int_v.u32_arr[0], sour + 4);                                                                                                                                         \
            SET_FOR_ABCD(double_int_v.u32_arr[1], sour);                                                                                                                                             \
            re_len = 8;                                                                                                                                                                              \
            break;                                                                                                                                                                                   \
        case ORDER_GHEFCDAB:                                                                                                                                                                         \
            SET_FOR_CDAB(double_int_v.u32_arr[0], sour);                                                                                                                                             \
            SET_FOR_CDAB(double_int_v.u32_arr[1], sour + 4);                                                                                                                                         \
            re_len = 8;                                                                                                                                                                              \
            break;                                                                                                                                                                                   \
        case ORDER_BADCFEHG:                                                                                                                                                                         \
            SET_FOR_BADC(double_int_v.u32_arr[0], sour + 4);                                                                                                                                         \
            SET_FOR_BADC(double_int_v.u32_arr[1], sour);                                                                                                                                             \
            re_len = 8;                                                                                                                                                                              \
            break;                                                                                                                                                                                   \
        case ORDER_HGFEDCBA:                                                                                                                                                                         \
            SET_FOR_DCBA(double_int_v.u32_arr[0], sour);                                                                                                                                             \
            SET_FOR_DCBA(double_int_v.u32_arr[1], sour + 4);                                                                                                                                         \
            re_len = 8;                                                                                                                                                                              \
            break;                                                                                                                                                                                   \
        default:                                                                                                                                                                                     \
            proto_syslog(LOG_ERR, "SET FLOAT order error, order=%d", order);                                                                                                                         \
            break;                                                                                                                                                                                   \
        }                                                                                                                                                                                            \
    } while (0)

#define GET_BITS(tag_v, sour, order)                                     \
    do                                                                   \
    {                                                                    \
        tag_v.type = 0;                                                  \
        switch (order)                                                   \
        {                                                                \
        case ORDER_A:                                                    \
            tag_v.value.to_int = (uint8_t)GET_FOR_A(sour);               \
            break;                                                       \
        case ORDER_AB:                                                   \
            tag_v.value.to_int = (uint16_t)GET_FOR_AB(sour);             \
            break;                                                       \
        case ORDER_BA:                                                   \
            tag_v.value.to_int = (uint16_t)GET_FOR_BA(sour);             \
            break;                                                       \
        default:                                                         \
            proto_syslog(LOG_ERR, "GET_BITS order error, order=%d", order); \
            break;                                                       \
        }                                                                \
    } while (0)

#define GET_BITS_IN_U16(data, start_bit, bit_len) ((uint16_t)((data) << (16 - start_bit - bit_len)) >> (16 - bit_len))

#define SET_SINGLE_READ_TAG(p_tag, reg_info, reg_arr, start_addr, cur_ms)                                                           \
    do                                                                                                                              \
    {                                                                                                                               \
        tag_value_t tag_value = {.type = 0};                                                                                        \
        int index = (reg_info->reg_addr - start_addr) * 2;                                                                          \
        switch (reg_info->data_type)                                                                                                \
        {                                                                                                                           \
        case TYPE_BOOL:                                                                                                             \
            tag_value.type = 0;                                                                                                     \
            if (reg_info->reversal != 0)                                                                                            \
            {                                                                                                                       \
                tag_value.value.to_int = (*(reg_arr + (reg_info->reg_addr - start_addr)) != 0) ? 0 : 1;                             \
            }                                                                                                                       \
            else                                                                                                                    \
            {                                                                                                                       \
                tag_value.value.to_int = (*(reg_arr + reg_info->reg_addr - start_addr) != 0) ? 1 : 0;                               \
            }                                                                                                                       \
            proto_set_tag(p_tag, tag_value, cur_ms);                                                                                \
            break;                                                                                                                  \
                                                                                                                                    \
        case TYPE_INT:                                                                                                              \
            GET_INT(tag_value, reg_arr + index, reg_info->data_order);                                                              \
            proto_set_tag(p_tag, tag_value, cur_ms);                                                                                \
            break;                                                                                                                  \
                                                                                                                                    \
        case TYPE_UINT:                                                                                                             \
            GET_UINT(tag_value, reg_arr + index, reg_info->data_order);                                                             \
            proto_set_tag(p_tag, tag_value, cur_ms);                                                                                \
            break;                                                                                                                  \
                                                                                                                                    \
        case TYPE_FLOAT:                                                                                                            \
            GET_FLOAT(tag_value, reg_arr + index, reg_info->data_order);                                                            \
            proto_set_tag(p_tag, tag_value, cur_ms);                                                                                \
            break;                                                                                                                  \
                                                                                                                                    \
        case TYPE_BITS:                                                                                                             \
            GET_BITS(tag_value, reg_arr + index, reg_info->data_order);                                                             \
            if (reg_info->reversal != 0)                                                                                            \
            {                                                                                                                       \
                tag_value.value.to_int = GET_BITS_IN_U16(~tag_value.value.to_int, reg_info->bit_st, reg_info->bit_len);             \
            }                                                                                                                       \
            else                                                                                                                    \
            {                                                                                                                       \
                tag_value.value.to_int = GET_BITS_IN_U16(tag_value.value.to_int, reg_info->bit_st, reg_info->bit_len);              \
            }                                                                                                                       \
            proto_set_tag(p_tag, tag_value, cur_ms);                                                                                \
            break;                                                                                                                  \
                                                                                                                                    \
        default:                                                                                                                    \
            proto_syslog(LOG_ERR, "extract data error, date type isn't supported, data_type = %d, tag=%s, fun_code:%d, start_addr:%d", \
                      reg_info->data_type, p_tag->name, reg_info->func_code, reg_info->reg_addr);                                   \
            break;                                                                                                                  \
        }                                                                                                                           \
    } while (0);

#define SET_SINGLE_WRITE_TAG(p_tag, reg_info, reg_arr, start_addr, cur_ms)                                                          \
    do                                                                                                                              \
    {                                                                                                                               \
        tag_value_t tag_value = {.type = 0};                                                                                        \
        int index = (reg_info->reg_addr - start_addr) * 2;                                                                          \
        switch (reg_info->data_type)                                                                                                \
        {                                                                                                                           \
        case TYPE_BOOL:                                                                                                             \
            tag_value.type = 0;                                                                                                     \
            tag_value.value.to_int = (*(reg_arr + reg_info->reg_addr - start_addr) != 0) ? 1 : 0;                                   \
            proto_set_tag(p_tag, tag_value, cur_ms);                                                                                \
            break;                                                                                                                  \
        case TYPE_INT:                                                                                                              \
            GET_INT(tag_value, reg_arr + index, reg_info->data_order);                                                              \
            proto_set_tag(p_tag, tag_value, cur_ms);                                                                                \
            break;                                                                                                                  \
                                                                                                                                    \
        case TYPE_UINT:                                                                                                             \
            GET_UINT(tag_value, reg_arr + index, reg_info->data_order);                                                             \
            proto_set_tag(p_tag, tag_value, cur_ms);                                                                                \
            break;                                                                                                                  \
                                                                                                                                    \
        case TYPE_FLOAT:                                                                                                            \
            GET_FLOAT(tag_value, reg_arr + index, reg_info->data_order);                                                            \
            proto_set_tag(p_tag, tag_value, cur_ms);                                                                                \
            break;                                                                                                                  \
                                                                                                                                    \
        default:                                                                                                                    \
            proto_syslog(LOG_ERR, "extract data error, date type isn't supported, data_type = %d, tag=%s, fun_code:%d, start_addr:%d", \
                      reg_info->data_type, p_tag->name, reg_info->func_code, reg_info->reg_addr);                                   \
            break;                                                                                                                  \
        }                                                                                                                           \
    } while (0);

void channel_lock(channel_t *channel)
{
    pthread_mutex_lock(&channel->channel_lock);
}

void channel_unlock(channel_t *channel)
{
    pthread_mutex_unlock(&channel->channel_lock);
}

static int check_tag_state(channel_t *channel)
{
    #if 1
    long long cur_ms = clock_get_ms();
    long long data_ms = 0;
    #endif
    for (size_t i = 0; i < channel->devs_size; i++)
    {
        //只读的
        for (size_t j = 0; j < channel->devs[i]->tags.ro_tags_size; j++)
        {
            if (channel->devs[i]->tags.ro_tags[j]->force_set_cache.valid == true)
            {
                channel->devs[i]->tags.ro_tags[j]->state_code = TAG_STATE_GOOD;
                continue;
            }
            #if 1
            if((int8_t)channel->devs[i]->tags.ro_tags[j]->poll_list_type == -1) 
                channel->devs[i]->tags.ro_tags[j]->state_code = TAG_STATE_BAD;
            else if(channel->devs[i]->tags.ro_tags[j]->poll_list_type == 0)
                data_ms = 120*1000;
                

            if((cur_ms - channel->devs[i]->tags.ro_tags[j]->read_cache_time)>data_ms)
                channel->devs[i]->tags.ro_tags[j]->state_code = TAG_STATE_BAD;
            else
                channel->devs[i]->tags.ro_tags[j]->state_code = TAG_STATE_GOOD;
  
            #else
            channel->devs[i]->tags.ro_tags[j]->state_code = TAG_STATE_BAD;
            #endif
          
            if(channel->devs[i]->tags.ro_tags[j]->state_code!=channel->devs[i]->tags.ro_tags[j]->last_state_code)
            {
                channel->devs[i]->tags.ro_tags[j]->last_state_code = channel->devs[i]->tags.ro_tags[j]->state_code;
                // reflush_data_up(channel->devs[i]->tags.ro_tags[j]); // 不用调回调了
            }
        }
        //可读可写的
        for (size_t j = 0; j < channel->devs[i]->tags.rw_tags_size; j++)
        {
            if (channel->devs[i]->tags.rw_tags[j]->force_set_cache.valid == true)
            {
                channel->devs[i]->tags.rw_tags[j]->state_code = TAG_STATE_GOOD;
                continue;
            }
            #if 1
            if((int8_t)channel->devs[i]->tags.rw_tags[j]->poll_list_type == -1) 
                channel->devs[i]->tags.rw_tags[j]->state_code = TAG_STATE_BAD;
            else if(channel->devs[i]->tags.rw_tags[j]->poll_list_type == 0)
                data_ms = 120*1000;

            if((cur_ms - channel->devs[i]->tags.rw_tags[j]->read_cache_time)>data_ms)
                channel->devs[i]->tags.rw_tags[j]->state_code = TAG_STATE_BAD;
            else
                channel->devs[i]->tags.rw_tags[j]->state_code = TAG_STATE_GOOD;
            #else
            channel->devs[i]->tags.rw_tags[j]->state_code = TAG_STATE_BAD;
            #endif
            if(channel->devs[i]->tags.rw_tags[j]->state_code!=channel->devs[i]->tags.rw_tags[j]->last_state_code)
            {
                channel->devs[i]->tags.rw_tags[j]->last_state_code = channel->devs[i]->tags.rw_tags[j]->state_code;
                // reflush_data_up(channel->devs[i]->tags.rw_tags[j]); // 不用调回调了
            }
        }
    }
    return 0;
}

//调用这个函数的时候如果是bool类型也要计算倍率和偏移，所以必须保证bool的倍率为1，偏移为0
int proto_set_tag(tag_t *tag, tag_value_t tag_v, uint64_t cur_ms)
{
    if (tag->force_set_cache.valid == false)
    {
        tag_value_t pval;
        tag_handle_fun funptr = tag->funptr;
        pval.type = tag->data_type;
        pval.value.to_float = 0;
        pval.value_last.to_float = 0;

        if (tag->data_type == 0)//存的类型
        {
            tag->last_read_cache.to_int = tag->read_cache.to_int;//更新last_read_cache
            if (fabsf(tag->scale - 1.0) < 0.0001)
            {
                pval.value.to_int = ((tag_v.type == 0) ? tag_v.value.to_int : tag_v.value.to_float) + (int)tag->offset;
            }
            else
            {
                pval.value.to_int = ((tag_v.type == 0) ? tag_v.value.to_int : tag_v.value.to_float) * tag->scale + tag->offset;
            }

            if (funptr != NULL)
                tag->read_cache.to_int = funptr(&pval);
            else
                tag->read_cache.to_int = pval.value.to_int;
             //proto_syslog(LOG_INFO, "TAG:%s->%d", tag->name, tag->read_cache.to_int);
        }
        else
        {
            tag->last_read_cache.to_float = tag->read_cache.to_float;//更新last_read_cache
            pval.value.to_float = ((tag_v.type == 0) ? tag_v.value.to_int : tag_v.value.to_float) * tag->scale + tag->offset;
            if (funptr != NULL)
                tag->read_cache.to_float = (double) funptr(&pval);
            else
                tag->read_cache.to_float = pval.value.to_float;
             //proto_syslog(LOG_INFO, "TAG:%s->%f", tag->name, tag->read_cache.to_float);
        }
        
    }
    else
    {
        if (tag->data_type == 0)//存的类型
        {
            tag->last_read_cache.to_int = tag->read_cache.to_int;
            tag->read_cache.to_int = tag->force_set_cache.value.to_int;
        }
        else
        {
            tag->last_read_cache.to_float = tag->read_cache.to_float;
            tag->read_cache.to_float = tag->force_set_cache.value.to_float;
        }
    }
    tag->last_read_time = tag->read_cache_time;
    tag->read_cache_time = cur_ms;
    reflush_data_up(tag);
    return 0;
}

int proto_set_tag_no_scale_offset(tag_t *tag, tag_value_t tag_v, uint64_t cur_ms)
{
    if (tag->force_set_cache.valid == false)
    {
        tag_value_t pval;
        tag_handle_fun funptr = tag->funptr;
        pval.type = tag->data_type;
        pval.value.to_float = 0;
        pval.value_last.to_float = 0;

        if (tag->data_type == 0)
        {
            tag->last_read_cache.to_int = tag->read_cache.to_int; // 更新last_read_cache
            pval.value.to_int = ((tag_v.type == 0) ? tag_v.value.to_int : tag_v.value.to_float);
            if (funptr != NULL)
                tag->read_cache.to_int = funptr(&pval);
            else
                tag->read_cache.to_int = pval.value.to_int;
        }
        else
        {
            tag->last_read_cache.to_float = tag->read_cache.to_float; // 更新last_read_cache
            pval.value.to_float = ((tag_v.type == 0) ? tag_v.value.to_int : tag_v.value.to_float);
            if (funptr != NULL)
                tag->read_cache.to_float = (double)funptr(&pval);
            else
                tag->read_cache.to_float = pval.value.to_float;
        }
    }
    else
    {
        if (tag->data_type == 0)
        {
            tag->last_read_cache.to_int = tag->read_cache.to_int;
            tag->read_cache.to_int = tag->force_set_cache.value.to_int;
        }
        else
        {
            tag->last_read_cache.to_float = tag->read_cache.to_float;
            tag->read_cache.to_float = tag->force_set_cache.value.to_float;
        }
    }
    tag->last_read_time = tag->read_cache_time;
    tag->read_cache_time = cur_ms;
    reflush_data_up(tag);
    return 0;
}

int modbus_update_tags(device_t *dev, uint8_t fun_code, uint16_t start_addr, void *buff, int buff_len)
{
    tag_t *p_tag = NULL;
    modbus_read_reg_info_t *read_reg = NULL;
    uint16_t end_addr =  ((fun_code == 0x01) || (fun_code == 0x02)) ? start_addr + buff_len : start_addr + (buff_len >> 1);
    uint8_t *reg_arr = (uint8_t *)buff;

    uint64_t cur_ms = clock_get_ms();
    for (size_t i = 0; i < dev->tags.ro_tags_size; i++)
    {
        read_reg = (modbus_read_reg_info_t *)dev->tags.ro_tags[i]->proto_ptr;
        p_tag = dev->tags.ro_tags[i];
        if ((read_reg->func_code == fun_code) && (read_reg->reg_addr >= start_addr) && (read_reg->reg_addr < end_addr))
        {
            SET_SINGLE_READ_TAG(p_tag, read_reg, reg_arr, start_addr, cur_ms);
        }
    }

    for (size_t i = 0; i < dev->tags.rw_tags_size; i++)
    {
        modbus_write_reg_info_t *write_reg = (modbus_write_reg_info_t *)dev->tags.rw_tags[i]->proto_ptr;
        p_tag = dev->tags.rw_tags[i];
        if ((write_reg->func_code == fun_code) && (write_reg->reg_addr >= start_addr) && (write_reg->reg_addr < end_addr))
        {
            SET_SINGLE_WRITE_TAG(p_tag, write_reg, reg_arr, start_addr, cur_ms);
        }
    }
    return 0;
}

int modbus_update_tags_new(modbus_poll_info_t *poll_info, void *buff, int buff_len)
{
    tag_t *p_tag = NULL;
    modbus_read_reg_info_t *read_reg = NULL;
    modbus_write_reg_info_t *write_reg = NULL;
    uint16_t start_addr = poll_info->reg_addr;
    uint8_t *reg_arr = (uint8_t *)buff;

    uint64_t cur_ms = clock_get_ms();
    for (int t = 0; t < poll_info->tag_sum; t++)
    {
        read_reg = (modbus_read_reg_info_t *)poll_info->p_tag[t]->proto_ptr;
        p_tag = poll_info->p_tag[t];

        if (read_reg->head != 0)
        {
            SET_SINGLE_READ_TAG(p_tag, read_reg, reg_arr, start_addr, cur_ms);
        }
        else
        {
            write_reg = (modbus_write_reg_info_t *)poll_info->p_tag[t]->proto_ptr;
            SET_SINGLE_WRITE_TAG(p_tag, write_reg, reg_arr, start_addr, cur_ms);
        }
    }
    return 0;
}
static int tag_2_modbus_data(tag_t *write_tag, modbus_read_reg_info_t *write_reg, uint8_t *dest_buff)
{
    int use_len = 0;
    switch (write_reg->data_type)
    {
    case TYPE_BOOL:
        *dest_buff = (write_tag->write_cache.to_int != 0) ? 1 : 0;
        use_len = 1;
        break;

    case TYPE_INT:
        SET_INT(write_tag, dest_buff, use_len, write_reg->data_order);
        break;

    case TYPE_UINT:
        SET_UINT(write_tag, dest_buff, use_len, write_reg->data_order);
        break;

    case TYPE_FLOAT:
        SET_FLOAT(write_tag, dest_buff, use_len, write_reg->data_order);
        break;

    default:
        proto_syslog(LOG_ERR, "write reg type error, [funcode:%d, addr:0x%x, type=%d]", write_reg->func_code, write_reg->reg_addr, write_reg->data_type);
        break;
    }

    return use_len;
}

void *HistoryCycleTask(void *data)
{
    UA_SERVER_PARAM *ua_param = (UA_SERVER_PARAM *)data;
    time_t ts = 0;
    while (1)
    {
        ts = time(NULL) - get_pcs_ctrl_var()->history_data_rolling * 60 * 60; // 前history_data_rolling小时的时间戳
        proto_syslog(LOG_NOTICE, "sqlite3 delete by time:%lu", ts);
        sqlite3_get_table_list_and_del_by_time((sqlite3 *)ua_param->setting.historizingBackend.context, ts);
        sleep(1);
    }
    return NULL;
}

#if 1 //各个设备适配器的回调函数
static inline uint8_t to_bcd(uint8_t dat)
{
    return (((dat / 10) << 4) | (dat % 10));
}

int meter_at_at182g_adp_handle(tag_t *tag, void *p)
{
    int re = -99;
    device_t *dev = (device_t *)tag->dev_ptr;
    if (tag->write_cache.to_int < MIN_TIMESTAMP)
    {
        proto_syslog(LOG_ERR, "dev:%s, synchronous time fail, timestamp error", dev->no);
        return -100;
    }

    time_t now = tag->write_cache.to_int; // time(NULL);
    switch (dev->proto_type)
    {
    case PROTOCOL_TYPE_MODBUS_TCP:
    case PROTOCOL_TYPE_MODBUS_RTU:
    case PROTOCOL_TYPE_MODBUS_RTU_RCRC:
    {
        struct tm target_time = {};  
        localtime_r(&now, &target_time);
        struct tm *local = &target_time;
        uint8_t date_time[8] = {};
        int year = local->tm_year + 1900;
        date_time[0] = to_bcd(year / 100);
        date_time[1] = to_bcd(year % 100);
        date_time[2] = to_bcd(local->tm_mon + 1);
        date_time[3] = to_bcd(local->tm_mday);
        date_time[4] = (local->tm_wday == 0) ? 7 : local->tm_wday;
        date_time[5] = to_bcd(local->tm_hour);
        date_time[6] = to_bcd(local->tm_min);
        date_time[7] = to_bcd(local->tm_sec);
        re = modbus_write_registers((modbus_t *)p, 0x501a, sizeof(date_time) / 2, (const uint16_t *)date_time);
        if (sizeof(date_time) / 2 != re)
        {
            proto_syslog(LOG_ERR, "dev:%s, synchronous time fail, ERROR:%s, errno=%d", dev->no, modbus_strerror(errno), errno);
            re = -1;
        }
        else
        {
            re = 0;
            proto_syslog(LOG_NOTICE, "dev:%s,synchronous time success", dev->no);
        }
    }
    break;

    case PROTOCOL_TYPE_DLT645_2007:
    case PROTOCOL_TYPE_DLT645_1997:
        re = dlt645_sync_date_time((dlt645_hd_t *)p, now);
        if (re == 0)
        {
            proto_syslog(LOG_NOTICE, "dev:%s,synchronous time success", dev->no);
        }
        break;

    default:
        proto_syslog(LOG_NOTICE, "dev[%s] adapter handle protocol is not supported", dev->no);
        re = -88;
    }
    return re;
}

int meter_acrel_dtsd1352_adp_handle(tag_t *tag, void *p)
{
    int re = -99;
    device_t *dev = (device_t *)tag->dev_ptr;
    if (tag->write_cache.to_int < MIN_TIMESTAMP)
    {
        proto_syslog(LOG_ERR, "dev:%s, synchronous time fail, timestamp error", dev->no);
        return -100;
    }
    time_t now = tag->write_cache.to_int; // time(NULL);
    switch (dev->proto_type)
    {
    case PROTOCOL_TYPE_MODBUS_TCP:
    case PROTOCOL_TYPE_MODBUS_RTU:
    case PROTOCOL_TYPE_MODBUS_RTU_RCRC:
    {
        struct tm target_time = {};  
        localtime_r(&now, &target_time);
        struct tm *local = &target_time;
        uint8_t date_time[6] = {};
        int year = local->tm_year + 1900;
        date_time[5] = year % 100;
        date_time[4] = local->tm_mon + 1;
        date_time[3] = local->tm_mday;
        date_time[2] = local->tm_hour;
        date_time[1] = local->tm_min;
        date_time[0] = local->tm_sec;
        re = modbus_write_registers((modbus_t *)p, 0x501a, sizeof(date_time) / 2, (const uint16_t *)date_time);
        if (sizeof(date_time) / 2 != re)
        {
            proto_syslog(LOG_ERR, "dev:%s, synchronous time fail, ERROR:%s, errno=%d", dev->no, modbus_strerror(errno), errno);
            re = -1;
        }
        else
        {
            re = 0;
            proto_syslog(LOG_NOTICE, "dev:%s,synchronous time success", dev->no);
        }
    }
    break;

    case PROTOCOL_TYPE_DLT645_2007:
    case PROTOCOL_TYPE_DLT645_1997:
        re = dlt645_sync_date_time((dlt645_hd_t *)p, now);
        if (re == 0)
        {
            proto_syslog(LOG_NOTICE, "dev:%s,synchronous time success", dev->no);
        }
        break;

    default:
        proto_syslog(LOG_NOTICE, "dev[%s] adapter handle protocol is not supported", dev->no);
        re = -88;
    }
    return re;
}

int meter_dlt645_common_adp_handle(tag_t *tag, void *p)
{
    int re = -99;
    device_t *dev = (device_t *)tag->dev_ptr;
    if (tag->write_cache.to_int < MIN_TIMESTAMP)
    {
        proto_syslog(LOG_ERR, "dev:%s, synchronous time fail, timestamp error", dev->no);
        return -100;
    }
    time_t now = tag->write_cache.to_int; // time(NULL);
    switch (dev->proto_type)
    {
    case PROTOCOL_TYPE_DLT645_2007:
    case PROTOCOL_TYPE_DLT645_1997:
        re = dlt645_sync_date_time((dlt645_hd_t *)p, now);
        if (re == 0)
        {
            proto_syslog(LOG_NOTICE, "dev:%s,synchronous time success", dev->no);
        }
        break;

    default:
        proto_syslog(LOG_NOTICE, "dev[%s] adapter handle protocol is not supported", dev->no);
        re = -88;
    }
    return re;
}

#endif

static int modbus_write_tag(modbus_t *ctx, tag_t *write_tag)
{
    int re = -1;
    uint8_t reg_buff[16] = {0};
    int len = 0;

    modbus_read_reg_info_t *write_reg = (modbus_read_reg_info_t *)write_tag->proto_ptr;
    device_t *dev = (device_t *)write_tag->dev_ptr;
    modbus_set_slave(ctx, *((int *)dev->proto_ptr));
    channel_lock(write_tag->chan_ptr);
    if (write_reg  != NULL)
    {
        len = tag_2_modbus_data(write_tag, write_reg, reg_buff);

        switch (write_reg->func_code)
        {
        case MODBUS_FC_READ_COILS:
        case MODBUS_FC_WRITE_SINGLE_COIL:
            if (len == 1)
            {
                proto_syslog(LOG_NOTICE, "dev:%s, modbus_write_bit, [tag:%s], addr = %d, value = %d",
                             dev->no, write_tag->name, write_reg->reg_addr, reg_buff[0]);
                re = modbus_write_bit(ctx, write_reg->reg_addr, reg_buff[0]);
                if (re != 1)
                {
                    proto_syslog(LOG_ERR, "dev:%s, modbus write bit error, [tag:%s], ret = %d, addr = %d, ERROR:%s, errno=%d",
                                 dev->no, write_tag->name, re, write_reg->reg_addr, modbus_strerror(errno), errno);
                    BUSINESS_LOG(BLOG_ERR, DAS_ID_MODBUS_7, dev->no, "[数采]: 设备(%s) 写线圈失败, tag: %s, value: %d, 错误原因(%s)", dev->no, write_tag->name, reg_buff[0], modbus_strerror(errno));
                }
                else
                {
                    re = 0;
                    proto_syslog(LOG_NOTICE, "dev:%s, modbus write bit success, [tag:%s]", dev->no, write_tag->name);
                    BUSINESS_LOG(BLOG_NOTICE, DAS_ID_MODBUS_7, dev->no, "[数采]: 设备(%s) 写线圈成功, tag: %s, value: %d", dev->no, write_tag->name, reg_buff[0]);
                }
            }
            else
            {
                proto_syslog(LOG_ERR, "dev:%s, tag extract error [tag:%s], len = %d",
                             dev->no, write_tag->name, len);
            }
            break;

        case MODBUS_FC_READ_HOLDING_REGISTERS:
            len = len >> 1; // 转为寄存器个数
            if (len > 0)
            {
                proto_syslog(LOG_NOTICE, "dev:%s, modbus write registers, [tag:%s], addr = %d",
                             dev->no, write_tag->name, write_reg->reg_addr);
                if (len == 1)
                {
                    int regv = (reg_buff[1] << 8) | reg_buff[0];
                    proto_syslog(LOG_NOTICE, "val:%d",regv);
                    re = modbus_write_register(ctx, write_reg->reg_addr, regv); 
                    
                }
                else
                {
                    re = modbus_write_registers(ctx, write_reg->reg_addr, len, (const uint16_t *)reg_buff); // 直接使用0x10命令码，为了支持浮点数一次写入
                }
                if (len != re)
                {
                    proto_syslog(LOG_ERR, "dev:%s, modbus write registers error, [tag:%s], ret = %d, addr = %d, ERROR:%s, errno=%d",
                                dev->no, write_tag->name, re, write_reg->reg_addr, modbus_strerror(errno), errno);
                    BUSINESS_LOG(BLOG_ERR, DAS_ID_MODBUS_7, dev->no, "[数采]: 设备(%s) 写寄存器失败, tag: %s, addr: %d, 错误原因(%s)", dev->no, write_tag->name, write_reg->reg_addr, modbus_strerror(errno));
                }
                else
                {
                    re = 0;
                    proto_syslog(LOG_NOTICE, "dev:%s, modbus write registers success, [tag:%s]", dev->no, write_tag->name);
                    BUSINESS_LOG(BLOG_NOTICE, DAS_ID_MODBUS_7, dev->no, "[数采]: 设备(%s) 写寄存器成功, tag: %s, addr: %d", dev->no, write_tag->name, write_reg->reg_addr);
                }
            }
            else
            {
                proto_syslog(LOG_ERR, "dev:%s, tag extract error [tag:%s], len = %d",
                             dev->no, write_tag->name, len);
            }
            break;
        case MODBUS_FC_WRITE_SINGLE_REGISTER:
            len = len >> 1; // 转为寄存器个数
            if (len == 1)
            {
                proto_syslog(LOG_NOTICE, "dev:%s, modbus write registers, [tag:%s], addr = %d",
                             dev->no, write_tag->name, write_reg->reg_addr);
                int regv = (reg_buff[1] << 8) | reg_buff[0];
                proto_syslog(LOG_NOTICE, "val:%d", regv);
                re = modbus_write_register(ctx, write_reg->reg_addr, regv);
                if (len != re)
                {
                    proto_syslog(LOG_ERR, "dev:%s, modbus write registers error, [tag:%s], ret = %d, addr = %d, ERROR:%s, errno=%d",
                                 dev->no, write_tag->name, re, write_reg->reg_addr, modbus_strerror(errno), errno);
                    BUSINESS_LOG(BLOG_ERR, DAS_ID_MODBUS_7, dev->no, "[数采]: 设备(%s) 写寄存器失败, tag: %s, addr: %d, 错误原因(%s)", dev->no, write_tag->name, write_reg->reg_addr, modbus_strerror(errno));
                }
                else
                {
                    re = 0;
                    proto_syslog(LOG_NOTICE, "dev:%s, modbus write registers success, [tag:%s]", dev->no, write_tag->name);
                    BUSINESS_LOG(BLOG_ERR, DAS_ID_MODBUS_7, dev->no, "[数采]: 设备(%s) 写寄存器成功, tag: %s, addr: %d", dev->no, write_tag->name, write_reg->reg_addr);
                }
            }
            else
            {
                proto_syslog(LOG_ERR, "dev:%s, tag extract error [tag:%s], len = %d",
                             dev->no, write_tag->name, len);
            }
            break;
        case MODBUS_FC_WRITE_MULTIPLE_REGISTERS:
            len = len >> 1; // 转为寄存器个数
            if (len > 0)
            {
                proto_syslog(LOG_NOTICE, "dev:%s, modbus write registers, [tag:%s], addr = %d",
                             dev->no, write_tag->name, write_reg->reg_addr);
                re = modbus_write_registers(ctx, write_reg->reg_addr, len, (const uint16_t *)reg_buff);
                if (len != re)
                {
                    proto_syslog(LOG_ERR, "dev:%s, modbus write registers error, [tag:%s], ret = %d, addr = %d, ERROR:%s, errno=%d",
                                dev->no, write_tag->name, re, write_reg->reg_addr, modbus_strerror(errno), errno);
                    BUSINESS_LOG(BLOG_ERR, DAS_ID_MODBUS_7, dev->no, "[数采]: 设备(%s) 写寄存器失败, tag: %s, addr: %d, 错误原因(%s)", dev->no, write_tag->name, write_reg->reg_addr, modbus_strerror(errno));
                }
                else
                {
                    re = 0;
                    proto_syslog(LOG_NOTICE, "dev:%s, modbus write registers success, [tag:%s]", dev->no, write_tag->name);
                    BUSINESS_LOG(BLOG_ERR, DAS_ID_MODBUS_7, dev->no, "[数采]: 设备(%s) 写寄存器成功, tag: %s, addr: %d", dev->no, write_tag->name, write_reg->reg_addr);
                }
            }
            else
            {
                proto_syslog(LOG_ERR, "dev:%s, tag extract error [tag:%s], len = %d",
                             dev->no, write_tag->name, len);
            }
            break;

        default:
            proto_syslog(LOG_ERR, "This funcode(0x%02x) isn't supported write", write_reg->func_code);
            BUSINESS_LOG(BLOG_ERR, DAS_ID_MODBUS_7, NULL, "[数采]: 功能码 %02x 不支持写");
            re = -1;
            break;
        }
    }
    else
    {
        re = (dev->adp != NULL) ? dev->adp->handle(write_tag, ctx) : -3;
    }
    channel_unlock(write_tag->chan_ptr);
    return re;
}

static void sync_wtag_read_cache(write_list_t *w_info)
{
    tag_t *write_tag = &w_info->w_tag;
    tag_t *psour = w_info->pw_tag;
    // write_tag同步当前的数据
    write_tag->last_read_cache.to_float = psour->read_cache.to_float;
    write_tag->read_cache.to_float = write_tag->write_cache.to_float;
    // 将数据在备份回来到psour
    psour->read_cache.to_float = write_tag->read_cache.to_float;
    psour->last_read_cache.to_float = write_tag->last_read_cache.to_float;
    // 使用write_tag取回调函数
    reflush_data_up(write_tag);
    psour->inited = write_tag->inited;
    //
    psour->read_cache_time = clock_get_ms();
}

static int find_tag_id_by_dev_and_node(device_t *dev, char *node_name)
{
    for(int i = 1; i < dev->tags.sys_info_tags_size; i++)
    {
        if(strcmp(dev->tags.sys_info_tags[i]->name, node_name) == 0)
        {
            return i;
        }
    }
    return -1;
}



static int check_write(modbus_t *ctx, channel_t *channel, int time_out)
{
    uint8_t buff[256];
    int rx_len = 0;
    device_t *dev = NULL;
    int ret = sem_wait_time(&channel->write_sem, 0, time_out);
    if (ret == 0)
    {
        do
        {
            write_list_t *w_info = get_data_from_list(channel);
            if (w_info == NULL)
            {
                proto_syslog(LOG_ERR, "get data from list is null");
                return -1;
            }
            tag_t *write_tag = &w_info->w_tag;
            if (write_tag == NULL)
            {
                proto_syslog(LOG_ERR, "get data from list is null");
                free(w_info);
                return -2;
            }
            dev = (device_t *)write_tag->dev_ptr;
            if (dev->write_gap > 0)
            {
                usleep(dev->write_gap * 1000);
            }
            proto_syslog(LOG_INFO, "get data from list is [%s.%s(type:%d, double=%f, int=%ld)]",
                         dev->no, write_tag->name, write_tag->data_type, write_tag->write_cache.to_float, write_tag->write_cache.to_int);
            rx_len = modbus_write_tag(ctx, write_tag);
            if (0 != rx_len)
            {
                write_tag->last_state_code = write_tag->state_code;
                write_tag->state_code = TAG_STATE_BAD;
                int tag_id = find_tag_id_by_dev_and_node(dev, write_tag->name);
                if(tag_id != -1)
                    refresh_special_state(dev, write_tag->write_cache.to_int, tag_id);
                check_online_state(dev);
                proto_syslog(LOG_ERR, "modbus write tag error, tag:%s, ret = %d", write_tag->name, rx_len);
                BUSINESS_LOG(BLOG_ERR, DAS_ID_MODBUS_5, write_tag->dev_tag_name, "[数采]: %s 通道写tag: %s 失败,", channel->channel, write_tag->name);
                free(w_info);
                return -3;
            }
            else
            {
                BUSINESS_LOG(BLOG_NOTICE, DAS_ID_MODBUS_5, write_tag->dev_tag_name, "[数采]: %s 通道写tag: %s 成功,", channel->channel, write_tag->name);
                write_tag->last_state_code = write_tag->state_code;
                write_tag->state_code = TAG_STATE_GOOD;
                if (write_tag->no_fetch == 0) // 需要回读
                {
                    proto_syslog(LOG_DEBUG, "fetch [channel:%s][dev:%s]", channel->channel, dev->no);
                    modbus_read_reg_info_t *write_reg = (modbus_read_reg_info_t *)write_tag->proto_ptr;
                    switch (write_reg->data_order)
                    {
                    case ORDER_A:
                    case ORDER_AB:
                    case ORDER_BA:
                        rx_len = 1;
                        break;
                    case ORDER_ABCD:
                    case ORDER_DCBA:
                    case ORDER_BADC:
                    case ORDER_CDAB:
                        rx_len = 2;
                        break;
                    case ORDER_ABCDEFGH:
                    case ORDER_GHEFCDAB:
                    case ORDER_BADCFEHG:
                    case ORDER_HGFEDCBA:
                        rx_len = 4;
                        break;

                    default:
                        break;
                    }
                    channel_lock(write_tag->chan_ptr);
                    rx_len = modbus_proto_read(ctx, *((int *)dev->proto_ptr), write_reg->func_code, write_reg->reg_addr, rx_len, buff); // 读数据
                    channel_unlock(write_tag->chan_ptr);
                    if (rx_len > 0)
                    {
                        // proto_syslog(LOG_NOTICE, "[channel:%s][dev:%s] modbus TAGS:", channel->channel, dev->no);
                        uint64_t cur_ms = clock_get_ms();
                        SET_SINGLE_WRITE_TAG(w_info->pw_tag, write_reg, buff, write_reg->reg_addr, cur_ms);
                        // modbus_update_tags(dev, poll_info->func_code, poll_info->reg_addr, buff, rx_len);//根据返回数据更新tag
                        refresh_online_state(dev);
                        BUSINESS_LOG(BLOG_ERR, DAS_ID_MODBUS_6, dev->no, "[数采]: %s 通道 %s 设备回读正常", channel->channel, dev->no);
                    }
                    else
                    {
                        check_online_state(dev);
                        proto_syslog(LOG_ERR, "[channel:%s][dev:%s] :modbus read error(rx_len = %d), ERROR info:%s, errno=%d", channel->channel, dev->no, rx_len, modbus_strerror(errno), errno);
                        BUSINESS_LOG(BLOG_ERR, DAS_ID_MODBUS_6, dev->no, "[数采]: %s 通道 %s 设备回读失败(无回复数据), 失败原因(%s)", channel->channel, dev->no, modbus_strerror(errno));
                        return -4;
                    }
                }
                else
                {
                    sync_wtag_read_cache(w_info);
                }
            }
            free(w_info);
        } while (sem_wait_time(&channel->write_sem, 0, 0) == 0);
    }

    return 0;
}

void setmodebustimeout(modbus_t *ctx, int timeout_ms)
{
    modbus_set_response_timeout(ctx, timeout_ms / 1000, (timeout_ms % 1000) * 1000);
}

static int modbus_broadcast(RAW_SOCK *raw, uint8_t dev_addr, modbus_poll_info_t *poll_info, void *buff, int len)
{
    struct
    {
        uint16_t poll_index; // 模板poll中的位置
        uint16_t data_len;   // 数据长度
        uint8_t dev_addr;    // 设备地址
        uint8_t data[256];   // 数据
    } send_pack = {
        // 没有校验，传输通道保证了数据准确性
        .poll_index = poll_info->index,
        .data_len = len,
        .dev_addr = dev_addr,
    };

    memcpy(send_pack.data, buff, len);
    // proto_syslog_hex(LOG_ERR, &send_pack, len + 5, "send_pack");
    return send_protocol_data(raw, (char *)&send_pack, len + 5);
}
#ifdef EN_EMS
static int contain_meter(channel_t *channel)
{
    device_t *dev = NULL;
    for (size_t i = 0; i < channel->chunks_interv_size; i++)
    {
        dev = channel->chunks_interv[i]->dev_ptr;
        if ((dev->is_meter != 0) && (dev->dev_child_type != NULL) && ((strcmp(dev->dev_child_type, "GRID") == 0) || (strcmp(dev->dev_child_type, "TRAS_LV") == 0))) // 是电表
        {
            return 1;
        }
    }

    return 0;
}
#endif
static void iec104_ConnectionHandler(void * parameter, CS104_Connection connection, CS104_ConnectionEvent event)
{
    channel_t * chan;
    iec104_proto * iecp;
    const char * msg = "unknown";

    chan = (channel_t *) parameter;
    iecp = (chan != NULL) ? (iec104_proto *) chan->proto_ptr : NULL;
    if (iecp == NULL)
		return;

	pthread_mutex_lock(&iecp->connect_lock);
    switch (event) {
    case CS104_CONNECTION_OPENED:
        iecp->connected = 1;
        msg = "connected";
        break;

    case CS104_CONNECTION_CLOSED:
        iecp->connected = 0;
        msg = "disconnected";
        break;

    case CS104_CONNECTION_STARTDT_CON_RECEIVED:
        iecp->connected = 2;
        msg = "Received STARTDT_CON";
        break;

    case CS104_CONNECTION_STOPDT_CON_RECEIVED:
        iecp->connected = 1;
        msg = "Received STOPDT_CON";
        break;

    default:
        break;
    }
	pthread_mutex_unlock(&iecp->connect_lock);
    proto_syslog(LOG_ERR, "IEC104 connection handler, event: %s", msg);
}

#ifndef IEC104_READ_VERBOSE_LOGERR
#define IEC104_READ_VERBOSE_LOGERR 0
#endif

static bool iec104_ASDUReceivedHandler(void * parameter, int address, CS101_ASDU asdu)
{
    channel_t * chan;
    iec104_proto * iecp;
    int i, numel, typebad, comaddr;
    int type = CS101_ASDU_getTypeID(asdu);
    uint64_t cur_ms = clock_get_ms();

    chan = (channel_t *) parameter;
    if (chan == NULL)
        return false;
    iecp = (iec104_proto *) chan->proto_ptr;
    comaddr = CS101_ASDU_getCA(asdu);

    typebad = 0; /* try to keep it as zero */
    numel = CS101_ASDU_getNumberOfElements(asdu);
    for (i = 0; i < numel; i++) {
        int addr, valok;
        tag_value_t tagval;
        InformationObject io = CS101_ASDU_getElement(asdu, i);

        addr = InformationObject_getObjectAddress(io);
        valok = 0;
        tagval.type = 0;
        tagval.value.to_float = 0;
        tagval.value_last.to_float = 0;

        switch (type) {
        case M_SP_NA_1: { /* 1 - Single point information (BOOLEAN) */
            SinglePointInformation t = (SinglePointInformation) io;
            bool bval = SinglePointInformation_getValue(t);
            SinglePointInformation_destroy(t);

            valok = 1;
            tagval.type = 0; /* integer type */
            tagval.value.to_int = (int) bval;
#if IEC104_READ_VERBOSE_LOGERR
            proto_syslog(LOG_ERR, "addr: %d M_SP_NA_1 value: %d", addr, (int) bval);
#endif
            break;
        }

        case M_ME_NC_1: { /* 13 - Short measured value (FLOAT32) */
            MeasuredValueShort t = (MeasuredValueShort) io;
            float fval = MeasuredValueShort_getValue(t);
            MeasuredValueShort_destroy(t);

            valok = 1;
            tagval.type = 1; /* float type */
            tagval.value.to_float = (double) fval;
#if IEC104_READ_VERBOSE_LOGERR
            proto_syslog(LOG_ERR, "addr: %d M_ME_NC_1 value: %f", addr, fval);
#endif
            break;
        }

        case M_IT_NA_1: { /* 15 - Integrated totals (INT32 with quality indicators) */
            IntegratedTotals t = (IntegratedTotals) io;
            int ival = BinaryCounterReading_getValue(IntegratedTotals_getBCR(t));
            IntegratedTotals_destroy(t);

            valok = 1;
            tagval.type = 0; /* integer type */
            tagval.value.to_int = ival;
#if IEC104_READ_VERBOSE_LOGERR
            proto_syslog(LOG_ERR, "addr: %d M_IT_NA_1 value: %d", addr, ival);
#endif
            break;
        }

        case M_ME_NB_1: { /* 11 - Scaled measured value (-32768 ... +32767) */
            MeasuredValueScaled t = (MeasuredValueScaled) io;
            int ival = MeasuredValueScaled_getValue(t);
            MeasuredValueScaled_destroy(t);

            valok = 1;
            tagval.type = 0; /* integer type */
            tagval.value.to_int = ival;
#if IEC104_READ_VERBOSE_LOGERR
            proto_syslog(LOG_ERR, "addr: %d M_ME_NB_1 value: %d", addr, ival);
#endif
            break;
        }

        case M_ME_NA_1: { /* 9 - Normalized measured value (-1.0 ... +1.0) */
            MeasuredValueNormalized t = (MeasuredValueNormalized) io;
            float fval = MeasuredValueNormalized_getValue(t);
            MeasuredValueNormalized_destroy(t);

            valok = 1;
            tagval.type = 1; /* float type */
            tagval.value.to_float = (double) fval;
#if IEC104_READ_VERBOSE_LOGERR
            proto_syslog(LOG_ERR, "addr: %d M_ME_NA_1 value: %f", addr, fval);
#endif
            break;
        }

        case M_ME_TE_1: { /* 35 - measured scaled values with CP56Time2a timestamp */
            int ival = MeasuredValueScaled_getValue((MeasuredValueScaled) io);
            MeasuredValueScaledWithCP56Time2a_destroy((MeasuredValueScaledWithCP56Time2a) io);

            valok = 1;
            tagval.type = 0; /* integer type */
            tagval.value.to_int = ival;
#if IEC104_READ_VERBOSE_LOGERR
            proto_syslog(LOG_ERR, "addr: %d M_ME_TE_1 value: %d", addr, ival);
#endif
            break;
        }

        case M_ME_TF_1: { /* 36 - Short measured value (FLOAT32) with CP56Time2a */
            MeasuredValueShortWithCP24Time2a t = (MeasuredValueShortWithCP24Time2a) io;
            float fval = MeasuredValueShort_getValue((MeasuredValueShort) io);
            MeasuredValueShortWithCP24Time2a_destroy(t);

            valok = 1;
            tagval.type = 1; /* float type */
            tagval.value.to_float = (double) fval;
#if IEC104_READ_VERBOSE_LOGERR
            proto_syslog(LOG_ERR, "addr: %d M_ME_TF_1 value: %f", addr, fval);
#endif
            break;
        }

        case C_SC_NA_1: /* 45 - Send control command */
        case C_SE_NA_1: /* 48 - Setpoint command, normalized value (-1.0 ... +1.0) */
        case C_SE_NB_1: /* 49 - Set-point command, scaled value */
        case C_SE_NC_1: /* 50 - Set-point command, short floating point number */
        case M_EI_NA_1: /* 70 - End of Initialization */
        case C_IC_NA_1: /* 100 - interrogation command */
        case C_CI_NA_1: /* 101 - counter interrogation command */
            /* DO NOT forget to free object: */
            InformationObject_destroy(io);
            break;

        default:
            typebad++;
            /* DO NOT forget to free object: */
            InformationObject_destroy(io);
            break;
        }

        if (valok != 0) {
            tag_t * ptag = NULL;

            if (chan->devs_size > 0) {
                int j, k, p, q;
                device_t * devp;
                uint32_t * paddr;

                k = chan->devs_size;
                if (k > CHANNEL_DEVS_MAX)
                    k = CHANNEL_DEVS_MAX;

                /* fast checking of IEC104 tag search */
                if (iecp != NULL && iecp->last_device != NULL) {
                    for (j = 0; j < k; ++j) {
                        int offs, off1;
                        tag_t * tagp;
                        const iec104_reg_info_t * reginfo;

                        devp = chan->devs[j];
                        if (devp != iecp->last_device)
                            continue;

                        offs = iecp->last_ro_index;
                        for (off1 = 0; off1 <= 4; ++off1) {
                            int offset = offs + off1;
                            if (offset < 0 || offset > devp->tags.ro_tags_size)
                                break;

                            tagp = devp->tags.ro_tags[offset];
                            if (tagp == NULL)
                                continue;

                            reginfo = (const iec104_reg_info_t *) tagp->proto_ptr;
                            if (reginfo == NULL)
                                continue;

                            if (reginfo->ioa_addr == addr) {
                                ptag = tagp;
                                if (iecp != NULL) {
                                    iecp->last_device = devp;
                                    iecp->last_ro_index = offset;
                                }
                                if (ptag != NULL) {
                                proto_set_tag(ptag, tagval, cur_ms);
                                } else {
                                    proto_syslog(LOG_DEBUG, "No tag found at IEC IOA: %d, type: %d", addr, type);
                                }
                                ptag = NULL;
                            }
                        }
                    }
                }

                /* slowly find the device with IEC104 common-address `comaddr */
                for (j = 0; j < k; ++j) {
                    devp = chan->devs[j];
                    if (devp == NULL || devp->proto_ptr == NULL)
                        continue;

                    paddr = (uint32_t *) devp->proto_ptr;
                    if (*paddr != (uint32_t) comaddr)
                        continue;

                    q = devp->tags.ro_tags_size;
                    for (p = 0; p < q; ++p) {
                        tag_t * tagp;
                        const iec104_reg_info_t * reginfo;
                        tagp = devp->tags.ro_tags[p];
                        if (tagp == NULL)
                            continue;

                        reginfo = (const iec104_reg_info_t *) tagp->proto_ptr;
                        if (reginfo == NULL)
                            continue;

                        if (reginfo->ioa_addr == addr) {
                            ptag = tagp;
                            if (iecp != NULL) {
                                iecp->last_device = devp;
                                iecp->last_ro_index = p;
                            }
                            if (ptag != NULL) {
                                proto_set_tag(ptag, tagval, cur_ms);
                            } else {
                                proto_syslog(LOG_DEBUG, "No tag found at IEC IOA: %d, type: %d", addr, type);
                            }
                            ptag = NULL;
                            // goto next;
                        }
                    }
                }
            }

// next:
//             if (ptag != NULL) {
//                 proto_set_tag(ptag, tagval, cur_ms);
//             } else {
//                 proto_syslog(LOG_DEBUG, "No tag found at IEC IOA: %d, type: %d", addr, type);
//             }
        }
    }

    if (typebad != 0) {
        proto_syslog(LOG_ERR, "Error, invalid IEC104 ASDU type: %d", type);
    }
    return true;
}

/*
 * Function used for IEC104 遥控
 * value: must be 0 or 1
 */
static int iec104_single_cmd(CS104_Connection con104, int ioa, int caddr, bool value)
{
    bool ret;
    SingleCommand sc;

    ret = value != 0;
    sc = SingleCommand_create(NULL, ioa, ret, false, 0);
    if (sc == NULL) {
        proto_syslog(LOG_ERR, "Failed to create IEC104 single-command, addr: %d", ioa);
        return -1;
    }

    ret = CS104_Connection_sendProcessCommand(con104, C_SC_NA_1,
        CS101_COT_ACTIVATION, caddr, (InformationObject) sc);
    SingleCommand_destroy(sc);
    if (!ret) {
        proto_syslog(LOG_ERR, "Failed to write IEC104 single-command, addr: %d", ioa);
        return -1;
    }
    return 0;
}

/*
 * Function used for IEC104 遥调,
 * supported range: scaled value (-32768 ... +32767)
 */
static int iec104_setpoint_cmd(CS104_Connection con104, int ioa, int caddr, int value)
{
    bool ret;
    SetpointCommandScaled sc;

    sc = SetpointCommandScaled_create(NULL, ioa, value, false, 0);
    if (sc == NULL) {
        proto_syslog(LOG_ERR, "Failed to create IEC104 setpoint-command, addr: %d", ioa);
        return -1;
    }

    ret = CS104_Connection_sendProcessCommand(con104, C_SE_NB_1,
        CS101_COT_ACTIVATION, caddr, (InformationObject) sc);
    SetpointCommandScaled_destroy(sc);
    if (!ret) {
        proto_syslog(LOG_ERR, "Failed to write IEC104 setpoint-command, addr: %d", ioa);
        return -1;
    }
    return 0;
}

/*
 * Setpoint command, short value (FLOAT32)
 */
static int iec104_setfloat_cmd(CS104_Connection con104, int ioa, int caddr, float value)
{
    bool ret;
    SetpointCommandShort sc;

    sc = SetpointCommandShort_create(NULL, ioa, value, false, 0);
    if (sc == NULL) {
        proto_syslog(LOG_ERR, "Failed to create IEC104 SetpointCommandShort, addr: %d, value: %f",
            ioa, value);
        return -1;
    }

    ret = CS104_Connection_sendProcessCommand(con104, C_SE_NC_1,
        CS101_COT_ACTIVATION, caddr, (InformationObject) sc);
    SetpointCommandShort_destroy(sc);
    if (!ret) {
        proto_syslog(LOG_ERR, "Failed to write IEC104 SetpointCommandShort, addr: %d, value: %f",
            ioa, value);
        return -1;
    }
    return 0;
}

/*
 * Setpoint command, normalized value (-1.0 ... +1.0)
 */
static int iec104_setfloatn_cmd(CS104_Connection con104, int ioa, int caddr, float value)
{
    bool ret;
    SetpointCommandNormalized sc;

    sc = SetpointCommandNormalized_create(NULL, ioa, value, false, 0);
    if (sc == NULL) {
        proto_syslog(LOG_ERR, "Failed to create IEC104 SetpointCommandNormalized, addr: %d, value: %f",
            ioa, value);
        return -1;
    }

    ret = CS104_Connection_sendProcessCommand(con104, C_SE_NA_1,
        CS101_COT_ACTIVATION, caddr, (InformationObject) sc);
    SetpointCommandNormalized_destroy(sc);
    if (!ret) {
        proto_syslog(LOG_ERR, "Failed to write IEC104 SetpointCommandNormalized, addr: %d, value: %f",
            ioa, value);
        return -1;
    }
    return 0;
}

static int iec104_read_poll(CS104_Connection con104, chunk_interv_t * chkp, int idx)
{
    int caddr;
    device_t * pdev;
    uint32_t * paddr;
    uint64_t nowms;
    proto_iec104_poll * poll_info;

    // get device
    pdev = chkp->dev_ptr;
    if (pdev == NULL) {
        proto_syslog(LOG_ERR, "invalid NULL device at chunks_interv: %d", idx);
        return -1;
    }
    paddr = (uint32_t *) pdev->proto_ptr;
    if (paddr == NULL) {
        proto_syslog(LOG_ERR, "Invalid common address for IEC104 at %d", idx);
        return -1;
    }

    caddr = (int) *paddr;

    poll_info = (proto_iec104_poll *) chkp->chunk_proto_ptr;
    if (poll_info == NULL) {
        proto_syslog(LOG_ERR, "invalid NULL poll_info at index: %d", idx);
        return -1;
    }

    nowms = clock_get_ms();
    if ((nowms - chkp->last_poll) >= chkp->poll_interv) {
        bool bret;

        chkp->last_poll = nowms;
        if (poll_info->iec104_id == 0) // IEC104 总召唤:
            bret = CS104_Connection_sendInterrogationCommand(con104, CS101_COT_ACTIVATION, caddr, IEC60870_QOI_STATION);
        else // 发电度量召唤:
            bret = CS104_Connection_sendCounterInterrogationCommand(con104,
                CS101_COT_ACTIVATION, caddr, IEC60870_QCC_FRZ_READ + IEC60870_QCC_RQT_GENERAL);
        if (bret == 0) {
            proto_syslog(LOG_ERR, "IEC104 InterrogationCommand has failed: %d, common-addr: %d", (int) bret, caddr);
        }
        refresh_online_state(pdev);
        return 1;
    }
    return 0;
}

static int iec104_write_poll(CS104_Connection con104, channel_t * channel, int msec)
{
    int ret, err, caddr;
    device_t * pdev;
    tag_t * write_tag;
    write_list_t * w_info;
    iec104_reg_info_t * reginfo;

    err = 0;
    w_info = NULL;
    write_tag = NULL;

wait_again:
    ret = sem_wait_time(&channel->write_sem, msec / 1000, msec % 1000);
    if (ret < 0) {
        err = errno;
        if (err == ETIMEDOUT)
            return 0;

        proto_syslog(LOG_ERR, "semaphore wait has failed for IEC104 write: %s",
            strerror(err));
        return -1;
    }

    msec = 0; // wait zero milliseconds next time
    w_info = get_data_from_list(channel);
    if (w_info == NULL) {
        proto_syslog(LOG_ERR, "get data from list is null");
        return -1;
    }

    write_tag = &w_info->w_tag;
    reginfo = (iec104_reg_info_t *) write_tag->proto_ptr;
    if (reginfo == NULL) {
        proto_syslog(LOG_ERR, "Error, failed to get IEC104 tag reg-info");
        BUSINESS_LOG(BLOG_ERR, DAS_ID_IEC104_6, NULL, "[数采]: 请检查 IEC104 设备寄存器地址是否正确");
        free(w_info);
        w_info = NULL;
        return -1;
    }

    pdev = (device_t *) write_tag->dev_ptr;
    if (pdev == NULL) {
        proto_syslog(LOG_ERR, "Error, invalid device pointer for write_tag.");
        BUSINESS_LOG(BLOG_ERR, DAS_ID_IEC104_7, NULL, "[数采]: 请检查 IEC104 设备写寄存器tag是否正确");
        free(w_info);
        w_info = NULL;
        goto wait_again;
    }

    do {
        uint32_t * paddr = (uint32_t *) pdev->proto_ptr;
        if (paddr == NULL) {
            proto_syslog(LOG_ERR, "Invalid common address for IEC104 device: %p", pdev);
            free(w_info);
            w_info = NULL;
            return -1;
        }
        caddr = (int) *paddr;
    } while (0);

    proto_syslog(LOG_ERR, "IEC104 data list is [%s.%s(type:%d, double=%f, long=%ld)]",
        pdev->no, write_tag->name, write_tag->data_type, write_tag->write_cache.to_float, write_tag->write_cache.to_int);

    ret = -1;
    switch (reginfo->ioa_type) {
    case C_SC_NA_1: /* 单点值，取值范围：0或1两个值 */
        if (write_tag->data_type == 0)
            ret = iec104_single_cmd(con104, reginfo->ioa_addr, caddr, write_tag->write_cache.to_int != 0);
        break;

    case C_SE_NB_1: /* 点位值，取值范围：-32768 ... +32767 */
        if (write_tag->data_type == 0)
            ret = iec104_setpoint_cmd(con104, reginfo->ioa_addr, caddr, write_tag->write_cache.to_int);
        break;

    case C_SE_NC_1: /* 32位浮点，对应C语言中的 float */
        if (write_tag->data_type != 0)
            ret = iec104_setfloat_cmd(con104, reginfo->ioa_addr, caddr, write_tag->write_cache.to_float);
        break;

    case C_SE_NA_1: /* 均一化浮点，取值范围：-1.0 ... +1.0 */
        if (write_tag->data_type != 0)
            ret = iec104_setfloatn_cmd(con104, reginfo->ioa_addr, caddr, write_tag->write_cache.to_float);
        break;

    default:
        ret = -2;
        proto_syslog(LOG_ERR, "Error, invalid IEC104 reginfo, type: %d, ioa: %d",
            reginfo->ioa_type, reginfo->ioa_addr);
        BUSINESS_LOG(BLOG_ERR, DAS_ID_IEC104_8, NULL, "[数采]: IOA: %d, %d 不支持的数据类型, 请检查模板协议数据类型是否正确");
        break;
    }

    if (ret < 0) {
        w_info->pw_tag->last_state_code = w_info->pw_tag->state_code;
        w_info->pw_tag->state_code = TAG_STATE_BAD;
        check_online_state(pdev);
        proto_syslog(LOG_ERR, "IEC104 write tag error, tag:%s, ret = %d", write_tag->name, ret);
        BUSINESS_LOG(BLOG_ERR, DAS_ID_IEC104_9, NULL, "[数采]: IEC104 写tag: %s 失败", write_tag->name);
    } else {
        sync_wtag_read_cache(w_info);
    }

    free(w_info);
    w_info = NULL;
    goto wait_again;
}

static int iec104_client_channel(channel_t * channel)
{
    uint64_t now, then;
    iec104_proto * iecp;
    CS104_Connection con104;
	uint16_t connected;

    con104 = NULL;
    iecp = (iec104_proto *) channel->proto_ptr;
    if (iecp == NULL || iecp->iec104_ip == NULL || iecp->iec104_port == 0) {
        proto_syslog(LOG_CRIT, "[数采]:Invalid empty protocol pointer for IEC104: %p", iecp);
        BUSINESS_LOG(BLOG_ERR, DAS_ID_IEC104_1, channel->channel, "[数采]: %s 通道配置错误, 请检查 IEC104 IP地址和端口(默认2404)是否配置正确", channel->channel);
        sleep(5);
        return -1;
    }
    else
    {
        BUSINESS_LOG(BLOG_ERR, DAS_ID_IEC104_1, channel->channel, "[数采]: %s 通道配置正常", channel->channel);
    }

    /* update IEC104 connection flag to false */
    iecp->connected = 0;
    proto_syslog(LOG_INFO, "Trying to connect to IEC104 %s:%u...",
        iecp->iec104_ip, (unsigned int) iecp->iec104_port);
    con104 = CS104_Connection_create(iecp->iec104_ip, iecp->iec104_port);
    if (con104 == NULL) {
        proto_syslog(LOG_ERR, "Failed to create IEC104 connection.");
        BUSINESS_LOG(BLOG_ERR, DAS_ID_IEC104_2, channel->channel, "[数采]: IEC104 %s 通道创建失败", channel->channel);
        sleep(5);
        return -1;
    }
    else
    {
        BUSINESS_LOG(BLOG_ERR, DAS_ID_IEC104_2, channel->channel, "[数采]: IEC104 %s 通道创建成功", channel->channel);
    }

	pthread_mutex_init(&iecp->connect_lock, NULL);
    CS104_Connection_setConnectTimeout(con104, 3000);
    CS104_Connection_setConnectionHandler(con104, iec104_ConnectionHandler, channel);
    CS104_Connection_setASDUReceivedHandler(con104, iec104_ASDUReceivedHandler, channel);
    CS104_Connection_connectAsync(con104);

    now = sysuptime();
    /* wait until IEC104 connection is established: */
    for (;;) {
		pthread_mutex_lock(&iecp->connect_lock);
		connected = iecp->connected;
		pthread_mutex_unlock(&iecp->connect_lock);
        if (connected != 0)
            break;
        sleep(1);
        then = sysuptime();
        if ((then - now) >= 5000)
            break;
    }

    /* IEC104 client main loop */
    while (1) {
		pthread_mutex_lock(&iecp->connect_lock);
		connected = iecp->connected;
		if (connected == 0) {
			pthread_mutex_unlock(&iecp->connect_lock);
			break;
		}
        CS104_Connection_sendStartDT(con104);
		pthread_mutex_unlock(&iecp->connect_lock);

        now = sysuptime();
        /* wait until IEC104 connection is ready for transfer: */
        for (;;) {
			pthread_mutex_lock(&iecp->connect_lock);
			connected = iecp->connected;
			pthread_mutex_unlock(&iecp->connect_lock);
            if (connected == 2)
                break;
            sleep(1);
            then = sysuptime();
            if ((then - now) >= 5000)
                break;
        }

        while (1) {
			pthread_mutex_lock(&iecp->connect_lock);
			if (iecp->connected != 2) {
				pthread_mutex_lock(&iecp->connect_lock);
				break;
			}
			
            int i, j;
            j = channel->chunks_interv_size;
            for (i = 0; i < j; ++i) {
                chunk_interv_t * chkp;

                // get chunk interval struct pointer
                chkp = channel->chunks_interv[i];
                if (chkp == NULL) {
                    proto_syslog(LOG_ERR, "invalid NULL chunks_interv at index: %d", i);
                    continue;
                }
                iec104_read_poll(con104, chkp, i);
            }

            iec104_write_poll(con104, channel, 1000);
			pthread_mutex_unlock(&iecp->connect_lock);
        }

    }

    CS104_Connection_destroy(con104);
	pthread_mutex_destroy(&iecp->connect_lock);
    proto_syslog(LOG_ERR, "Error, IEC104 connection has closed: %s:%u",
        iecp->iec104_ip, (unsigned int) iecp->iec104_port);
    return 0;
}

char *get_cnt_id(void)
{
#ifdef EN_EMS
    return get_cabinet_info()->cnt_id;
#else
    return NULL;
#endif
}
static volatile int go_rx_broadcast_data = 1; 

int get_is_rx_broadcast_data(void)
{
    return go_rx_broadcast_data > 0;
}

static void *check_broadcast_data(void *arg)
{
    channel_t *channel = (channel_t *)arg;
    RAW_SOCK *raw = protocol_Init_with_group(get_pcs_ctrl_var()->net_node, get_cnt_id(), channel->channel);
    if (raw == NULL)
    {
        proto_syslog(LOG_CRIT, "raw init error:%s", channel->channel);
        return NULL;
    }
    ETH_HEADER *hd = NULL;
    while (channel->some_sign != 0)
    {
        if ((OK == raw_data_recv(raw)) && (raw->len > sizeof(ETH_HEADER)))
        {
            hd = (ETH_HEADER *)raw->rcv_buffer;
            proto_syslog(LOG_INFO, "receive broadcast data from %02x:%02x:%02x:%02x:%02x:%02x", hd->ether_shost[0], hd->ether_shost[1], hd->ether_shost[2], hd->ether_shost[3], hd->ether_shost[4], hd->ether_shost[5]);
            if (memcmp(raw->rcv_buffer + sizeof(ETH_HEADER) + 1, raw->channel_name, strlen(raw->channel_name)) == 0)
            {
                if (memcmp(hd->ether_shost, raw->smac_hex, MAC_HEX_LEN) < 0)
                {
                    go_rx_broadcast_data++;
                    proto_syslog(LOG_INFO, "go_rx_broadcast_data:%d", go_rx_broadcast_data);
                }
            }
        }
    }
    return NULL;
}

static int check_broadcast_data_start(channel_t *channel)
{
    pthread_attr_t attr;
    go_rx_broadcast_data = 0;
    if (0 != pthread_attr_init(&attr))
    {
        return -3;
    }

    if (0 != pthread_attr_setdetachstate(&attr, PTHREAD_CREATE_DETACHED))
    {
        pthread_attr_destroy(&attr);
        return -4;
    }
    pthread_t thread;
    int ret = pthread_create(&thread, &attr, check_broadcast_data, channel);
    pthread_attr_destroy(&attr);
    if (ret != 0)
    {
        ems_syslog(LOG_ERR, "created pthread err[%s]", strerror(errno));
        return -2;
    }
    else
    {
        channel->some_sign = 1;
        pthread_setname_np(thread, "check_bdcast");
    }
    return 0;
}

static int check_broadcast_data_stop(channel_t *channel)
{
    channel->some_sign = 0; // 直接用这个成员标识读取实际接口的线程是否结束
    return 0;
}

/* Raw data callback functions */
static void modbus_raw_send_callback(modbus_t *ctx, const uint8_t *data, int data_length, void *channel)
{
    // printf("modbus_raw_send_callback: channel=%p ========", channel);
    // for(int i = 0; i < data_length; i++)
    // {
    //     printf("%02x ", data[i]);
    // }
    // printf("\n");
    channel_t *ch = channel;
    
    /* Fill send_buf in channel structure */
    if (ch->send_buf != NULL) {
        free(ch->send_buf);
    }
    ch->send_buf = malloc(data_length);
    if (ch->send_buf != NULL) {
        memcpy(ch->send_buf, data, data_length);
        ch->send_buf_size = data_length;
    }
    
}

static void modbus_raw_receive_callback(modbus_t *ctx, const uint8_t *data, int data_length, void *channel)
{
    // printf("modbus_raw_receive_callback: channel=%p ========", channel);
    // for(int i = 0; i < data_length; i++)
    // {
    //     printf("%02x ", data[i]);
    // }
    // printf("\n");
    channel_t *ch = channel;
    
    /* Fill recv_buf in channel structure */
    if (ch->recv_buf != NULL) {
        free(ch->recv_buf);
    }
    ch->recv_buf = malloc(data_length);
    if (ch->recv_buf != NULL) {
        memcpy(ch->recv_buf, data, data_length);
        ch->recv_buf_size = data_length;
    }
    
}

static int modbus_channel(channel_t *channel)
{
    long long cur_ms = 0;
    modbus_t *ctx = NULL;
    uint8_t buff[256] = {0};
    device_t *dev = NULL;
    int err_counter = 0;
    modbus_poll_info_t *poll_info = NULL;
    int rx_len = 0;
    int32_t fgap = -1;
    uint64_t lastread = sysuptime() - 10000;
    RAW_SOCK *raw = protocol_Init_with_group(get_pcs_ctrl_var()->net_node,  get_cnt_id(), channel->channel);
    if (raw == NULL)
    {
        proto_syslog(LOG_CRIT, "[数采]:creat %s error!", channel->channel);
        BUSINESS_LOG(BLOG_ERR, DAS_ID_MODBUS_1, channel->channel, "[数采]: %s 通道创建失败", channel->channel);
        goto modbus_channel_exit;
    }
    BUSINESS_LOG(BLOG_NOTICE, DAS_ID_MODBUS_1, channel->channel, "[数采]: %s 通道创建成功", channel->channel);

    ctx = modbus_proto_init(channel->proto_ptr);
    if (ctx == NULL)
    {
        proto_syslog(LOG_CRIT, "[数采]:modbus_proto_init error in the channel:[%s], ERROR:%s, errno=%d", channel->channel, modbus_strerror(errno), errno);
        BUSINESS_LOG(BLOG_ERR, DAS_ID_MODBUS_2, channel->channel, "[数采]: %s 通道初始化失败, 失败原因(%s)", channel->channel, modbus_strerror(errno));
        goto modbus_channel_exit;
    }
    modbus_set_raw_send_callback(ctx, modbus_raw_send_callback,(void *)channel);
    modbus_set_raw_receive_callback(ctx, modbus_raw_receive_callback,(void *)channel);

    channel->s = get_fd_by_ctx(ctx);//获取通道的描述符
    if (contain_meter(channel) != 0)
    {
        check_broadcast_data_start(channel);
    }
    proto_syslog(LOG_INFO, "modbus_proto_init success, in the channel:[%s]", channel->channel);
    BUSINESS_LOG(BLOG_NOTICE, DAS_ID_MODBUS_2, channel->channel, "[数采]: %s 通道初始化成功", channel->channel);
    while (1)
    {
        int contine_interv = 0;
        // ret = cond_wait(&channel->cond,&channel->mutex,0,INVITETIME);
        
        if(check_write(ctx, channel, 20)<0)
            err_counter++;
        do
        {
            contine_interv = 0;
            check_tag_state(channel);
            cur_ms = clock_get_ms();
            for (size_t j = 0; j < channel->devs_size; j++)
            {
            dev = channel->devs[j];
            chunk_interv_t *chunks_interv = NULL;
            for (size_t i = 0; (dev != NULL) && (i < dev->chunks_interv_size); i++)
            {
                if(check_write(ctx, channel, 0)<0)
                    err_counter++;

                chunks_interv = dev->chunks_interv[i];
                poll_info = (modbus_poll_info_t *)chunks_interv->chunk_proto_ptr;
                if ((cur_ms - chunks_interv->last_poll) >= chunks_interv->poll_interv)
                {
                    chunks_interv->last_poll = cur_ms;
                    setmodebustimeout(ctx, dev->timeout);
                    proto_syslog(LOG_DEBUG, "[channel:%s][dev:%s] modbus read dev_addr=%d, code=%d, addr=%d, len=%d, interv=%u, cur_ms=%lld",
                              channel->channel, dev->no, *((int *)dev->proto_ptr), poll_info->func_code, poll_info->reg_addr, poll_info->size, chunks_interv->poll_interv, cur_ms);
                    channel_lock(channel);
					if (fgap != dev->write_gap)
						fgap = dev->write_gap;
                    frame_wgap_delay(lastread, fgap);
                    rx_len = modbus_proto_read(ctx, *((int *)dev->proto_ptr), poll_info->func_code, poll_info->reg_addr, poll_info->size, buff); // 读数据
                    lastread = sysuptime();
                    channel_unlock(channel);
                    if (rx_len > 0)
                    {
                        BUSINESS_LOG(BLOG_NOTICE, DAS_ID_MODBUS_3, channel->channel, "[数采]: %s 通信正常", channel->channel);
                        // proto_syslog(LOG_DEBUG, "[channel:%s][dev:%s] modbus TAGS:", channel->channel, dev->no);
#if 0
                        modbus_update_tags(dev, poll_info->func_code, poll_info->reg_addr, buff, rx_len); // 根据返回数据更新tag
#else
                        if ((dev->is_meter != 0)&&(dev->dev_child_type!=NULL)&&((strcmp(dev->dev_child_type,"GRID") == 0)||(strcmp(dev->dev_child_type,"TRAS_LV") == 0))) // 是电表
                        {
                            if (go_rx_broadcast_data != 0)
                            {
                                proto_syslog(LOG_WARNING, "%s:go to rx broadcast data", channel->channel);
                                goto modbus_channel_exit;
                            }
                            if (0 > modbus_broadcast(raw, *((int *)dev->proto_ptr), poll_info, buff, rx_len))
                            {
                                proto_syslog(LOG_ERR, "modbus broadcast data error");
                                BUSINESS_LOG(BLOG_ERR, DAS_ID_MODBUS_4, dev->no, "[数采]: 电表(%s) 广播数据失败, 请检查网络状态是否正常",dev->no);
                            }
                            else
                            {
                                BUSINESS_LOG(BLOG_NOTICE, DAS_ID_MODBUS_4, dev->no, "[数采]: 电表(%s) 广播数据正常",dev->no);
                            }
                        }
                        modbus_update_tags_new(poll_info, buff, rx_len); // 优化查找后的更新tag
                        abnormal_log_match_record(poll_info->p_tag, poll_info->tag_sum, channel->send_buf, channel->send_buf_size, channel->recv_buf, channel->recv_buf_size);
#endif
                        refresh_online_state(dev);
                        err_counter = 0;
                        
                    }
                    else
                    {
                        check_online_state(dev);
                        proto_syslog(LOG_ERR, "[channel:%s][dev:%s] :modbus read error(rx_len = %d), ERROR info:%s, errno=%d", channel->channel, dev->no, rx_len, modbus_strerror(errno), errno);
                        BUSINESS_LOG(BLOG_ERR, DAS_ID_MODBUS_3, channel->channel, "[数采]: %s 通信失败, 失败原因(%s)", channel->channel, modbus_strerror(errno));
                        if ((errno < MODBUS_ENOBASE) && (errno != ETIMEDOUT ))//非超时的io错误
                        {
                            BUSINESS_LOG(BLOG_ERR, DAS_ID_MODBUS_3, channel->channel, "[数采]: %s 通信出现系统级IO错误, 请检查 物理连接、设备状态、参数配置、设备地址 是否正常", channel->channel);
                            goto modbus_channel_exit;
                        }
                        err_counter++;
                        dev = NULL;// 为了去给下一个设备发送命令
                    }

                    if (err_counter > COMMUNIT_RETRY_MAX)
                    {
                        goto modbus_channel_exit;
                    }
                    contine_interv = 1;
                }
            }
            }
        } while ((contine_interv != 0));

        // interv空闲
        if (channel->chunks_nonstop_size > 0) // 有配置
        {
            size_t nonstop_index = channel->chunks_nonstop_index;
            dev = channel->chunks_nonstop[nonstop_index]->dev_ptr;
            poll_info = (modbus_poll_info_t *)channel->chunks_nonstop[nonstop_index]->chunk_proto_ptr;

            proto_syslog(LOG_DEBUG, "[channel:%s][dev:%s]  modbus read dev_addr=%d, code=%d, addr=%d, len=%d",
                      channel->channel, dev->no, *((int *)dev->proto_ptr), poll_info->func_code, poll_info->reg_addr, poll_info->size);
            channel_lock(channel);

            if (fgap != dev->write_gap)
                fgap = dev->write_gap;
            frame_wgap_delay(lastread, fgap);
            rx_len = modbus_proto_read(ctx, *((int *)dev->proto_ptr), poll_info->func_code, poll_info->reg_addr, poll_info->size, buff); // 读数据
            lastread = sysuptime();

            channel_unlock(channel);
            if (rx_len > 0)
            {
                // proto_syslog(LOG_DEBUG, "[channel:%s][dev:%s] modbus TAGS:", channel->channel, dev->no);
                refresh_online_state(dev);
                modbus_update_tags(dev, poll_info->func_code, poll_info->reg_addr, buff, rx_len); // 根据返回数据更新tag
                err_counter = 0;
                BUSINESS_LOG(BLOG_NOTICE, DAS_ID_MODBUS_3, channel->channel, "[数采]: %s 通信正常", channel->channel);   
            }
            else
            {
                check_online_state(dev);
                proto_syslog(LOG_ERR, "[channel:%s][dev:%s] modbus read error(rx_len = %d), ERROR info:%s, errno=%d", channel->channel, dev->no, rx_len, modbus_strerror(errno), errno);
                BUSINESS_LOG(BLOG_ERR, DAS_ID_MODBUS_3, NULL, "[数采]: %s 通信失败(无回复数据), 失败原因(%s)", channel->channel, modbus_strerror(errno));                        
                if ((errno < MODBUS_ENOBASE) && (errno != ETIMEDOUT ))//非超时的io错误
                {
                    BUSINESS_LOG(BLOG_ERR, DAS_ID_MODBUS_3, channel->channel, "[数采]: %s 通信出现系统级IO错误, 请检查 物理连接、设备状态、参数配置、设备地址 是否正常", channel->channel);
                    goto modbus_channel_exit;
                }
                err_counter++;
            }
            if (err_counter > COMMUNIT_RETRY_MAX)
            {
                goto modbus_channel_exit;
            }
            nonstop_index++;
            if (nonstop_index >= channel->chunks_nonstop_size)
            {
                channel->chunks_nonstop_index = 0;
            }
            else
            {
                channel->chunks_nonstop_index = nonstop_index;
            }
        }
    }
modbus_channel_exit:
    check_broadcast_data_stop(channel);
    if (ctx != NULL)
    {
        modbus_proto_deinit(ctx);
    }
    if (raw != NULL)
    {
        raw_config_deinit(raw);
    }
    ctx = NULL;
    return -2;
}

//
int dlt645_update_tags(device_t *dev, dlt645_poll_info_t *poll_info, void *buff, int it_sum, channel_t *channel)
{
    tag_t *p_tag = NULL;
    dlt645_rw_info_t *read_reg = NULL;
    int32_t *p_int = (int32_t *)buff;
    double *p_double = (double *)buff;
    uint64_t cur_ms = clock_get_ms();
    tag_value_t tag_value = {.type = 0};
    if (poll_info->pos == 0) // 整形
    {
        tag_value.type = 0;
        for (size_t i = 0; i < dev->tags.ro_tags_size; i++)
        {
            read_reg = (dlt645_rw_info_t *)dev->tags.ro_tags[i]->proto_ptr;
            p_tag = dev->tags.ro_tags[i];
            if (read_reg->data_id == poll_info->data_id)
            {
                if (read_reg->index < it_sum)
                {
                    tag_value.value.to_int = *(p_int + read_reg->index);
                    proto_set_tag(p_tag, tag_value, cur_ms);   
                    abnormal_log_match_record(&p_tag, 1, channel->send_buf, channel->send_buf_size, channel->recv_buf, channel->recv_buf_size);
                }
                else
                {
                    proto_syslog(LOG_WARNING, "TAG:%s, index warning, index=%d", p_tag->name, read_reg->index);
                }
            }
        }

        for (size_t i = 0; i < dev->tags.rw_tags_size; i++)
        {
            read_reg = (dlt645_rw_info_t *)dev->tags.rw_tags[i]->proto_ptr;
            p_tag = dev->tags.rw_tags[i];
            if (read_reg->data_id == poll_info->data_id)
            {
                if (read_reg->index < it_sum)
                {
                    tag_value.value.to_int = *(p_int + read_reg->index);
                    proto_set_tag(p_tag, tag_value, cur_ms);
                    abnormal_log_match_record(&p_tag, 1, channel->send_buf, channel->send_buf_size, channel->recv_buf, channel->recv_buf_size);
                }
                else
                {
                    proto_syslog(LOG_WARNING, "TAG:%s, index warning, index=%d", p_tag->name, read_reg->index);
                }
            }
        }
    }
    else
    {
        tag_value.type = 1; 
        for (size_t i = 0; i < dev->tags.ro_tags_size; i++)
        {
            read_reg = (dlt645_rw_info_t *)dev->tags.ro_tags[i]->proto_ptr;
            p_tag = dev->tags.ro_tags[i];
            if (read_reg->data_id == poll_info->data_id)
            {
                if (read_reg->index < it_sum)
                {
                    tag_value.value.to_float = *(p_double + read_reg->index);
                    proto_set_tag(p_tag, tag_value, cur_ms);  
                    abnormal_log_match_record(&p_tag, 1, channel->send_buf, channel->send_buf_size, channel->recv_buf, channel->recv_buf_size);
                }
                else
                {
                    proto_syslog(LOG_WARNING, "TAG:%s, index warning, index=%d", p_tag->name, read_reg->index);
                }
            }
        }
    }
    return 0;
}

static int dlt645_check_write(dlt645_hd_t *hd, channel_t *channel, int time_out)
{
    device_t *dev = NULL;
    dlt645_rw_info_t *rw_info = NULL;
    int ret = sem_wait_time(&channel->write_sem, 0, time_out);
    if (ret == 0)
    {
        write_list_t *w_info = get_data_from_list(channel);
        if (w_info == NULL)
        {
            proto_syslog(LOG_ERR, "get data from list is null");
            return -1;
        }
        tag_t *write_tag = &w_info->w_tag;
        if (write_tag == NULL)
        {
            proto_syslog(LOG_ERR, "get data from list is null");
            free(w_info);
            return -2;
        }
        dev = (device_t *)write_tag->dev_ptr;
        rw_info = (dlt645_rw_info_t *)write_tag->proto_ptr;
        if (rw_info != NULL)
        {
            proto_syslog(LOG_ERR, "============================================================");
            // dlt645_write_data(hd, ((dlt645_dev_attr *)dev->proto_ptr)->bcd_addr, &it, &data);
        }
        else
        {
            ret = (dev->adp != NULL) ? dev->adp->handle(write_tag, hd) : -3;
            if (ret != 0)
            {
                w_info->pw_tag->last_state_code = w_info->pw_tag->state_code;
                w_info->pw_tag->state_code = TAG_STATE_BAD;
            }
            else
            {
                sync_wtag_read_cache(w_info);
            }
        }
        free(w_info);
    }
    return -2;
}


int send_protocol_data_cb(char *data,int len,void *contex)
{
    return send_protocol_data(contex,data,len);
}

void *net_get_dlt_task(void *arg)
{
    channel_t *channel = (channel_t *)arg;
    RAW_SOCK *raw = NULL;
RETURY: 
    if ((raw = protocol_Init_with_group(get_pcs_ctrl_var()->net_node, get_cnt_id(), MAC_CHANNEL_NAME))==NULL)
    {
        sleep(10);
        goto RETURY;
    }
    char data_645[256] = {0};
    int len = 0;
    char data_tmp[256] = {0};
    unsigned int data_id = 0;
    device_t *dev = NULL;
    dlt645_poll_info_t *poll_info = NULL;
    while (1)
    {
START:        
        len = 0;
        memset(data_tmp,0,sizeof(data_tmp));
        if(read_protocol_data(raw,data_645,&len)==0)
        {
            proto_syslog(LOG_INFO, "dlt read len:%d",len);
            if(len>0)
            {
                proto_syslog_hex(LOG_INFO,data_645,len,"net dlt data:");
                dlt645_2007_head *head = (dlt645_2007_head *)data_645;
                dlt645_1997_head *head_1997 = (dlt645_1997_head *)data_645;
                proto_syslog(LOG_INFO, "head_1997->com.ctrl_code:%02x",head_1997->com.ctrl_code);
                if (head_1997->com.ctrl_code == 0x81) // 这儿不严谨，后期需优化
                {
                    data_id = ((unsigned char)(head_1997->di[1]-0x33)<<8)|(unsigned char)(head_1997->di[0]-0x33);
                }
                else
                {
                    data_id = ((unsigned char)(head->di[3]-0x33)<<24)|((unsigned char)(head->di[2]-0x33)<<16)|((unsigned char)(head->di[1]-0x33)<<8)|((unsigned char)head->di[0]-0x33);
                    proto_syslog(LOG_DEBUG, "ctrl_code hex:%02x %02x %02x %02x",head->di[3],head->di[2],head->di[1],head->di[0]);
                }
                proto_syslog(LOG_INFO, "ctrl_code:%02x",data_id);
                for (size_t i = 0; i < channel->chunks_interv_size; i++)
                {
                    poll_info = (dlt645_poll_info_t *)channel->chunks_interv[i]->chunk_proto_ptr;
                    if (poll_info->data_id != data_id)
                        continue;
                    dev = channel->chunks_interv[i]->dev_ptr;
                    int ret = dlt645_analysis(poll_info,data_645,len,((dlt645_dev_attr *)dev->proto_ptr)->bcd_addr,data_tmp,sizeof(data_tmp));
                    if (ret > 0)
                    {
                        // proto_syslog(LOG_DEBUG, "[channel:%s][dev:%s] dlt645 TAGS:", channel->channel, dev->no);
                        dlt645_update_tags(dev, poll_info, data_tmp, ret, channel);
                        refresh_online_state(dev);
                    }
                    else
                    {
                        proto_syslog(LOG_NOTICE, "%s:dlt645_analysis error", channel->channel);
                        check_online_state(dev);
                    }
                    check_tag_state(channel);
                    goto START;
                }
                for (size_t i = 0; i < channel->chunks_nonstop_size; i++)
                {
                    poll_info = (dlt645_poll_info_t *)channel->chunks_nonstop[i]->chunk_proto_ptr;
                    if (poll_info->data_id != data_id)
                        continue;
                    int ret = dlt645_analysis(poll_info,data_645,len,((dlt645_dev_attr *)dev->proto_ptr)->bcd_addr,data_tmp,sizeof(data_tmp));
                    if (ret > 0)
                    {
                        // proto_syslog(LOG_DEBUG, "[channel:%s][dev:%s] dlt645 TAGS:", channel->channel, dev->no);
                        dlt645_update_tags(dev, poll_info, data_tmp, ret, channel);
                        refresh_online_state(dev);
                    }
                    else
                    {
                        proto_syslog(LOG_ERR, "%s:dlt645 read error(ret = %d), info:%s, errno:%d", channel->channel, ret, strerror(errno), errno);
                        check_online_state(dev);
                    }
                    check_tag_state(channel);
                    goto START;
                }
                proto_syslog(LOG_INFO, "can't find pr");         
            }
        }
        else
        {
            for (size_t i = 0; i < channel->devs_size; i++)
            {
                check_online_state(channel->devs[i]); 
            }
        }
        check_tag_state(channel);
    }
    
}

#if EN_EMS
typedef struct modbus_meter_dev_poll
{
    uint8_t dev_addr;
    modbus_poll_info_t **poll;
    int poll_len;
} modbus_meter_dev_poll;

static int find_meter_sum(channel_t *channel, uint8_t *addr)
{
    int tmp = 0;
    for (int i = 0; i < channel->devs_size; i++)
    {
        if (channel->devs[i]->is_meter)
        {
            addr[tmp] = *((int *)channel->devs[i]->proto_ptr);
            tmp++;
        }
    }
    return tmp;
}

static modbus_meter_dev_poll new_dev_poll(channel_t *channel, uint8_t dev_addr)
{
    int tmp = 0;
    modbus_meter_dev_poll ret = {.dev_addr = dev_addr};
    // 查找poll个数
    for (int i = 0; i < channel->chunks_interv_size; i++)
    {
        if (*((int *)channel->chunks_interv[i]->dev_ptr->proto_ptr) == dev_addr)
        {
            tmp++;
        }
    }
    ret.poll_len = tmp;
    ret.poll = calloc(tmp, sizeof(modbus_poll_info_t **));
    if (ret.poll == NULL)
    {
        ret.dev_addr = 0xff; // 使用一个无效的地址
    }
    // 赋值
    tmp = 0;
    for (int i = 0; i < channel->chunks_interv_size; i++)
    {
        if (*((int *)channel->chunks_interv[i]->dev_ptr->proto_ptr) == dev_addr)
        {
            ret.poll[tmp++] = (modbus_poll_info_t *)channel->chunks_interv[i]->chunk_proto_ptr;
        }
    }
    return ret;
}

static int delete_dev_poll(modbus_meter_dev_poll *p)
{
    free(p->poll);
    return 0;
}

static int find_modbus_meter_dev_poll(channel_t *channel, modbus_meter_dev_poll **p_out)
{
    int tmp = 0;
    uint8_t addr[CHANNEL_DEVS_MAX] = {0};
    modbus_meter_dev_poll *dev_poll = NULL;
    modbus_poll_info_t *md_poll = NULL;
    tmp = find_meter_sum(channel, addr);
    proto_syslog(LOG_NOTICE, "found %d meters in channel %s", tmp, channel->channel);
    if (tmp < 1)
    {
        return -1;
    }
    dev_poll = calloc(tmp, sizeof(modbus_meter_dev_poll));
    if (dev_poll == NULL)
    {
        proto_syslog(LOG_NOTICE, "calloc error for channel %s", channel->channel);
        return -2;
    }
    // 检查该通道电表个数
    for (int i = 0; i < tmp; i++)
    {
        dev_poll[i] = new_dev_poll(channel, addr[i]);
    }
    // 打印
    proto_syslog(LOG_NOTICE, "===================================");
    for (int i = 0; i < tmp; i++)
    {
        proto_syslog(LOG_NOTICE, "dev_poll[%d]:addr=0x%02x", i, dev_poll->dev_addr);
        for (int j = 0; j < dev_poll[i].poll_len; j++)
        {
            md_poll = dev_poll->poll[j];
            proto_syslog(LOG_NOTICE, "index:%d,fun_code=0x%02x,reg_addr=0x%x len=%d", md_poll->index, md_poll->func_code, md_poll->reg_addr, md_poll->size);
        }
    }
    *p_out = dev_poll;
    return tmp;
}

static int net_get_modbus_task(channel_t *channel)
{
    RAW_SOCK *raw = NULL;
    modbus_poll_info_t *poll_info = NULL;
    modbus_meter_dev_poll *dev_poll = NULL;
    int len = 0;
    int err_counter = 0;
    int dev_poll_sum = 0;
    time_t last_t = 0;
    time_t cur_t = time(NULL);
    struct
    {
        uint16_t poll_index; // 通道中poll的位置
        uint16_t data_len;   // 数据长度
        uint8_t dev_addr;    // 设备地址
        uint8_t data[256];   // 数据
    } rx_pack;

    for (int i = 0; i < 5; i++)
    {
        raw = protocol_Init_with_group(get_pcs_ctrl_var()->net_node, get_cnt_id(), channel->channel);
        if (raw != NULL)
        {
            break;
        }
        sleep(1);
    }
    int ch_broadcast_timeout = get_pcs_ctrl_var()->ch_broadcast_timeout;
    if (raw != NULL)
    {
        proto_syslog(LOG_NOTICE, "net_get_modbus_task start, channel:%s", channel->channel);
        dev_poll_sum = find_modbus_meter_dev_poll(channel, &dev_poll);
        if(dev_poll != NULL)
        {
            while (1)
            {
                len = 0;
                if (read_protocol_data(raw, (char *)&rx_pack, &len) == 0)
                {
                    if (len > 0)
                    {
                        proto_syslog_hex(LOG_DEBUG, (char *)&rx_pack, len, "get modbus broadcast data:");
                        proto_syslog(LOG_NOTICE, "poll_index=%d,dev_addr=%d,data_len=%d", rx_pack.poll_index, rx_pack.dev_addr, rx_pack.data_len);
                        for (int i = 0; i < dev_poll_sum; i++)
                        {
                            if ((dev_poll[i].dev_addr == rx_pack.dev_addr) && (rx_pack.poll_index < channel->chunks_interv_size))
                            {
                                poll_info = (modbus_poll_info_t *)channel->chunks_interv[rx_pack.poll_index]->chunk_proto_ptr;
                                modbus_update_tags_new(poll_info, rx_pack.data, rx_pack.data_len);
                                refresh_online_state(poll_info->p_tag[0]->dev_ptr);
                                err_counter = 0;
                            }
                        }
                    }
                    check_tag_state(channel);
                }
                else
                {
                    cur_t = time(NULL);
                    if ((cur_t - last_t) > 0) // 每一秒检测一次
                    {
                        last_t = cur_t;
                        for (int i = 0; i < dev_poll_sum; i++)
                        {
                            check_online_state(dev_poll[i].poll[0]->p_tag[0]->dev_ptr);
                        }
                        check_tag_state(channel);
                        err_counter++;
                        if (err_counter > ch_broadcast_timeout)
                        {
                            break;
                        }
                    }
                }
            }

            for (int i = 0; i < dev_poll_sum; i++)
            {
                delete_dev_poll(dev_poll + i);
            }
            free(dev_poll);
        }
        raw_config_deinit(raw);
    }

    return -1;
}

#endif

/* Raw data callback functions */
static void dlt645_raw_send_callback(dlt645_hd_t *hd, const uint8_t *data, int data_length, void *channel)
{
    channel_t *ch = channel;
    
    /* Fill send_buf in channel structure */
    if (ch->send_buf != NULL) {
        free(ch->send_buf);
    }
    ch->send_buf = malloc(data_length);
    if (ch->send_buf != NULL) {
        memcpy(ch->send_buf, data, data_length);
        ch->send_buf_size = data_length;
    }
    
}

static void dlt645_raw_receive_callback(dlt645_hd_t *hd, const uint8_t *data, int data_length, void *channel)
{
    channel_t *ch = channel;
    
    /* Fill recv_buf in channel structure */
    if (ch->recv_buf != NULL) {
        free(ch->recv_buf);
    }
    ch->recv_buf = malloc(data_length);
    if (ch->recv_buf != NULL) {
        memcpy(ch->recv_buf, data, data_length);
        ch->recv_buf_size = data_length;
    }
    
}
static int dlt645_channel(channel_t *channel)
{
    long long cur_ms = 0;
    int err_counter = 0;
    dlt645_poll_info_t *poll_info = NULL;
    uint8_t rx_buff[256] = {0};
    device_t *dev = NULL;
    dlt645_hd_t *hd = NULL;
    if (channel->proto_type == PROTOCOL_TYPE_DLT645_1997)
    {
        hd = dlt645_proto_init((dlt645_proto_t *)channel->proto_ptr, DLT645_1997);
    }
    else
    {
        hd = dlt645_proto_init((dlt645_proto_t *)channel->proto_ptr, DLT645_2007);
    }
    if (hd == NULL)
    {
        proto_syslog(LOG_CRIT, "[数采]:dlt645_proto_init error in the channel:[%s]", channel->channel);
        BUSINESS_LOG(BLOG_ERR, DAS_ID_DLT645_1, channel->channel, "[数采]: %s 通道初始化失败", channel->channel);
        return -1;
    }
    RAW_SOCK *raw = channel->some_ptr;
    if (raw == NULL)
    {
        if ((raw = protocol_Init_with_group(get_pcs_ctrl_var()->net_node,  get_cnt_id(), MAC_CHANNEL_NAME))==NULL)
        {
            proto_syslog(LOG_WARNING, "creat error!");
        }
        channel->some_ptr = raw;
    }
    dlt645_set_raw_send_callback(hd, dlt645_raw_send_callback,(void *)channel);
    dlt645_set_raw_receive_callback(hd, dlt645_raw_receive_callback,(void *)channel);
    channel->s = hd->fd;
    proto_syslog(LOG_INFO, "dlt645_proto_init success, in the channel:[%s]", channel->channel);
    BUSINESS_LOG(BLOG_NOTICE, DAS_ID_DLT645_1, channel->channel, "[数采]: %s 通道初始化成功", channel->channel);
    while (1)
    {
        uint8_t contine_interv = 0;
        dlt645_check_write(hd, channel, 20);
        do
        {
            contine_interv = 0;
            cur_ms = clock_get_ms();
            check_tag_state(channel);
            for (size_t i = 0; i < channel->chunks_interv_size; i++)
            {
                dlt645_check_write(hd, channel, 0);
                dev = channel->chunks_interv[i]->dev_ptr;
                poll_info = (dlt645_poll_info_t *)channel->chunks_interv[i]->chunk_proto_ptr;
                if ((cur_ms - channel->chunks_interv[i]->last_poll) >= channel->chunks_interv[i]->poll_interv)
                {
                    channel->chunks_interv[i]->last_poll = cur_ms;
                    dlt645_set_time(hd, dev->timeout, dev->interval);
                    proto_syslog(LOG_DEBUG, "[channel:%s][dev:%s] dlt645 read dev_addr = %s, data_id=0x%08x, interv=%u, cur_ms=%lld",
                              channel->channel, dev->no, dev->addr, poll_info->data_id, channel->chunks_interv[i]->poll_interv, cur_ms);

                    channel_lock(channel);
                    if (hd->frame_gap != dev->write_gap)
                        hd->frame_gap = dev->write_gap;

                    void *send_protocol_data_cb_t = NULL;
                    if((dev->dev_child_type!=NULL)&&((strcmp(dev->dev_child_type,"GRID") == 0)||(strcmp(dev->dev_child_type,"TRAS_LV") == 0)))
                        send_protocol_data_cb_t = send_protocol_data_cb;
                    int ret = dlt645_proto_read(hd, ((dlt645_dev_attr *)dev->proto_ptr)->bcd_addr,send_protocol_data_cb_t,raw,poll_info, rx_buff, sizeof(rx_buff)); // 读数据
                    channel_unlock(channel);
                    if (ret > 0)
                    {
                        BUSINESS_LOG(BLOG_NOTICE, DAS_ID_DLT645_2, channel->channel, "[数采]: %s 通道回复正常", channel->channel);
                        // proto_syslog(LOG_DEBUG, "[channel:%s][dev:%s] dlt645 TAGS:", channel->channel, dev->no);
                        dlt645_update_tags(dev, poll_info, rx_buff, ret, channel);
                        refresh_online_state(dev);
                        err_counter = 0;
                    }
                    else if (ret == 0)
                    {
                        proto_syslog(LOG_ERR, "%s:dlt645 read timeout", channel->channel);
                        BUSINESS_LOG(BLOG_ERR, DAS_ID_DLT645_2, channel->channel, "[数采]: %s 通道回复超时", channel->channel);
                        err_counter++;
                        check_online_state(dev);
                    }
                    else
                    {
                        proto_syslog(LOG_ERR, "%s:dlt645 read error(ret = %d), info:%s, errno:%d", channel->channel, ret, strerror(errno), errno);
                        BUSINESS_LOG(BLOG_ERR, DAS_ID_DLT645_2, channel->channel, "[数采]: %s 通道回复错误, 错误原因(%d)", channel->channel, strerror(errno));
                        check_online_state(dev);
                        goto dlt645_channel_exit;
                    }
                    contine_interv = 1;
                }

                if (err_counter > COMMUNIT_RETRY_MAX)
                {
                    goto dlt645_channel_exit;
                }
            }
        } while ((contine_interv != 0));

        // interv空闲
        if (channel->chunks_nonstop_size > 0) // 有配置
        {
            size_t nonstop_index = channel->chunks_nonstop_index;
            dev = channel->chunks_nonstop[nonstop_index]->dev_ptr;
            poll_info = (dlt645_poll_info_t *)channel->chunks_nonstop[nonstop_index]->chunk_proto_ptr;
            proto_syslog(LOG_DEBUG, "[channel:%s][dev:%s] dlt645 read dev_addr = %s, data_id=0x%08x, interv=%u, cur_ms=%llu",
                      channel->channel, dev->no, dev->addr, poll_info->data_id, channel->chunks_interv[nonstop_index]->poll_interv, cur_ms);
            channel_lock(channel);
            if (hd->frame_gap != dev->write_gap)
                hd->frame_gap = dev->write_gap;
            int ret = dlt645_proto_read(hd, ((dlt645_dev_attr *)dev->proto_ptr)->bcd_addr,send_protocol_data_cb,raw,poll_info, rx_buff, sizeof(rx_buff)); // 读数据
            channel_unlock(channel);
            if (ret > 0)
            {
                BUSINESS_LOG(BLOG_NOTICE, DAS_ID_DLT645_2, channel->channel, "[数采]: %s 通道回复正常", channel->channel);
                // proto_syslog(LOG_DEBUG, "[channel:%s][dev:%s] dlt645 TAGS:", channel->channel, dev->no);
                dlt645_update_tags(dev, poll_info, rx_buff, ret, channel);
                refresh_online_state(dev);
                err_counter = 0;
            }
            else if (ret == 0)
            {
                proto_syslog(LOG_ERR, "%s:dlt645 read timeout", channel->channel);
                BUSINESS_LOG(BLOG_ERR, DAS_ID_DLT645_2, channel->channel, "[数采]: %s 通道回复超时", channel->channel);
                err_counter++;
                check_online_state(dev);
            }
            else
            {
                proto_syslog(LOG_ERR, "%s:dlt645 read error(ret = %d), info:%s, errno:%d", channel->channel, ret, strerror(errno), errno);
                BUSINESS_LOG(BLOG_ERR, DAS_ID_DLT645_2, channel->channel, "[数采]: %s 通道回复错误, 错误原因(%d)", channel->channel, strerror(errno));
                check_online_state(dev);
                if (errno > 0)//io错误
                {
                    BUSINESS_LOG(BLOG_ERR, DAS_ID_DLT645_2, channel->channel, "[数采]: %s 通道出现系统级IO错误, 请检查 物理连接或设备 是否正常", channel->channel);
                    goto dlt645_channel_exit;
                }
                else//协议错误
                {
                    BUSINESS_LOG(BLOG_ERR, DAS_ID_DLT645_2, channel->channel, "[数采]: %s 通道出现协议错误, 请检查 回复数据 是否正确", channel->channel);
                    err_counter++;
                }
            }

            if (err_counter > COMMUNIT_RETRY_MAX)
            {
                goto dlt645_channel_exit;
            }

            nonstop_index++;
            if (nonstop_index >= channel->chunks_nonstop_size)
            {
                channel->chunks_nonstop_index = 0;
            }
            else
            {
                channel->chunks_nonstop_index = nonstop_index;
            }
        }
    }
dlt645_channel_exit:
    if (hd != NULL)
    {
        dlt645_proto_deinit(hd);
    }
    hd = NULL;

    return -1;
}

static int ups_update_tags(device_t *dev, int poll_index, const char *cmd, char *rx_str)
{
    tag_t *p_tag = NULL;
    uint64_t cur_ms = clock_get_ms();
    tag_value_t tag_value = {.type = 0};
    ups_read_info_t *ext_info = NULL;

    char *s = rx_str + 1; // 第一个字符忽略
    char *arr[CMD_MAX_REG_SUM] = {0};
    int it_sum = 0;

    while ((s != NULL) && (*s != '\0'))
    {
        if (it_sum < CMD_MAX_REG_SUM)
        {
            arr[it_sum++] = strsep(&s, " ");
        }
        else
        {
            proto_syslog(LOG_WARNING, "ups cmd[%s] receive data split error", cmd);
        }
    }

    for (size_t i = 0; i < dev->tags.ro_tags_size; i++)
    {
        ext_info = (ups_read_info_t *)dev->tags.ro_tags[i]->proto_ptr;
        p_tag = dev->tags.ro_tags[i];
        //if (strcmp(ext_info->cmd, cmd) == 0)
        if (poll_index == p_tag->poll_list_num)
        {
            if (ext_info->pos < it_sum)
            {
                switch (ext_info->type)
                {
                case UPS_DATA_FLOAT:
                    tag_value_t.type = 1;
                    sscanf(arr[ext_info->pos], "%lf", &tag_value_t.value.to_float);
                    proto_set_tag(p_tag, tag_value, cur_ms);
                    break;

                case UPS_DATA_BIT:
                    tag_value_t.type = 0;
                    tag_value_t.value.to_int = arr[ext_info->pos][7 - ext_info->bit_pos] - '0';
                    proto_set_tag(p_tag, tag_value, cur_ms);
                    break;

                default:
                    proto_syslog(LOG_ERR, "TAG:%s, ext type error, type=%d", p_tag->name, ext_info->type);
                    break;
                }
            }
            else
            {
                proto_syslog(LOG_WARNING, "TAG:%s, index warning, index=%d", p_tag->name, ext_info->pos);
            }
        }
    }
    return 0;
}

static int ups_write_tag(ups_hd_t *hd, tag_t *tag)
{
    char cmd[32] = {0};
    ups_write_info_t *write_info = (ups_write_info_t *)tag->proto_ptr;

    if (write_info->need_par != 0)
    {
        if (tag->data_type == 0)
        {
            snprintf(cmd, sizeof(cmd), "%s%ld", write_info->cmd, tag->write_cache.to_int);
        }
        else
        {
            snprintf(cmd, sizeof(cmd), "%s%.1f", write_info->cmd, tag->write_cache.to_float);
        }

        return ups_write(hd, cmd);
    }
    else
    {
        return ups_write(hd, write_info->cmd);
    }
}

static int ups_check_write(ups_hd_t *hd, channel_t *channel, int time_out)
{
    device_t *dev = NULL;

    int ret = sem_wait_time(&channel->write_sem, 0, time_out);
    if (ret == 0)
    {
        do
        {
            write_list_t *w_info = get_data_from_list(channel);
            if (w_info == NULL)
            {
                proto_syslog(LOG_ERR, "get data from list is null");
                return -1;
            }
            tag_t *write_tag = &w_info->w_tag;
            if (write_tag == NULL)
            {
                proto_syslog(LOG_ERR, "get data from list is null");
                free(w_info);
                return -2;
            }
            dev = (device_t *)write_tag->dev_ptr;
            proto_syslog(LOG_INFO, "get data from list is [%s.%s(type:%d, double=%f, long=%ld)]",
                         dev->no, write_tag->name, write_tag->data_type, write_tag->write_cache.to_float, write_tag->write_cache.to_int);
            ups_set_time(hd, dev->timeout, dev->interval);
            channel_lock(channel);
            ret = ups_write_tag(hd, write_tag);
            channel_unlock(channel);
            if (0 != ret)
            {
                w_info->pw_tag->last_state_code = w_info->pw_tag->state_code;
                w_info->pw_tag->state_code = TAG_STATE_BAD;
                // write_tag->write_cache.to_float = write_tag->read_cache.to_float; // 写入错误，将写cache写回去，方便上位机查看
                //  写入错误没有其他处理,后期把质量码加上,方便知道数据是否可用
                check_online_state(dev);
                free(w_info);
                proto_syslog(LOG_ERR, "ups write tag error, tag:%s, ret = %d", write_tag->name, ret);
                BUSINESS_LOG(BLOG_ERR, DAS_ID_UPS_3, write_tag->dev_tag_name, "[数采]: %s 通道写 tag: %s 失败", write_tag->name);
                return -3;
            }
            else
            {
                BUSINESS_LOG(BLOG_NOTICE, DAS_ID_UPS_3, write_tag->dev_tag_name, "[数采]: %s 通道写 tag: %s 成功", channel->channel, write_tag->name);
                w_info->pw_tag->last_state_code = w_info->pw_tag->state_code;
                w_info->pw_tag->state_code = TAG_STATE_GOOD;
                sync_wtag_read_cache(w_info);
            }
            free(w_info);
        } while (sem_wait_time(&channel->write_sem, 0, 0) == 0);
    }
    return 0;
}

static int ups_channel(channel_t *channel)
{
    long long cur_ms = 0;
    int err_counter = 0;
    ups_poll_info_t *poll_info = NULL;
    char rx_buff[200] = {0};
    device_t *dev = NULL;
    ups_hd_t *hd = ups_proto_init((ups_proto_t *)channel->proto_ptr);

    if (hd == NULL)
    {
        proto_syslog(LOG_CRIT, "[数采]:dlt645_proto_init error in the channel:[%s]", channel->channel);
        BUSINESS_LOG(BLOG_ERR, DAS_ID_UPS_1, channel->channel, "[数采]: %s 通道初始化失败", channel->channel);
        return -1;
    }
    channel->s = hd->fd;
    proto_syslog(LOG_INFO, "dlt645_proto_init success, in the channel:[%s]", channel->channel);
    BUSINESS_LOG(BLOG_NOTICE, DAS_ID_UPS_1, channel->channel, "[数采]: %s 通道初始化成功", channel->channel);

    while (1)
    {
        uint8_t contine_interv = 0;
        do
        {
            check_tag_state(channel);
            if(ups_check_write(hd, channel, 100)<0)
                err_counter++;
            contine_interv = 0;
            cur_ms = clock_get_ms();
            for (size_t i = 0; i < channel->chunks_interv_size; i++)
            {
                if(ups_check_write(hd, channel, 0)<0)
                    err_counter++;
                dev = channel->chunks_interv[i]->dev_ptr;
                poll_info = (ups_poll_info_t *)channel->chunks_interv[i]->chunk_proto_ptr;
                if ((cur_ms - channel->chunks_interv[i]->last_poll) >= channel->chunks_interv[i]->poll_interv)
                {
                    channel->chunks_interv[i]->last_poll = cur_ms;
                    ups_set_time(hd, dev->timeout, dev->interval);
                    proto_syslog(LOG_DEBUG, "[channel:%s][dev:%s] ups read cmd = %s, interv=%u, cur_ms=%lld",
                              channel->channel, dev->no, poll_info->cmd, channel->chunks_interv[i]->poll_interv, cur_ms);
                    channel_lock(channel);
                    int ret = ups_read(hd, poll_info->cmd, rx_buff, sizeof(rx_buff));
                    channel_unlock(channel);
                    if (ret > 0)
                    {
                        BUSINESS_LOG(BLOG_NOTICE, DAS_ID_UPS_2, channel->channel, "[数采]: %s 通道回复正常", channel->channel);
                        // proto_syslog(LOG_DEBUG, "[channel:%s][dev:%s] dlt645 TAGS:", channel->channel, dev->no);
                        ups_update_tags(dev, i, poll_info->cmd, rx_buff);
                        refresh_online_state(dev);
                        err_counter = 0;
                    }
                    else if (ret == 0)
                    {
                        proto_syslog(LOG_ERR, "%s:ups read timeout", channel->channel);
                        BUSINESS_LOG(BLOG_ERR, DAS_ID_UPS_2, channel->channel, "[数采]: %s 通道回复超时", channel->channel);
                        err_counter++;
                        check_online_state(dev);
                    }
                    else
                    {
                        proto_syslog(LOG_ERR, "%s:ups read error(ret = %d), info:%s, errno:%d", channel->channel, ret, strerror(errno), errno);
                        BUSINESS_LOG(BLOG_ERR, DAS_ID_UPS_2, channel->channel, "[数采]: %s 通道回复失败", channel->channel);
                        check_online_state(dev);
                        if (errno > 0) // io错误
                        {
                            BUSINESS_LOG(BLOG_ERR, DAS_ID_UPS_2, channel->channel, "[数采]: %s 通道出现系统级IO错误, 请检查 物理连接或设备 是否正常", channel->channel);			
                            goto ups_channel_exit;
                        }
                        else // 协议错误
                        {
                            BUSINESS_LOG(BLOG_ERR, DAS_ID_UPS_2, channel->channel, "[数采]: %s 通道出现协议错误, 请检查 回复数据 是否正确", channel->channel);
                            err_counter++;
                        }
                    }
                    contine_interv = 1;
                }

                if (err_counter > COMMUNIT_RETRY_MAX)
                {
                    goto ups_channel_exit;
                }
            }
        } while ((contine_interv != 0));
    }
ups_channel_exit:
    if (hd != NULL)
    {
        ups_free(hd);
    }
    hd = NULL;

    return -1;
}

#define EMS_DEV_INDEX 0 //ems在通道中的索引 
channel_t *ems_channel_pr;
void reflush_all_ems_data()
{
    channel_t *channel = ems_channel_pr;
    uint64_t cur_ms = clock_get_ms();
    for (size_t i = 0; i < channel->devs[EMS_DEV_INDEX]->tags.ro_tags_size; i++)
    {
        tag_t* tag = channel->devs[EMS_DEV_INDEX]->tags.ro_tags[i];
        if (tag->data_type == TYPE_TAG_FLOAT)
        {
            tag->last_read_cache.to_float = tag->read_cache.to_float;
        }
        else
        {
            tag->last_read_cache.to_int = tag->read_cache.to_int;
        }
        tag->last_read_time = tag->read_cache_time;
        tag->read_cache_time = cur_ms;
        reflush_data_up(tag);
    }
    // for (size_t i = 0; i < channel->devs[EMS_DEV_INDEX]->tags.rw_tags_size; i++)
    // {
    //     tag_t* tag = channel->devs[EMS_DEV_INDEX]->tags.rw_tags[i];
    //     if (tag->data_type == TYPE_TAG_FLOAT)
    //     {
    //         tag->last_read_cache.to_float = tag->read_cache.to_float;
    //         tag->read_cache.to_float = tag->write_cache.to_float;
    //     }
    //     else
    //     {
    //         tag->last_read_cache.to_int = tag->read_cache.to_int;
    //         tag->read_cache.to_int = tag->write_cache.to_int;
    //     }
    //     tag->last_read_time = tag->read_cache_time;
    //     tag->read_cache_time = cur_ms;
    //     reflush_data_up(tag);
    // }
    //refresh_online_state(channel->devs[EMS_DEV_INDEX]);
}

int ems_channel_is_idle(void)
{
    channel_t *chan_ptr = ems_channel_pr;
    while(chan_ptr == NULL)
    {
        chan_ptr = ems_channel_pr;
        ems_syslog(LOG_WARNING, "ems_channel_pr is null");
        sleep(1);
    }
    int ret = 0;
    pthread_rwlock_wrlock(&chan_ptr->list_rwlock);
    ret = list_empty(&chan_ptr->wirte_list);
    pthread_rwlock_unlock(&chan_ptr->list_rwlock);
    return ret;
}

static int ems_channel(void *channel_context) //ems的通道 
{

    channel_t *channel = channel_context;
    ems_channel_pr = channel;
    time_t now = time(NULL);
    time_t latest_time = now;
    proto_syslog(LOG_NOTICE, "ems data in the channel:[%s]", channel->channel);
    refresh_online_state(channel->devs[EMS_DEV_INDEX]);
    check_online_state(channel->devs[EMS_DEV_INDEX]);
    while (1)
    {
        int ret = sem_wait_time(&channel->write_sem, 10, 0);//TODO 这个只做写的话时间可以长点 如果需要定时刷新数据把时间改成需要的 但是再读的时候记得
        if (ret == 0)
        {
            write_list_t *w_info = get_data_from_list(channel);
            if (w_info == NULL)
            {
                proto_syslog(LOG_ERR, "get data from list is null");
                continue;
            }
            tag_t *write_tag = &w_info->w_tag;
            if (write_tag != NULL)
            {
                sync_wtag_read_cache(w_info);
            }
            free(w_info);
        }
        else
        {
            // for (size_t i = 0; i < channel->devs[EMS_DEV_INDEX]->tags.ro_tags_size; i++)
            // {
            //     uint64_t cur_ms = clock_get_ms();
            //     channel->devs[EMS_DEV_INDEX]->tags.ro_tags[i]->read_cache_time = cur_ms;
            //     reflush_data_up(channel->devs[EMS_DEV_INDEX]->tags.ro_tags[i]);
            // }
            // for (size_t i = 0; i < channel->devs[EMS_DEV_INDEX]->tags.rw_tags_size; i++)
            // {
            //     uint64_t cur_ms = clock_get_ms();
            //     channel->devs[EMS_DEV_INDEX]->tags.rw_tags[i]->read_cache_time = cur_ms;
            //     reflush_data_up(channel->devs[EMS_DEV_INDEX]->tags.rw_tags[i]);
            // }
        }
        //
        now = time(NULL);
        if (latest_time != now)
        {
            latest_time = now;
            refresh_online_state(channel->devs[EMS_DEV_INDEX]);
            check_online_state(channel->devs[EMS_DEV_INDEX]);
        }
    }
    return 0;
}

static int lc_channel(void *channel_context) // lc的通道
{

    tag_t *write_tag = NULL;
    time_t now = time(NULL);
    time_t latest_time = now;
    channel_t *channel = channel_context;
    ems_channel_pr = channel;
    proto_syslog(LOG_NOTICE, "lc data in the channel:[%s]", channel->channel);
    for (int i = 0; i < channel->devs_size; i++)
    {
        refresh_online_state(channel->devs[i]);
        check_online_state(channel->devs[i]);
        proto_syslog(LOG_NOTICE, "refresh_online_state:[%s]", channel->devs[i]->no);
    }
    while (1)
    {
        int ret = sem_wait_time(&channel->write_sem, 1, 0);
        if (ret == 0)
        {
            write_list_t *w_info = get_data_from_list(channel);
            if (w_info == NULL)
            {
                proto_syslog(LOG_ERR, "get data from list is null");
                continue;
            }
            write_tag = &w_info->w_tag;
            proto_syslog(LOG_INFO, "get data from list is [%s.%s(type:%d, double=%f, int=%ld)]",
                         "LC", write_tag->name, write_tag->data_type, write_tag->write_cache.to_float, write_tag->write_cache.to_int);
            if (write_tag != NULL)
            {
                sync_wtag_read_cache(w_info);
            }
            free(w_info);
        }
        //
        now = time(NULL);
        if (latest_time != now)
        {
            latest_time = now;
            for (int i = 0; i < channel->devs_size; i++)
            {
                refresh_online_state(channel->devs[i]);
                check_online_state(channel->devs[i]);
            }
        }
    }
    return 0;
}

int bmser_data_report(const bmser_cb_par *p)
{
    device_t *dev_ptr = NULL;
    channel_t *channel = (channel_t *)p->top;
    const block_data *block = p->proto_data;
    uint32_t block_id = TOGGLE_ENDIAN_32(block->block_id);
    uint16_t start_sub_id = TOGGLE_ENDIAN_16(block->start_sub_id);
    uint16_t bytes = TOGGLE_ENDIAN_16(block->bytes);
    uint16_t tmp = 0;
    proto_syslog(LOG_INFO, "channel:[%s] blockId:0x%08x, startSubId=0x%04x, SubNum:0x%04x, bytes:0x%04x\n", channel->channel, block_id, start_sub_id,
              TOGGLE_ENDIAN_16(block->sub_num), bytes);
    bmser_poll_info_t *poll = NULL;
    bmser_parser_info_t *parse_info = NULL;
    tag_value_t tag_value = {.type = 0};
    
    for (int i = 0; i < channel->chunks_interv_size; i++)
    {
        poll = (bmser_poll_info_t *)channel->chunks_interv[i]->chunk_proto_ptr;
        //proto_syslog(LOG_DEBUG, "blockId:0x%08x, startSubId=0x%04x", poll->block_id, poll->sub_id);
        if ((poll->block_id == block_id) && (poll->sub_id == start_sub_id))
        {
            proto_syslog_hex(LOG_INFO, block->content, bytes, "block->content:%d", bytes);
            for (int t = 0; t < poll->tag_sum; t++)//后面使用bmser_bin_extra替换掉
            {
                parse_info = (bmser_parser_info_t *)poll->p_tag[t]->proto_ptr;
                // proto_syslog(LOG_DEBUG, "block_id=0x%08x, sub_id=0x%04x, byte_index = %d, type=%d", parse_info->block_id, parse_info->sub_id, parse_info->byte_index, parse_info->type);
                switch (parse_info->type)
                {
                case BMSER_DATA_U16:
                    tag_value.value.to_int = (uint16_t)GET_FOR_BA((&block->content[parse_info->byte_index]));
                    break;
                case BMSER_DATA_S16:
                    tag_value.value.to_int = (int16_t)GET_FOR_BA((&block->content[parse_info->byte_index]));
                    break;
                case BMSER_DATA_U32:
                    tag_value.value.to_int = (uint32_t)GET_FOR_BADC((&block->content[parse_info->byte_index]));
                    break;
                case BMSER_DATA_S32:
                    tag_value.value.to_int = (int32_t)GET_FOR_BADC((&block->content[parse_info->byte_index]));
                    break;
                case BMSER_DATA_BIT16:
                    tmp = GET_FOR_BA((&block->content[parse_info->byte_index]));
                    tag_value.value.to_int = GET_BITS_IN_U16(tmp, parse_info->bit_pos, parse_info->bit_len);
                    break;
                default:
                    proto_syslog(LOG_ERR, "parse_info->type error, type = %d", parse_info->type);
                    return -1;
                    break;
                }
                proto_syslog(LOG_INFO, "TAG:%s->byte_index:0x%x[%d]", poll->p_tag[t]->name, parse_info->byte_index, parse_info->byte_index);
                proto_set_tag(poll->p_tag[t], tag_value, clock_get_ms());
                dev_ptr = (device_t*)poll->p_tag[t]->dev_ptr;
            }

            if (dev_ptr != NULL)
            {
                refresh_online_state(dev_ptr);
            }
        }
    }
    
    return 0;
}

static int tag_to_bmser_buff(tag_t *tag, void *buff, int max_len)
{
    uint8_t *data_buff = (uint8_t *)buff;
    int len = 0;
    bmser_parser_info_t *parse_info = (bmser_parser_info_t *)tag->proto_ptr;
    uint32_t value = (tag->data_type == 0) ? (tag->write_cache.to_int - tag->offset) / tag->scale : (tag->write_cache.to_float - tag->offset) / tag->scale;

    switch (parse_info->type)
    {
    case BMSER_DATA_U16:
    case BMSER_DATA_S16:
        SET_FOR_BA((uint16_t)value, data_buff);
        len = 2;
        break;

    case BMSER_DATA_U32:
    case BMSER_DATA_S32:
        SET_FOR_BADC(value, data_buff);
        len = 4;
        break;

    case BMSER_DATA_BIT16:
        proto_syslog(LOG_WARNING, "type is not supported, type = %d", parse_info->type);
        break;

    default:
        proto_syslog(LOG_ERR, "type is not supported, type = %d", parse_info->type);
        return -1;
        break;
    }

    return len;
}

int bmser_write_tag(bmser_protocol *bmser, tag_t *tag)
{
    static uint32_t trans_id = 0x12345678;
    uint8_t data_buff[4] = {0};
    int len = 0;
    block_data *block = NULL;
    bmser_parser_info_t *parse_info = (bmser_parser_info_t *)tag->proto_ptr;
    device_t *dev = (device_t *)tag->dev_ptr;

    len = tag_to_bmser_buff(tag, data_buff, sizeof(data_buff));
    if (len <= 0)
    {
        proto_syslog(LOG_ERR, "tag to bmser buff error,tag:%s", tag->name);
        return -3;
    }

    block = new_block_data(parse_info->block_id, parse_info->sub_id, data_buff, 1, len);
    if (block == NULL)
    {
        proto_syslog(LOG_ERR, "new bmser block data err");
        return -1;
    }

    bmser_bin_pack pp = {.trans_id = trans_id++, .dev_len = 0, .cmd_len = 0};
    uint8_t *pack = NULL;
    int ret_len = new_bmser_proto_pack(&pp, &block, 1, &pack);
    free(block);
    if (ret_len > 0)
    {
        bmser_dev_attr *dev_attr = (bmser_dev_attr *)dev->proto_ptr;
        bmser_protocol_send_pack(bmser, dev_attr->addr.str, pack, ret_len);
        proto_syslog_hex(LOG_NOTICE, pack, ret_len, "bmser tx:%d", ret_len);
        free(pack);
        return 0;
    }
    else
    {
        proto_syslog(LOG_ERR, "new bmser proto pack error");
        return -2;
    }
}

static int bmser_bin_channel(channel_t *channel)
{
    device_t *dev = NULL;
    bmser_proto_s *proto_if_info = (bmser_proto_s*)channel->proto_ptr;
    bmser_protocol *bmser = new_bmser_protocol(proto_if_info->ip, proto_if_info->port, bmser_data_report, channel);
    if (bmser == NULL)
    {
        proto_syslog(LOG_ERR, "new_bmser_protocol error");
        BUSINESS_LOG(BLOG_ERR, DAS_ID_BMSER_BIN_1, channel->channel, "[数采]: %s 通道创建失败", channel->channel);
        return -1;
    }

    for (int i = 0; i < channel->devs_size; i++)
    {
        dev = channel->devs[i];
        bmser_dev_attr *dev_attr = (bmser_dev_attr*)dev->proto_ptr;
        if (0 == bmser_protocol_add_dev(bmser, dev_attr->addr.str))
        {
            proto_syslog(LOG_NOTICE, "bmser protocol add dev success,dev:%s, dev_addr:%s", dev->no, dev_attr->addr.str);
            BUSINESS_LOG(BLOG_NOTICE, DAS_ID_BMSER_BIN_2, dev->no, "[数采]: %s 通道增加设备 %s 成功", channel->channel, dev->no);
        }
        else
        {
            proto_syslog(LOG_NOTICE, "bmser protocol add dev failed,dev:%s", dev->no);
            BUSINESS_LOG(BLOG_ERR, DAS_ID_BMSER_BIN_2, dev->no, "[数采]: %s 通道增加设备 %s 失败", channel->channel, dev->no);
        }
    }

    proto_syslog(LOG_ERR, "new_bmser_protocol success, bmser protocol start");
    BUSINESS_LOG(BLOG_NOTICE, DAS_ID_BMSER_BIN_1, channel->channel, "[数采]: %s 通道创建成功", channel->channel);
    if (bmser_protocol_start(bmser) < 0)
    {
        proto_syslog(LOG_ERR, "bmser protocol start error");
        BUSINESS_LOG(BLOG_ERR, DAS_ID_BMSER_BIN_3, channel->channel, "[数采]: %s 通道启动失败", channel->channel);
    }
    else
    {
        proto_syslog(LOG_NOTICE, "bmser protocol start success");
        BUSINESS_LOG(BLOG_NOTICE, DAS_ID_BMSER_BIN_3, channel->channel, "[数采]: %s 通道启动成功", channel->channel);
        while (1)
        {
            check_tag_state(channel);
            int ret = sem_wait_time(&channel->write_sem, 1, 100);
            if (ret == 0)
            {
                write_list_t *w_info = get_data_from_list(channel);
                if (w_info == NULL)
                {
                    proto_syslog(LOG_ERR, "get data from list is null");
                    return -1;
                }
                tag_t *write_tag = &w_info->w_tag;
                if (write_tag == NULL)
                {
                    proto_syslog(LOG_ERR, "get data from list is null");
                    free(w_info);
                    return -2;
                }
                dev = (device_t *)write_tag->dev_ptr;
                proto_syslog(LOG_INFO, "get data from list is [%s.%s(type:%d, double=%f, long=%ld)]",
                             dev->no, write_tag->name, write_tag->data_type, write_tag->write_cache.to_float, write_tag->write_cache.to_int);
                ret = bmser_write_tag(bmser, write_tag);
                if (0 != ret)
                {
                    w_info->pw_tag->last_state_code = w_info->pw_tag->state_code;
                    w_info->pw_tag->state_code = TAG_STATE_BAD;
                    // write_tag->write_cache.to_float = write_tag->read_cache.to_float; // 写入错误，将写cache写回去，方便上位机查看
                    //  写入错误没有其他处理,后期把质量码加上,方便知道数据是否可用
                    proto_syslog(LOG_ERR, "ups write tag error, tag:%s, ret = %d", w_info->pw_tag->name, ret);
                    BUSINESS_LOG(BLOG_ERR, DAS_ID_BMSER_BIN_4, write_tag->dev_tag_name, "[数采]: %s 通道写 tag: %s 失败", channel->channel, w_info->pw_tag->name);
                    free(w_info);
                    return -3;
                }
                else
                {
                    BUSINESS_LOG(BLOG_NOTICE, DAS_ID_BMSER_BIN_4, write_tag->dev_tag_name, "[数采]: %s 通道写 tag: %s 成功", channel->channel, w_info->pw_tag->name);
                    w_info->pw_tag->last_state_code = w_info->pw_tag->state_code;
                    w_info->pw_tag->state_code = TAG_STATE_GOOD;
                    sync_wtag_read_cache(w_info);
                }
                free(w_info);
            }
        }
    }
    delete_bmser_protocol(bmser);
    return -1;
}
#ifdef EN_BMSER_DB_CTROL_LIB

static int bmser_bin_extra(uint32_t dev_addr, bmser_poll_info_t const *poll, void *buff, int maxLen, const uint8_t *valid_flag)
{
    tag_value_t tag_value = {.type = 0};
    bmser_parser_info_t *parse_info = NULL;
    uint16_t tmp = 0;
    uint8_t *content = (uint8_t *)buff;
    for (int t = 0; t < poll->tag_sum; t++)
    {
        parse_info = (bmser_parser_info_t *)poll->p_tag[t]->proto_ptr;
        if (valid_flag[parse_info->byte_index / 2] == 0) // 数据无效
        {
            proto_syslog(LOG_NOTICE, "bmser extra notice dev_addr:0x%x, sub_id:0x%x", dev_addr, parse_info->sub_id);
            continue;
        }
        // proto_syslog(LOG_DEBUG, "block_id=0x%08x, sub_id=0x%04x, byte_index = %d, type=%d", parse_info->block_id, parse_info->sub_id, parse_info->byte_index, parse_info->type);
        switch (parse_info->type)
        {
        case BMSER_DATA_U16:
            tag_value.value.to_int = (uint16_t)GET_FOR_BA(content + parse_info->byte_index);
            break;
        case BMSER_DATA_S16:
            tag_value.value.to_int = (int16_t)GET_FOR_BA(content + parse_info->byte_index);
            break;
        case BMSER_DATA_U32:
            tag_value.value.to_int = (uint32_t)GET_FOR_BADC(content + parse_info->byte_index);
            break;
        case BMSER_DATA_S32:
            tag_value.value.to_int = (int32_t)GET_FOR_BADC(content + parse_info->byte_index);
            break;
        case BMSER_DATA_BIT16:
            tmp = GET_FOR_BA(content + parse_info->byte_index);
            tag_value.value.to_int = GET_BITS_IN_U16(tmp, parse_info->bit_pos, parse_info->bit_len);
            break;
        default:
            proto_syslog(LOG_ERR, "parse_info->type error, type = %d", parse_info->type);
            return -1;
            break;
        }
        proto_syslog(LOG_INFO, "TAG:%s->byte_index:0x%x[%d]", poll->p_tag[t]->name, parse_info->byte_index, parse_info->byte_index);
        proto_set_tag(poll->p_tag[t], tag_value, clock_get_ms());
    }

    // dev_ptr = (device_t *)poll->p_tag[0]->dev_ptr;
    // if (dev_ptr != NULL)
    // {
    //     refresh_online_state(dev_ptr);
    // }
    return 0;
}


static int bmser_bin_lib_channel_check_write(channel_t *channel, int time_out)
{
    uint8_t write_buff[4] = {0};
    int ret = sem_wait_time(&channel->write_sem, 0, time_out);
    if (ret == 0)
    {
        write_list_t *w_info = get_data_from_list(channel);
        if (w_info == NULL)
        {
            proto_syslog(LOG_ERR, "get data from list is null");
            return -1;
        }
        tag_t *write_tag = &w_info->w_tag;
        if (write_tag == NULL)
        {
            proto_syslog(LOG_ERR, "get data from list is null");
            free(w_info);
            return -2;
        }
        device_t *dev = (device_t *)write_tag->dev_ptr;
        bmser_dev_attr *attr = (bmser_dev_attr *)dev->proto_ptr;
        bmser_parser_info_t *parse_info = (bmser_parser_info_t *)write_tag->proto_ptr;
        proto_syslog(LOG_INFO, "get data from list is [%s.%s(type:%d, double=%f, long=%ld)]",
                     dev->no, write_tag->name, write_tag->data_type, write_tag->write_cache.to_float, write_tag->write_cache.to_int);
        ret = tag_to_bmser_buff(write_tag, write_buff, sizeof(write_buff));
        if (ret <= 0)
        {
            proto_syslog(LOG_ERR, "tag to bmser buff error,tag:%s", write_tag->name);
            free(w_info);
            return -3;
        }

        ret = bmser_write(write_buff, attr->addr.no, parse_info->sub_id, ret >> 1);
        if (0 != ret)
        {
            w_info->pw_tag->last_state_code = w_info->pw_tag->state_code;
            w_info->pw_tag->state_code = TAG_STATE_BAD;
            // write_tag->write_cache.to_float = write_tag->read_cache.to_float; // 写入错误，将写cache写回去，方便上位机查看
            //  写入错误没有其他处理,后期把质量码加上,方便知道数据是否可用
            proto_syslog(LOG_ERR, "bmser write tag error, tag:%s, ret = %d", w_info->pw_tag->name, ret);
        }
        else
        {
            w_info->pw_tag->last_state_code = w_info->pw_tag->state_code;
            w_info->pw_tag->state_code = TAG_STATE_GOOD;
            sync_wtag_read_cache(w_info);
        }
        free(w_info);
    }

    return 0;
}
void refresh_online_state(device_t *dev_ptr);
void set_dev_offline(device_t *dev_ptr) ;
// 检查BMSER设备在线状态
int  check_bmser_online_state(device_t* dev)
{
    static uint8_t  bmser_rx_buff[4]  = {0};
    static uint8_t  bmser_rx_valid[4] = {0};
    bmser_dev_attr* attr              = (bmser_dev_attr*)dev->proto_ptr;
    // 读取 0x0137 和 0x0138
    int       tmp  = bmser_read(bmser_rx_buff, bmser_rx_valid, sizeof(bmser_rx_buff), attr->addr.no, BMSER_ONLINE_ADDR, BMSER_ONLINE_POINT_COUNT);
    long long time = get_time_stamp_ms();
    if (tmp == 0)
    {
        uint32_t v = GET_FOR_BADC(bmser_rx_buff);
        if (v == attr->online) // 如果在线状态没有变化，则超时判定离线
        {
            if (time - attr->last_time > BMSER_CHECK_ONLINE_TIME)
            {
                set_dev_offline(dev);
            }
        }
        else
        {
            attr->online    = v;
            attr->last_time = time;
            refresh_online_state(dev);
        }
    }
    else
    {
        if (time - attr->last_time > BMSER_CHECK_ONLINE_TIME)
        {
            set_dev_offline(dev);
        }
    }
    return 0;
}

static int bmser_bin_lib_channel(channel_t *channel)
{
    device_t *dev = NULL;
    bmser_poll_info_t *poll = NULL;
    long long cur_ms = 0;
    bmser_dev_attr *attr = NULL;
    int err_counter = 0;
    static uint8_t bmser_rx_buff[4096] = {0};
    static uint8_t bmser_rx_valid[4096] = {0};
    int tmp = bmser_init();
    if (tmp != 0)
    {
        proto_syslog(LOG_ERR, "bmser_init error, ret=%d", tmp);
        return -1;
    }

    proto_syslog(LOG_ERR, "bmser_init success, bmser lib protocol start");
    while (1)
    {
        check_tag_state(channel);
        int contine_interv = 0;
        bmser_bin_lib_channel_check_write(channel, 20);
        do
        {
            contine_interv = 0;

            cur_ms = clock_get_ms();
            for (size_t i = 0; i < channel->chunks_interv_size; i++)
            {
                bmser_bin_lib_channel_check_write(channel, 0);
                if ((cur_ms - channel->chunks_interv[i]->last_poll) >= channel->chunks_interv[i]->poll_interv)
                {
                    channel->chunks_interv[i]->last_poll = cur_ms;
                    dev = channel->chunks_interv[i]->dev_ptr;
                    attr = (bmser_dev_attr *)dev->proto_ptr;
                    poll = (bmser_poll_info_t *)channel->chunks_interv[i]->chunk_proto_ptr;
                    proto_syslog(LOG_DEBUG, "cur_ms=%llu, dev_addr:0x%08x, start_sub_id:%d, sub_sum:%d", cur_ms, attr->addr.no, poll->sub_id, poll->sub_num);
                    tmp = bmser_read(bmser_rx_buff, bmser_rx_valid, sizeof(bmser_rx_buff), attr->addr.no, poll->sub_id, poll->sub_num);
                    if (tmp != 0)
                    {
                        proto_syslog(LOG_ERR, "bmser_read error,ret=%d, dev_addr:0x%08x, sub_id:%d", tmp, attr->addr.no, poll->sub_id);
                        if (++err_counter > COMMUNIT_RETRY_MAX)
                        {
                            proto_syslog(LOG_ERR, "bmser bin lib channel return");
                            return -2;
                        }
                    }
                    else
                    {
                        err_counter = 0;
                        tmp = poll->sub_num * 2;
                        proto_syslog_hex(LOG_DEBUG, bmser_rx_buff, tmp, "bmser_read:%d", tmp);
                        bmser_bin_extra(attr->addr.no, poll, bmser_rx_buff, sizeof(bmser_rx_buff), bmser_rx_valid);
                        contine_interv = 1;
                    }
                    check_bmser_online_state(dev);
                }
            }
        } while (contine_interv != 0);
    }
    return -1;
}
#else
static int bmser_bin_lib_channel(channel_t *channel)
{
    while(1)
    {
         proto_syslog(LOG_ERR, "channel %s loop", channel->channel);
         sleep(3);
    }
    return -1;
}
#endif
void set_priority()
{
    int policy,ret;  
    struct sched_param param1; 

    ret = pthread_getschedparam(pthread_self(), &policy, &param1); 
    if(ret<0)
        printf ("getschedparam:%s",strerror(errno));
    // if (policy == SCHED_FIFO){  
    //     printf("policy:SCHED_FIFO\n");  
    // }  
    // else if (policy == SCHED_OTHER){  
    //     printf("policy:SCHED_OTHER\n");  
    // }  
    // else if (policy == SCHED_RR){  
    //     printf("policy:SCHED_RR\n");  
    // } 
    // printf("thread_func priority is %d\n", param1.sched_priority);
    policy = SCHED_RR;
    param1.sched_priority = 51;
    ret = pthread_setschedparam(pthread_self(), policy, &param1);
    if(ret<0)
        printf("setschedparam:%s",strerror(errno));
 
}
#if 0
void check_write_sleep(channel_t *channel)
{
    proto_syslog(LOG_ERR, "sem_wait_time");
    int ret = sem_wait_time(&channel->write_sem, 3, 0);
    if (ret == 0)
    {
        proto_syslog(LOG_ERR, "get sem");
        write_list_t *w_info = get_data_from_list(channel);
        if (w_info == NULL)
        {
            proto_syslog(LOG_ERR, "get data from list is null");
            return;
        }
        tag_t *write_tag = &w_info->w_tag;
        if (write_tag == NULL)
        {
            proto_syslog(LOG_ERR, "get data from list is null");
        }
        else
        {
            write_tag->last_state_code = write_tag->state_code;
            write_tag->state_code = TAG_STATE_BAD;
            reflush_data_up(write_tag);
        }
        free(w_info);
    }
}
#endif

int read_data_timeout(int fd, unsigned char *buff,int buffsize, int timeout)
{
    int s_rc;
    int rx_len = 0;
    struct timeval tv;
    tv.tv_sec = timeout / 1000;
    tv.tv_usec = (timeout % 1000) * 1000;
    fd_set rset;
    FD_ZERO(&rset);
    FD_SET(fd, &rset);
    int len_tmp = 0;
    do
    {
        while ((s_rc = select(fd + 1, &rset, NULL, NULL, &tv)) == -1)
        {
            if (errno == EINTR)
            {
                proto_syslog(LOG_ERR, "A non blocked signal was caught\n");
                FD_ZERO(&rset);
                FD_SET(fd, &rset);
            }
            else
            {
                return -1;
            }
        }

        if (s_rc == 0)
        {
            errno = ETIMEDOUT;
            break;
        }
        len_tmp = read(fd, buff + rx_len,buffsize);
        if (len_tmp == 0)
        {
            break;
        }
        rx_len += len_tmp;
        tv.tv_sec = 0;
        tv.tv_usec = 20000;
    } while (true);
    ems_syslog_hex(LOG_INFO,buff,rx_len,"BIN RX:");
    return rx_len;
}

//返回 -1 未找到设备 -2 读写超时 大于0 读取到的数据长度
int send_data_and_get_return_by_dev_no(char *dev_no,int timeout,unsigned char *send_buf,int send_len,unsigned char *recv_buf,int recv_buf_size)
{
    int ret = -1;
    proto_forward_t *pro_pr = get_proto_forward_var();
    channel_t *chan_pr=NULL;
    device_t *dev_pr=NULL;
    if(find_dev_channel_pr_by_dev_no(pro_pr,dev_no,&chan_pr,&dev_pr)==0)
    {
        channel_lock(chan_pr);
        if(chan_pr->s>0)
        {
            ems_syslog_hex(LOG_INFO,send_buf,send_len,"BIN TX:");
            if (write(chan_pr->s, send_buf, send_len) < 0)
            {
                proto_syslog(LOG_ERR, "dev_no[%s] send data error!", dev_no);
                ret = -2;
            }
            else
            {
                int len = read_data_timeout(chan_pr->s,recv_buf,recv_buf_size,timeout);
                if(len <= 0)
                    ret = -2;
                else
                    ret = len;
            }
        }
        else
            ret = -2;
        channel_unlock(chan_pr);
    }
    
    return ret;

}

typedef struct can_attr
{
    int bitrate;
    int filter_len;
    struct can_filter filter[0];
} can_attr;

static can_attr *get_can_attr(const char *dev)
{
    can_attr *attr = NULL;
    char *cfg_f = DEV_PORT_CFG_FILE;
    char *f_data = read_file_data(cfg_f);
    if (f_data == NULL)
    {
        proto_syslog(LOG_ERR, "load uart config file error, file:%s", cfg_f);
        return attr;
    }

    cJSON *root = json_parse_string_with_comments(f_data);
    if (!root)
    {
        proto_syslog(LOG_ERR, "parse file to json obj error, file:%s", cfg_f);
        goto get_can_attr_exit;
    }

    cJSON *tmp_json = cJSON_GetObjectItemCaseSensitive(root, "dev_cfg");
    if (!tmp_json)
    {
        proto_syslog(LOG_ERR, "invalid dev_cfg");
        goto get_can_attr_exit;
    }

    cJSON *dev_json = cJSON_GetObjectItemCaseSensitive(tmp_json, "can");
    if (!dev_json)
    {
        proto_syslog(LOG_ERR, "invalid can");
        goto get_can_attr_exit;
    }

    cJSON *attr_arr_json = cJSON_GetObjectItemCaseSensitive(dev_json, "attr");
    if (!attr_arr_json)
    {
        proto_syslog(LOG_ERR, "invalid attr");
        goto get_can_attr_exit;
    }
    int cnt = cJSON_GetArraySize(attr_arr_json);
    cJSON *target_json = NULL;
    for (int i = 0; i < cnt; i++)
    {
        tmp_json = cJSON_GetArrayItem(attr_arr_json, i);
        target_json = cJSON_GetObjectItemCaseSensitive(tmp_json, "name");
        if (strcmp(target_json->valuestring, dev) == 0)
        {
            target_json = tmp_json;
            break;
        }
        target_json = NULL;
    }

    if (target_json == NULL)
    {
        proto_syslog(LOG_ERR, "no available config was found, %s", dev);
        goto get_can_attr_exit;
    }

    attr_arr_json = cJSON_GetObjectItemCaseSensitive(target_json, "filter");
    if (attr_arr_json == NULL)
    {
        goto get_can_attr_exit;
    }
    cnt = cJSON_GetArraySize(attr_arr_json);
    attr = calloc(1, sizeof(can_attr) + sizeof(struct can_filter) * cnt);
    attr->filter_len = cnt;

    tmp_json = cJSON_GetObjectItemCaseSensitive(target_json, "baud");
    if (tmp_json == NULL)
    {
        goto get_can_attr_exit;
    }
    attr->bitrate = tmp_json->valueint;

    for (int i = 0; i < cnt; i++)
    {
        tmp_json = cJSON_GetArrayItem(attr_arr_json, i);
        target_json = cJSON_GetObjectItemCaseSensitive(tmp_json, "id");
        dev_json = cJSON_GetObjectItemCaseSensitive(tmp_json, "mask");
        if ((target_json == NULL) || (dev_json == NULL))
        {
            free(attr);
            attr = NULL;
            break;
        }
        attr->filter[i].can_id = analysis_addr(target_json->valuestring);
        if (strlen(target_json->valuestring) > 5)
        {
            attr->filter[i].can_id |= 0x80000000;
        }
        attr->filter[i].can_mask = analysis_addr(dev_json->valuestring) & 0x9fffffff;// 不关心错误和远程帧标志位
    }

    proto_syslog(LOG_NOTICE, "%s bitrate=%d, filter_len=%d", dev, attr->bitrate, attr->filter_len);
    for (int i = 0; i < attr->filter_len; i++)
    {
        proto_syslog(LOG_NOTICE, "%s filter[%d]:{id:0x08%x, mask:0x%08x}", dev, i, attr->filter[i].can_id, attr->filter[i].can_mask);
    }

get_can_attr_exit:
    free(f_data);
    if (root != NULL)
    {
        cJSON_Delete(root);
    }
    return attr;
}

static int set_can_attr(int fd, const can_attr *p_attr)
{
    if (setsockopt(fd, SOL_CAN_RAW, CAN_RAW_FILTER, p_attr->filter, sizeof(struct can_filter) * p_attr->filter_len) != 0)
    {
        return -1;
    }

    return 0;
}

static int set_can_bitrate(const char *dev, int bitrate)
{
    char cmd[256] = {0};
    snprintf(cmd, sizeof(cmd), "ip link set %s down;ip link set %s up type can bitrate %d", dev, dev, bitrate);
    cmd_call(cmd, 5, NULL, 0);
    proto_syslog(LOG_NOTICE, "%s:%s", dev, cmd);
    return 0;
}

static int can_dev_map(char *hw_dev, int max_len)
{
    struct can_map
    {
        const char *sw_dev; // 软件can设备号
        const char *hw_dev; // 硬件丝印
    } can_map[] = {
        {.hw_dev = "can1", .sw_dev = "lnxcan1"},
        {.hw_dev = "can2", .sw_dev = "lnxcan2"},
        {.hw_dev = "can3", .sw_dev = "lnxcan3"},
        {.hw_dev = "can4", .sw_dev = "lnxcan4"},
        {.hw_dev = "can5", .sw_dev = "lnxcan5"},
    };
    for (int i = 0; i < sizeof(can_map) / sizeof(can_map[0]); i++)
    {
        if (strcmp(hw_dev, can_map[i].hw_dev) == 0)
        {
            snprintf(hw_dev, max_len, "%s", can_map[i].sw_dev);
            proto_syslog(LOG_WARNING, "hw_can[%s]->sw_can[%s]", can_map[i].hw_dev, can_map[i].sw_dev);
            return 0;
        }
    }

    proto_syslog(LOG_WARNING, "hw_can[%s]->sw_can[%s]", hw_dev, hw_dev);
    return -1;
}

static int new_can_socket(const char *dev_f)
{
    int fd = 0;
    int ret = 0;
    struct sockaddr_can addr = {0};
    struct ifreq ifr = {0};
    char dev[64] = {0};
    snprintf(dev, sizeof(dev), "%s", dev_f);
    str_to_lower_case(dev);
    can_dev_map(dev, sizeof(dev));// 根据丝印映射一下
    can_attr *attr = get_can_attr(dev_f);
    if (attr == NULL)
    {
        proto_syslog(LOG_ERR, "get can attr error:%s", dev_f);
        return -99;
    }
    else
    {
        set_can_bitrate(dev, attr->bitrate);
    }

    fd = socket(PF_CAN, SOCK_RAW | SOCK_NONBLOCK, CAN_RAW);
    if (fd < 0)
    {
        proto_syslog(LOG_ERR, "socket can failed:%s", dev);
        free(attr);
        return -1;
    }

    strcpy(ifr.ifr_name, dev);
    ret = ioctl(fd, SIOCGIFINDEX, &ifr);
    if (ret < 0)
    {
        proto_syslog(LOG_ERR, "%s ioctl SIOCGIFINDEX failed, ret=%d", dev, ret);
        goto new_can_if_err;
    }
    addr.can_family = PF_CAN;
    addr.can_ifindex = ifr.ifr_ifindex;
    ret = bind(fd, (struct sockaddr *)&addr, sizeof(addr));
    if (ret < 0)
    {
        proto_syslog(LOG_ERR, "%s bind failed, ret = %d", dev, ret);
        goto new_can_if_err;
    }

    set_can_attr(fd, attr);
    free(attr);

    return fd;
new_can_if_err:
    close(fd);
    return ret;
}

static int extract_can_data(proto_can_poll *poll, long long cur_ms,channel_t *channel)
{
    int tag_sum = poll->tag_sum;
    tag_t *p_tag = NULL;
    can_parser_info_t *parser_info = NULL;
    uint8_t *buff = poll->buff.buff;
    tag_value_t tag_value = {0};
    if (tag_sum > 0)
    {
        refresh_online_state((device_t *)(poll->p_tag[0]->dev_ptr));
    }
    for (int i = 0; i < tag_sum; i++)
    {
        p_tag = poll->p_tag[i];
        parser_info = (can_parser_info_t *)(p_tag->proto_ptr);
        // proto_syslog(LOG_DEBUG, "p_tag->name:%s, value_64:0x%016"PRIx64", data_mask:0x%016"PRIx64", data_target:0x%016"PRIx64"", p_tag->name, poll->buff.value_64, parser_info->data_mask, parser_info->data_target);
        if ((poll->buff.value_64 & parser_info->data_mask) != (parser_info->data_target & parser_info->data_mask))
        {
            continue;
        }

        switch (parser_info->data_type)
        {
        case TYPE_BOOL:
            tag_value.type = 0;
            if (parser_info->reverse != 0)
            {
                tag_value.value.to_int = (*(buff + parser_info->byte_index) != 0) ? 0 : 1;
            }
            else
            {
                tag_value.value.to_int = (*(buff + parser_info->byte_index) != 0) ? 1 : 0;
            }
            proto_set_tag(p_tag, tag_value, cur_ms);
            break;

        case TYPE_INT:
            GET_INT(tag_value, buff + parser_info->byte_index, parser_info->data_order);
            proto_set_tag(p_tag, tag_value, cur_ms);
            break;

        case TYPE_UINT:
            GET_UINT(tag_value, buff + parser_info->byte_index, parser_info->data_order);
            proto_set_tag(p_tag, tag_value, cur_ms);
            break;

        case TYPE_FLOAT:
            GET_FLOAT(tag_value, buff + parser_info->byte_index, parser_info->data_order);
            proto_set_tag(p_tag, tag_value, cur_ms);
            break;

        case TYPE_BITS:
            GET_BITS(tag_value, buff + parser_info->byte_index, parser_info->data_order);
            if (parser_info->reverse != 0)
            {
                tag_value.value.to_int = GET_BITS_IN_U16(~tag_value.value.to_int, parser_info->bit_pos, parser_info->bit_len);
            }
            else
            {
                tag_value.value.to_int = GET_BITS_IN_U16(tag_value.value.to_int, parser_info->bit_pos, parser_info->bit_len);
            }
            proto_set_tag(p_tag, tag_value, cur_ms);
            break;

        default:
            proto_syslog(LOG_ERR, "extract data error, date type isn't supported, data_type = %d, tag=%s, id:%d, index:%d",
                         parser_info->data_type, p_tag->name, parser_info->id, parser_info->byte_index);
            BUSINESS_LOG(BLOG_ERR, DAS_ID_CAN_4, p_tag->dev_tag_name, "[数采]: id: %d, tag: %s, %d 不支持的协议数据类型, 请检查模板", parser_info->id, p_tag->name, parser_info->data_type);
            break;
        }
    }
    abnormal_log_match_record(poll->p_tag, poll->tag_sum, poll->buff.buff, poll->dlc, (char *)poll->buff.buff, poll->dlc);

    return 0;
}

static int can_send(int can_fd, struct can_frame *frame)
{
    int len = write(can_fd, frame, sizeof(struct can_frame));
    proto_syslog_hex(LOG_DEBUG, frame->data, frame->can_dlc, "can tx:id:0x%08x dlc=%d, ret=%d", frame->can_id, frame->can_dlc, len);
    if (len != sizeof(struct can_frame))
    {
        len = -1;
    }
    return len;
}

static int can_tag_write_dev(int can_fd, tag_t *write_tag)
{
    int use_len = 0;
    channel_t *channel = (channel_t *)write_tag->chan_ptr;
    can_parser_info_t *parse_info = (can_parser_info_t *)write_tag->proto_ptr;
    proto_can_poll *poll = NULL;
    uint8_t *dest_buff = NULL;
    struct can_frame frame = {0};
    if (write_tag->poll_list_type == 0)
    {
        poll = (proto_can_poll *)channel->chunks_interv[write_tag->poll_list_num]->chunk_proto_ptr;
    }
    else
    {
        poll = (proto_can_poll *)channel->chunks_nonstop[write_tag->poll_list_num]->chunk_proto_ptr;
    }
    //
    dest_buff = &poll->buff.buff[parse_info->byte_index];
    switch (parse_info->data_type)
    {
    case TYPE_BOOL:
        *dest_buff = (write_tag->write_cache.to_int != 0) ? 1 : 0;
        use_len = 1;
        break;

    case TYPE_INT:
        SET_INT(write_tag, dest_buff, use_len, parse_info->data_order);
        break;

    case TYPE_UINT:
        SET_UINT(write_tag, dest_buff, use_len, parse_info->data_order);
        break;

    case TYPE_FLOAT:
        SET_FLOAT(write_tag, dest_buff, use_len, parse_info->data_order);
        break;

    default:
        proto_syslog(LOG_ERR, "write reg type error, [id:%d, index:0x%x, type=%d]", parse_info->id, parse_info->byte_index, parse_info->data_type);
        BUSINESS_LOG(BLOG_CRIT, DAS_ID_CAN_3, write_tag->dev_tag_name, "[数采]: id: %d, %d 不支持的数据类型, 请检查模板配置",  parse_info->id, parse_info->data_type);
        break;
    }

    if (use_len > 0)
    {
        if (parse_info->send_id == CAN_ERR_ID)
        {
            frame.can_id = poll->id;
            frame.can_dlc = poll->dlc;
            memcpy(frame.data, poll->buff.buff, poll->dlc);
            use_len = can_send(can_fd, &frame);
        }
        else
        {
            frame.can_id = parse_info->send_id;
            frame.can_dlc = use_len;
            memcpy(frame.data, dest_buff, use_len);
            use_len = can_send(can_fd, &frame);
        }
    }

    return use_len;
}

static int can_channel_check_write(int can_fd, channel_t *channel, int time_out)
{

    int ret = sem_wait_time(&channel->write_sem, 0, time_out);
    if (ret == 0)
    {
        write_list_t *w_info = get_data_from_list(channel);
        if (w_info == NULL)
        {
            proto_syslog(LOG_ERR, "get data from list is null");
            return -1;
        }
        tag_t *write_tag = &w_info->w_tag;
        if (write_tag == NULL)
        {
            free(w_info);
            proto_syslog(LOG_ERR, "get data from list is null");
            return -2;
        }
        device_t *dev = (device_t *)write_tag->dev_ptr;
        proto_syslog(LOG_INFO, "get data from list is [%s.%s(type:%d, double=%f, long=%ld)]",
                     dev->no, write_tag->name, write_tag->data_type, write_tag->write_cache.to_float, write_tag->write_cache.to_int);

        ret = can_tag_write_dev(can_fd, write_tag);
        if (0 > ret)
        {
            w_info->pw_tag->last_state_code = w_info->pw_tag->state_code;
            w_info->pw_tag->state_code = TAG_STATE_BAD;
            proto_syslog(LOG_ERR, "can write tag error, tag:%s, ret = %d", w_info->pw_tag->name, ret);
            BUSINESS_LOG(BLOG_ERR, DAS_ID_CAN_2, write_tag->name, "[数采]: %s 通道写tag: %s 失败", channel->channel, w_info->pw_tag->name);
        }
        else
        {
            BUSINESS_LOG(BLOG_NOTICE, DAS_ID_CAN_2, write_tag->name, "[数采]: %s 通道写tag: %s 成功", channel->channel, w_info->pw_tag->name);
            w_info->pw_tag->last_state_code = w_info->pw_tag->state_code;
            w_info->pw_tag->state_code = TAG_STATE_GOOD;
            sync_wtag_read_cache(w_info);
        }
        free(w_info);
    }

    return 0;
}

void can_set_raw_send_callback(proto_can_poll *poll, void (*cb)(proto_can_poll *poll, const uint8_t *data, int data_length, void *channel), void *channel)
{
    if (poll == NULL) {
        return;
    }
    
    poll->raw_send_cb = cb;
    poll->raw_send_channel = channel;
}

/* Set raw receive callback function */
void can_set_raw_receive_callback(proto_can_poll *poll, void (*cb)(proto_can_poll *poll, const uint8_t *data, int data_length, void *channel), void *channel)
{
    if (poll == NULL) {
        return;
    }
    
    poll->raw_receive_cb = cb;
    poll->raw_receive_channel = channel;
}

/* Raw data callback functions */
static void __attribute__((unused)) can_raw_send_callback(proto_can_poll *poll, const uint8_t *data, int data_length, void *channel)
{
    channel_t *ch = channel;
    
    /* Fill send_buf in channel structure */
    if (ch->send_buf != NULL) {
        free(ch->send_buf);
    }
    ch->send_buf = malloc(data_length);
    if (ch->send_buf != NULL) {
        memcpy(ch->send_buf, data, data_length);
        ch->send_buf_size = data_length;
    }
    
}

static void __attribute__((unused)) can_raw_receive_callback(proto_can_poll *poll, const uint8_t *data, int data_length, void *channel)
{
    channel_t *ch = channel;
    
    /* Fill recv_buf in channel structure */
    if (ch->recv_buf != NULL) {
        free(ch->recv_buf);
    }
    ch->recv_buf = malloc(data_length);
    if (ch->recv_buf != NULL) {
        memcpy(ch->recv_buf, data, data_length);
        ch->recv_buf_size = data_length;
    }
}

static int can_channel(channel_t *channel)
{
    uint32_t can_id = 0;
    int ret = -1;
    long long cur_ms = 0;
    long long h_ms = 0;
    struct can_frame frdup = {0};
    struct can_frame frame = {0};
    struct timeval tv = {0};
    proto_can_poll *poll = NULL;
    fd_set rset = {0};
    tv.tv_sec = 0;
    tv.tv_usec = 20 * 1000;
    int fd = new_can_socket(((can_proto_s *)(channel->proto_ptr))->dev);
    if (fd < 0)
    {
        proto_syslog(LOG_CRIT, "[数采]:new_can_socket err");
        BUSINESS_LOG(BLOG_ERR, DAS_ID_CAN_1, channel->channel, "[数采]: %s 通道创建失败", channel->channel);
        return -2;
    }
    else
    {
        BUSINESS_LOG(BLOG_NOTICE, DAS_ID_CAN_1, channel->channel, "[数采]: %s 通道创建成功", channel->channel);
    }
    // can_set_raw_send_callback(poll, can_raw_send_callback,(void *)channel);
    // can_set_raw_receive_callback(poll, can_raw_receive_callback,(void *)channel);
    channel->s = fd;

    while (1)
    {
        cur_ms = clock_get_ms();
        if ((cur_ms - h_ms) > 100)
        {
            h_ms = cur_ms;
            check_tag_state(channel);
            for (int i = 0; i < channel->devs_size; i++)
            {
                check_online_state(channel->devs[i]);
            }
        }
        can_channel_check_write(fd, channel, 0); // 检查UA写
        // 轮询写
        for (int i = 0; i < channel->chunks_interv_size; i++)
        {
            poll = (proto_can_poll *)channel->chunks_interv[i]->chunk_proto_ptr;
            if ((poll->loop_tx != 0) && ((cur_ms - channel->chunks_interv[i]->last_poll) > channel->chunks_interv[i]->poll_interv))
            {
                channel->chunks_interv[i]->last_poll = cur_ms;
                frame.can_id = poll->id;
                if (poll->remote_len < 0)
                {
                    frame.can_dlc = poll->dlc;
                    memcpy(frame.data, poll->buff.buff, poll->dlc);
                }
                else
                {
                    frame.can_dlc = poll->remote_len;
                    frame.can_id |= CAN_RTR_FLAG;
                }
                channel_lock(channel);
                if (can_send(fd, &frame) < 0)
                {
                    proto_syslog(LOG_ERR, "can send error,id:0x%08x", frame.can_id);
                }
                channel_unlock(channel);
            }
        }

        FD_ZERO(&rset);
        FD_SET(fd, &rset);
        tv.tv_sec = 0;
        tv.tv_usec = 20 * 1000;
        while ((ret = select(fd + 1, &rset, NULL, NULL, &tv)) == -1)
        {
            if (errno == EINTR)
            {
                proto_syslog(LOG_WARNING, "A non blocked signal was caught");
                FD_ZERO(&rset);
                FD_SET(fd, &rset);
            }
            else
            {
                break;
            }
        }
        if (ret == 0)
        {
            continue;
        }

        channel_lock(channel);
        ret = read(fd, &frdup, sizeof(frdup));
        channel_unlock(channel);
        if (ret < sizeof(frdup))
        {
            proto_syslog(LOG_ERR, "read failed\r\n");
            continue;
        }
        if (frdup.can_id & CAN_ERR_FLAG)
        {
            proto_syslog(LOG_ERR, "CAN device error\r\n");
            continue;
        }
        
        can_id = frdup.can_id;
        proto_syslog_hex(LOG_DEBUG, frdup.data, frdup.can_dlc, "can rx:id:0x%08x dlc=%d", can_id, frdup.can_dlc);

        cur_ms = clock_get_ms();
        for (int i = 0; i < channel->chunks_nonstop_size; i++)
        {
            poll = (proto_can_poll *)channel->chunks_nonstop[i]->chunk_proto_ptr;
            if ((poll->id ^ can_id) == 0)
            {
                memcpy(poll->buff.buff, frdup.data, frdup.can_dlc);
                // if (frdup.can_dlc > 0 && poll->raw_receive_cb != NULL) {
                //     proto_syslog(LOG_ERR,"poll->raw_receive_cb != NULL frdup.can_dlc = %d",frdup.can_dlc);
                //     poll->raw_receive_cb(poll, frdup.data, frdup.can_dlc, poll->raw_receive_channel);
                // }
                if (poll->fun != NULL) // 有自定义解析
                {
                    poll->fun(poll);
                    refresh_online_state(channel->chunks_nonstop[i]->dev_ptr);
                }
                else
                {
                    extract_can_data(poll, cur_ms,channel);
                }
            }
        }
        for (int i = 0; i < channel->chunks_interv_size; i++)
        {
            poll = (proto_can_poll *)channel->chunks_interv[i]->chunk_proto_ptr;
            if ((poll->id ^ can_id) == 0)
            {
                memcpy(poll->buff.buff, frdup.data, frdup.can_dlc);
                if (poll->fun != NULL) // 有自定义解析
                {
                    poll->fun(poll);
                    refresh_online_state(channel->chunks_interv[i]->dev_ptr);
                }
                else
                {
                    extract_can_data(poll, cur_ms,channel);
                }
            }
        }
    }

    close(fd);
    return -1;
}


static int extract_custom_data(proto_custom_poll *poll, uint8_t *buff, int rx_len, long long cur_ms)
{
    int tag_sum = poll->tag_sum;
    tag_t *p_tag = NULL;
    custom_parser_info_t *parser_info = NULL;

    tag_value_t tag_value = {0};
    if (tag_sum > 0)
    {
        refresh_online_state((device_t *)(poll->p_tag[0]->dev_ptr));
    }
	
	unsigned short crc_calculated = 0;
	int data_offset = 0;
	
    switch (poll->crc16){
		case CRC16_N:
	        break;
	
		case CRC16_X:
			break;

		case CRC16_C:            
			break;

		case CRC16_M:
	        data_offset = 0;
			crc_calculated = modbus_data_check_crc(buff + data_offset, rx_len - data_offset);
		    break;

		case CRC16_IBM:
			break;

		case CRC16_USB:
			break;
			
		case CRC16_DNP:
			break;

		case CRC16_DEC:
			break;

		default:
		    proto_syslog(LOG_ERR, "compute crc16 data error, crc16 type isn't supported, crc16_type = %d", poll->crc16);
            BUSINESS_LOG(BLOG_ERR, DAS_ID_CUSTOM_5, NULL, "[数采]: %d 不支持的检验类型, 请检查模板", poll->crc16);
            crc_calculated = 123;
	        break;
	}

	if (crc_calculated)
    {
        proto_syslog(LOG_ERR, "modbus crc error 0x%x", crc_calculated);
        return -1;
    }
	
    for (int i = 0; i < tag_sum; i++)
    {
        p_tag = poll->p_tag[i];
        parser_info = (custom_parser_info_t *)(p_tag->proto_ptr);
        switch (parser_info->data_type)
        {
        case TYPE_BOOL:
            tag_value.type = 0;
            if (parser_info->reverse != 0)
            {
                tag_value.value.to_int = (*(buff + parser_info->byte_index) != 0) ? 0 : 1;
            }
            else
            {
                tag_value.value.to_int = (*(buff + parser_info->byte_index) != 0) ? 1 : 0;
            }
            proto_set_tag(p_tag, tag_value, cur_ms);
            break;

        case TYPE_INT:
            GET_INT(tag_value, buff + parser_info->byte_index, parser_info->data_order);
            proto_set_tag(p_tag, tag_value, cur_ms);
            break;

        case TYPE_UINT:
            GET_UINT(tag_value, buff + parser_info->byte_index, parser_info->data_order);
            proto_set_tag(p_tag, tag_value, cur_ms);
            break;

        case TYPE_FLOAT:
            GET_FLOAT(tag_value, buff + parser_info->byte_index, parser_info->data_order);
            proto_set_tag(p_tag, tag_value, cur_ms);
            break;

        case TYPE_BITS:
            GET_BITS(tag_value, buff + parser_info->byte_index, parser_info->data_order);
            if (parser_info->reverse != 0)
            {
                tag_value.value.to_int = GET_BITS_IN_U16(~tag_value.value.to_int, parser_info->bit_pos, parser_info->bit_len);
            }
            else
            {
                tag_value.value.to_int = GET_BITS_IN_U16(tag_value.value.to_int, parser_info->bit_pos, parser_info->bit_len);
            }
            proto_set_tag(p_tag, tag_value, cur_ms);
            break;

        default:
            proto_syslog(LOG_ERR, "extract data error, date type isn't supported, data_type = %d, tag=%s, id:%d, index:%d",
                         parser_info->data_type, p_tag->name, parser_info->id, parser_info->byte_index);
            BUSINESS_LOG(BLOG_ERR, DAS_ID_CUSTOM_6, NULL, "[数采]: id: %d, tag: %s, %d 不支持的协议数据类型, 请检查模板", parser_info->id, p_tag->name, parser_info->data_type);
            break;
        }
    }

    return 0;
}

int custom_proto_read(custom_hd_t *hd, proto_custom_poll *poll, void *rx_buff, int rx_buff_len)
{
	int use_len = 0;
	uint8_t *dest_buff = NULL;

	for (int i = 0; i < poll->tag_sum; i++){ 
		tag_t *write_tag = poll->p_tag[i];
		custom_parser_info_t *parse_info = (custom_parser_info_t *)write_tag->proto_ptr;
		
	    dest_buff = &poll->buff[parse_info->byte_index];
	    switch (parse_info->data_type)
	    {
	    case TYPE_BOOL:
	        *dest_buff = (write_tag->write_cache.to_int != 0) ? 1 : 0;
			use_len = 1;
	        break;

	    case TYPE_INT:
	        SET_INT(write_tag, dest_buff, use_len, parse_info->data_order);
	        break;

	    case TYPE_UINT:
	        SET_UINT(write_tag, dest_buff, use_len, parse_info->data_order);
	        break;

	    case TYPE_FLOAT:
	        SET_FLOAT(write_tag, dest_buff, use_len, parse_info->data_order);
	        break;

	    default:
	        proto_syslog(LOG_ERR, "write reg type error, [id:%d, index:0x%x, type=%d]", parse_info->id, parse_info->byte_index, parse_info->data_type);
            BUSINESS_LOG(BLOG_ERR, DAS_ID_CUSTOM_7, write_tag->dev_tag_name, "[数采]: id: %d, %d 不支持的协议数据类型, 请检查模板", parse_info->id, parse_info->data_type);
	        break;
	    }
	}

	uint16_t val = 0; // 计算crc16
    switch (poll->crc16){
		case CRC16_N:
	        break;
	
		case CRC16_X:
			break;

		case CRC16_C:            
			break;

		case CRC16_M:
	        val = crc16_M(poll->buff, poll->len - 2);
		    poll->buff[poll->len - 1] = (val);     
            poll->buff[poll->len - 2] = (val) >> 8;
		    break;

		case CRC16_IBM:
			break;

		case CRC16_USB:
			break;
			
		case CRC16_DNP:
			break;

		case CRC16_DEC:
			break;

		default:
		    proto_syslog(LOG_ERR, "compute crc16 data error, crc16 type isn't supported, crc16_type = %d", poll->crc16);
	        break;
	}
    switch (poll->crc16){
		case CRC16_N:
	        break;
	
		case CRC16_X:
			break;

		case CRC16_C:            
			break;

		case CRC16_M:
	        val = crc16_M(poll->buff, poll->len - 2);
		    poll->buff[poll->len - 1] = (val);     
            poll->buff[poll->len - 2] = (val) >> 8;
		    break;

		case CRC16_IBM:
			break;

		case CRC16_USB:
			break;
			
		case CRC16_DNP:
			break;

		case CRC16_DEC:
			break;

		default:
		    proto_syslog(LOG_ERR, "compute crc16 data error, crc16 type isn't supported, crc16_type = %d", poll->crc16);
	        break;
	}
	
	use_len = custom_read(hd, poll->buff, poll->len, rx_buff, rx_buff_len);
    return use_len;
}

static int custom_tag_write_dev(custom_hd_t *hd, tag_t *w_tag, uint32_t *p_id, void *rx_buff, int rx_buff_len)
{
    int use_len = 0;
    channel_t *channel = (channel_t *)w_tag->chan_ptr;
    //custom_parser_info_t *parse_info = (custom_parser_info_t *)write_tag->proto_ptr;
    proto_custom_poll *poll = NULL;
    uint8_t *dest_buff = NULL;

    if (w_tag->poll_list_type == 0)
    {
        poll = (proto_custom_poll *)channel->chunks_interv[w_tag->poll_list_num]->chunk_proto_ptr;
    }
    else
    {
        poll = (proto_custom_poll *)channel->chunks_nonstop[w_tag->poll_list_num]->chunk_proto_ptr;
    }

	*p_id = poll->id | 0x80;
	for (int i = 0; i < poll->tag_sum; i++){ 
		tag_t *write_tag = poll->p_tag[i];
		custom_parser_info_t *parse_info = (custom_parser_info_t *)write_tag->proto_ptr;
		
	    dest_buff = &poll->buff[parse_info->byte_index];
	    switch (parse_info->data_type)
	    {
	    case TYPE_BOOL:
	        *dest_buff = (write_tag->write_cache.to_int != 0) ? 1 : 0;
			use_len = 1;
	        break;

	    case TYPE_INT:
	        SET_INT(write_tag, dest_buff, use_len, parse_info->data_order);
	        break;

	    case TYPE_UINT:
	        SET_UINT(write_tag, dest_buff, use_len, parse_info->data_order);
	        break;

	    case TYPE_FLOAT:
	        SET_FLOAT(write_tag, dest_buff, use_len, parse_info->data_order);
	        break;

	    default:
	        proto_syslog(LOG_ERR, "write reg type error, [id:%d, index:0x%x, type=%d]", parse_info->id, parse_info->byte_index, parse_info->data_type);
	        break;
	    }

	}

	uint16_t result = 0;
    switch (poll->crc16){
		case CRC16_N:
	        break;
	
		case CRC16_X:
			break;

		case CRC16_C:            
			break;

		case CRC16_M:
	        result = crc16_M(poll->buff, poll->len - 2);
		    poll->buff[poll->len - 1] = (result);     
            poll->buff[poll->len - 2] = (result) >> 8;
		    break;

		case CRC16_IBM:
			break;

		case CRC16_USB:
			break;
			
		case CRC16_DNP:
			break;

		case CRC16_DEC:
			break;

		default:
		    proto_syslog(LOG_ERR, "compute crc16 data error, crc16 type isn't supported, crc16_type = %d", poll->crc16);
	        break;
	}
	
	use_len = custom_write(hd, poll->buff, poll->len, rx_buff, rx_buff_len);
    return use_len;
}

static int custom_channel_check_write(custom_hd_t *hd, channel_t *channel, int time_out)
{
	long long cur_ms = 0;
	uint8_t rx_buff[512] = {0};

    int ret = sem_wait_time(&channel->write_sem, 0, time_out);
    if (ret == 0)
    {
        write_list_t *w_info = get_data_from_list(channel);
        if (w_info == NULL)
        {
            proto_syslog(LOG_ERR, "get data from list is null");
            return -1;
        }
        tag_t *write_tag = &w_info->w_tag;
        if (write_tag == NULL)
        {
            free(w_info);
            proto_syslog(LOG_ERR, "get data from list is null");
            return -2;
        }
		
        device_t *dev = (device_t *)write_tag->dev_ptr;
        proto_syslog(LOG_INFO, "get data from list is [%s.%s(type:%d, double=%f, long=%ld)]",
                     dev->no, write_tag->name, write_tag->data_type, write_tag->write_cache.to_float, write_tag->write_cache.to_int);

		uint32_t id = 0;
        ret = custom_tag_write_dev(hd, write_tag, &id, rx_buff, sizeof(rx_buff));
        if (ret < 0)
        {
            w_info->pw_tag->last_state_code = w_info->pw_tag->state_code;
            w_info->pw_tag->state_code = TAG_STATE_BAD;
            proto_syslog(LOG_ERR, "can write tag error, tag:%s, ret = %d", w_info->pw_tag->name, ret);
            BUSINESS_LOG(BLOG_ERR, DAS_ID_CUSTOM_4, write_tag->dev_tag_name, "[数采]: %s 通道写 tag: %s 失败", channel->channel, w_info->pw_tag->name);
        }
        else
        {
            BUSINESS_LOG(BLOG_NOTICE, DAS_ID_CUSTOM_4, write_tag->dev_tag_name, "[数采]: %s 通道写 tag: %s 成功", channel->channel, w_info->pw_tag->name);
			int rx_len = ret;
			cur_ms = clock_get_ms();
			for (int i = 0; i < channel->chunks_nonstop_size; i++)
			{
				proto_custom_poll *poll = (proto_custom_poll *)channel->chunks_nonstop[i]->chunk_proto_ptr;
				if (poll->id == id)
				{
					if (poll->fun != NULL) // 有自定义解析
					{
						poll->fun(poll);
						refresh_online_state(channel->chunks_nonstop[i]->dev_ptr);
					}
					else
					{
						extract_custom_data(poll, rx_buff, rx_len, cur_ms);
					}
				}
			}
			for (int i = 0; i < channel->chunks_interv_size; i++)
			{
				proto_custom_poll *poll = (proto_custom_poll *)channel->chunks_interv[i]->chunk_proto_ptr;
				if (poll->id == id)
				{
					if (poll->fun != NULL) // 有自定义解析
					{
						poll->fun(poll);
						refresh_online_state(channel->chunks_interv[i]->dev_ptr);
					}
					else
					{
						extract_custom_data(poll, rx_buff, rx_len, cur_ms);
					}
				}
			}
			refresh_online_state(dev);

            w_info->pw_tag->last_state_code = w_info->pw_tag->state_code;
            w_info->pw_tag->state_code = TAG_STATE_GOOD;
            sync_wtag_read_cache(w_info);
        }
        free(w_info);
    }

    return 0;
}

static int custom_channel(channel_t *channel)
{
	long long cur_ms = 0;
	int err_counter = 0;
	proto_custom_poll *poll_info = NULL;
	uint8_t rx_buff[512] = {0};
	device_t *dev = NULL;
	
	custom_hd_t *hd = custom_proto_init((custom_proto_t *)channel->proto_ptr);
	if (hd == NULL)
	{
		proto_syslog(LOG_CRIT, "[数采]:custom_proto_init error in the channel:[%s]", channel->channel);
        BUSINESS_LOG(BLOG_ERR, DAS_ID_CUSTOM_1, channel->channel, "[数采]: %s 通道初始化失败", channel->channel);
		return -1;
	}
	channel->s = hd->fd;
	proto_syslog(LOG_INFO, "custom_proto_init success, in the channel:[%s]", channel->channel);
    BUSINESS_LOG(BLOG_NOTICE, DAS_ID_CUSTOM_1, channel->channel, "[数采]: %s 通道初始化成功", channel->channel);
	
	while (1)
	{
	    int contine_interv = 0;
		do
        {
            contine_interv = 0;
			check_tag_state(channel);
			if(custom_channel_check_write(hd, channel, 100) < 0)
				err_counter++;
			
			cur_ms = clock_get_ms();
			for (size_t i = 0; i < channel->chunks_interv_size; i++)
			{
				if(custom_channel_check_write(hd, channel, 0) < 0)
					err_counter++;
				
				dev = channel->chunks_interv[i]->dev_ptr;
				poll_info = (proto_custom_poll *)channel->chunks_interv[i]->chunk_proto_ptr;
				if ((poll_info->loop_tx != 0) && (cur_ms - channel->chunks_interv[i]->last_poll) >= channel->chunks_interv[i]->poll_interv)
				{
					channel->chunks_interv[i]->last_poll = cur_ms;
					custom_set_time(hd, dev->timeout, dev->interval);
					channel_lock(channel);
					int ret = custom_proto_read(hd, poll_info, rx_buff, sizeof(rx_buff));
					channel_unlock(channel);
					if (ret > 0)
					{
                        BUSINESS_LOG(BLOG_NOTICE, DAS_ID_CUSTOM_2, channel->channel, "[数采]: %s 通道回复正常", channel->channel);
					    int rx_len = ret;
					    cur_ms = clock_get_ms();
						uint32_t id = poll_info->id | 0x80;
				        for (int i = 0; i < channel->chunks_interv_size; i++)
				        {
				            proto_custom_poll *poll = (proto_custom_poll *)channel->chunks_interv[i]->chunk_proto_ptr;
				            if (poll->id == id)
				            {
				                if (poll->fun != NULL) // 有自定义解析
				                {
				                    poll->fun(poll);
				                    refresh_online_state(channel->chunks_interv[i]->dev_ptr);
				                }
				                else
				                {
				                    extract_custom_data(poll, rx_buff, rx_len, cur_ms);
				                }
				            }
				        }
						refresh_online_state(dev);
						err_counter = 0;
					}
					else if (ret == 0)
					{
						proto_syslog(LOG_ERR, "%s:custom read timeout", channel->channel);
                        BUSINESS_LOG(BLOG_ERR, DAS_ID_CUSTOM_2, channel->channel, "[数采]: %s 通道通信超时", channel->channel);
						err_counter++;
						check_online_state(dev);
					}
					else
					{
						proto_syslog(LOG_ERR, "%s:custom read error(ret = %d), info:%s, errno:%d", channel->channel, ret, strerror(errno), errno);
						BUSINESS_LOG(BLOG_ERR, DAS_ID_CUSTOM_2, channel->channel, "[数采]: %s 通道回复错误, 错误原因(%s)", channel->channel, strerror(errno));
                        check_online_state(dev);
						if (errno > 0) // io错误
						{
                            BUSINESS_LOG(BLOG_ERR, DAS_ID_CUSTOM_2, channel->channel, "[数采]: %s 通道出现系统级IO错误, 请检查 物理连接或设备 是否正常", channel->channel);
							goto custom_channel_exit;
						}
						else // 协议错误
						{
                            BUSINESS_LOG(BLOG_ERR, DAS_ID_CUSTOM_2, channel->channel, "[数采]: %s 通道出现协议错误, 请检查 回复数据 是否正确", channel->channel);
							err_counter++;
						}
					}
					contine_interv = 1;
				}

				if (err_counter > COMMUNIT_RETRY_MAX)
				{
					goto custom_channel_exit;
				}
			}
	    } while ((contine_interv != 0));

		// interv空闲
		if (channel->chunks_nonstop_size > 0)
		{
			size_t nonstop_index = channel->chunks_nonstop_index;
			dev = channel->chunks_nonstop[nonstop_index]->dev_ptr;
			poll_info = (proto_custom_poll *)channel->chunks_nonstop[nonstop_index]->chunk_proto_ptr;
			if (poll_info->loop_tx != 0)
			{
				channel_lock(channel);
				int ret = custom_proto_read(hd, poll_info, rx_buff, sizeof(rx_buff)); // 读数据
				channel_unlock(channel);
				if (ret > 0)
				{
                    BUSINESS_LOG(BLOG_NOTICE, DAS_ID_CUSTOM_3, channel->channel, "[数采]: %s 通道回复正常", channel->channel);
					int rx_len = ret;
					cur_ms = clock_get_ms();
					uint32_t id = poll_info->id | 0x80;
					for (int i = 0; i < channel->chunks_nonstop_size; i++)
					{
						proto_custom_poll *poll = (proto_custom_poll *)channel->chunks_nonstop[i]->chunk_proto_ptr;
						if (poll->id == id)
						{
							if (poll->fun != NULL) // 有自定义解析
							{
								poll->fun(poll);
								refresh_online_state(channel->chunks_nonstop[i]->dev_ptr);
							}
							else
							{
								extract_custom_data(poll, rx_buff, rx_len, cur_ms);
							}
						}
					}
					refresh_online_state(dev);
					err_counter = 0;
				}
				else if (ret == 0)
				{
					proto_syslog(LOG_ERR, "%s:custom read timeout", channel->channel);
                    BUSINESS_LOG(BLOG_ERR, DAS_ID_CUSTOM_3, channel->channel, "[数采]: %s 通道通信超时", channel->channel);
					err_counter++;
					check_online_state(dev);
				}
				else
				{
					proto_syslog(LOG_ERR, "%s:custom read error(ret = %d), info:%s, errno:%d", channel->channel, ret, strerror(errno), errno);
                    BUSINESS_LOG(BLOG_ERR, DAS_ID_CUSTOM_3, channel->channel, "[数采]: %s 通道回复错误, 错误原因(%s)", channel->channel, strerror(errno));
					check_online_state(dev);
					if (errno > 0) // io错误
					{
                        BUSINESS_LOG(BLOG_ERR, DAS_ID_CUSTOM_2, channel->channel, "[数采]: %s 通道出现系统级IO错误, 请检查 物理连接或设备 是否正常", channel->channel);
						goto custom_channel_exit;
					}
					else // 协议错误
					{
                        BUSINESS_LOG(BLOG_ERR, DAS_ID_CUSTOM_2, channel->channel, "[数采]: %s 通道出现协议错误, 请检查 回复数据 是否正确", channel->channel);
						err_counter++;
					}
				}
			
				if (err_counter > COMMUNIT_RETRY_MAX)
				{
					goto custom_channel_exit;
				}
				
				nonstop_index++;
				if (nonstop_index >= channel->chunks_nonstop_size)
				{
					channel->chunks_nonstop_index = 0;
				}
				else
				{
					channel->chunks_nonstop_index = nonstop_index;
				}
			}
		}  
	}
	
custom_channel_exit:
	if (hd != NULL)
	{
		custom_proto_deinit(hd);
	}
	hd = NULL;

	return -1;
}

uint8_t calculateLCHKSUM(unsigned short lenid) {
    // 提取高4位、中间4位和低4位
    uint8_t high = (lenid >> 8) & 0x0F;   // D11-D8
    uint8_t middle = (lenid >> 4) & 0x0F; // D7-D4
    uint8_t low = lenid & 0x0F;           // D3-D0

    uint8_t sum = high + middle + low;
    uint8_t lchksum = (~sum + 1) & 0x0F;

    return lchksum;
}

unsigned short combineLength(unsigned short lenid, uint8_t lchksum) {
    // 高字节：LENID 的高4位 (D11-D8) + LCHKSUM (D3-D0)
    uint8_t highByte = ((lenid >> 8) & 0x0F) | (lchksum << 4);
    // 低字节：LENID 的低8位 (D7-D0)
    uint8_t lowByte = lenid & 0xFF;

    unsigned short length = (highByte << 8) | lowByte;
    return length;
}

static unsigned short calculateCHKSUM(const uint8_t *data, size_t length) {
    unsigned short sum = 0;
    for (size_t i = 0; i < length; i++) {
        sum += data[i];
    }

    // 对65536取模（实际上sum是unsigned short，范围已经是0-65535）
    // sum %= 65536;

    return (~sum + 1) & 0xFFFF;
}

// 创建数据帧
// LENGTH共2个字节， 由LENID和LCHKSUM组成， LENID表示INFO项的ASCII码字节数， 当LENID=0时， INFO为空，即无该项  --- infoLength 支持传入0
static uint8_t ydt1363_createFrame(uint8_t ver, uint8_t adr, uint8_t cid1, uint8_t cid2, const uint8_t *info, int infoLength, uint8_t *frame, uint32_t *frameLength) {

    unsigned short lenid = infoLength * 2;
    uint8_t lchksum = calculateLCHKSUM(lenid); //
    unsigned short length = combineLength(lenid, lchksum);

    // 计算帧总长度
    *frameLength = 1 + 12 + infoLength * 2 + 4 + 1; // SOI + 12 bytes header + INFO + CHKSUM + EOI
#if 0
    uint8_t *frame = (uint8_t *)malloc(*frameLength);
    if (!frame) {
        proto_syslog(LOG_ERR, "Failed to allocate memory for frame");
        return -1;
    }
#endif

    // 组帧
    size_t index = 0;
    frame[index++] = SOI;

    frame[index++] = toAsciiHigh(ver);
    frame[index++] = toAsciiLow(ver);

    frame[index++] = toAsciiHigh(adr);
    frame[index++] = toAsciiLow(adr);

    frame[index++] = toAsciiHigh(cid1);
    frame[index++] = toAsciiLow(cid1);

    frame[index++] = toAsciiHigh(cid2);
    frame[index++] = toAsciiLow(cid2);

    frame[index++] = toAsciiHigh(length >> 8);
    frame[index++] = toAsciiLow(length >> 8);
    frame[index++] = toAsciiHigh(length & 0xFF);
    frame[index++] = toAsciiLow(length & 0xFF);

    // 插入COMMAND INFO，注意：部分查询指令，INFO为空，即无该项，--- infoLength 支持传入0
    for (int i = 0; i < infoLength; i++) {
        frame[index++] = toAsciiHigh(info[i]);
        frame[index++] = toAsciiLow(info[i]);
    }

    unsigned short chksum = calculateCHKSUM(frame + 1, index - 1); // 计算CHKSUM
    frame[index++] = toAsciiHigh(chksum >> 8);
    frame[index++] = toAsciiLow(chksum >> 8);
    frame[index++] = toAsciiHigh(chksum & 0xFF);
    frame[index++] = toAsciiLow(chksum & 0xFF);

    frame[index++] = EOI;
    return 0;
}

// 解析数据帧
static int ydt1363_parseFrame(const uint8_t *frame, unsigned int length, uint8_t *info, unsigned short *infoLength) {
    if (length < 18 || frame[0] != SOI || frame[length - 1] != EOI) { // 数据帧 >= 18字节，SOI头，EOI尾
        proto_syslog(LOG_ERR, "Invalid frame format");
        return -1;
    }

    // 解帧
    uint8_t ver = fromAscii(frame[1], frame[2]);
    uint8_t adr = fromAscii(frame[3], frame[4]);
    uint8_t cid1 = fromAscii(frame[5], frame[6]);
    uint8_t cid2 = fromAscii(frame[7], frame[8]); 

    uint8_t rtn = cid2; // 返回码 RTN
    if (rtn != RTN_NORMAL) {
        switch (rtn) {
            case RTN_VER_ERR:
                proto_syslog(LOG_ERR, "RTN: 0x%02x (VER error, YD/T 1363.3-2014)", rtn);
                break;
            case RTN_CHKSUM_ERR:
                proto_syslog(LOG_ERR, "RTN: 0x%02x (CHKSUM error, YD/T 1363.3-2014)", rtn);
                break;
            case RTN_LCHKSUM_ERR:
                proto_syslog(LOG_ERR, "RTN: 0x%02x (LCHKSUM error, YD/T 1363.3-2014)", rtn);
                break;
            case RTN_CID2_INVALID:
                proto_syslog(LOG_ERR, "RTN: 0x%02x (CID2 invalid, YD/T 1363.3-2014)", rtn);
                break;
            case RTN_CMD_FORMAT_ERR:
                proto_syslog(LOG_ERR, "RTN: 0x%02x (cmd format error, YD/T 1363.3-2014)", rtn);
                break;
            case RTN_INVALID_DATA:
                proto_syslog(LOG_ERR, "RTN: 0x%02x (invalid data, YD/T 1363.3-2014)", rtn);
                break;
            default:
                if (rtn >= 0x80 && rtn <= 0xEF) {  // 0x80~0xEF为用户自定义错误码
                    proto_syslog(LOG_ERR, "RTN: 0x%02x (other error, user-defined YD/T 1363.3-2014)", rtn);
                } 
                else {
                    proto_syslog(LOG_ERR, "RTN: 0x%02x (unknown error, YD/T 1363.3-2014)", rtn);
                }
                break;
        }
        return -1;
    }
    
    uint8_t lengthHigh = fromAscii(frame[9], frame[10]);
    uint8_t lengthLow = fromAscii(frame[11], frame[12]);
    unsigned short lengthValue = (lengthHigh << 8) | lengthLow;

    unsigned short lenid = (lengthValue & 0xFF);        // 提取LENID的低8位 (D7-D0)
    uint8_t lchksum = (lengthValue >> 12) & 0x0F;       // 提取LCHKSUM的低4位 (D0-D3)
    unsigned short highLenid = (lengthValue >> 8) & 0x0F; // 提取LENID的高4位 (D11-D8)
    
    lenid |= (highLenid << 8); // 组合完整的LENID  -->12bit
    uint8_t calculatedLCHKSUM = calculateLCHKSUM(lenid);
    if (lchksum != calculatedLCHKSUM) {
        proto_syslog(LOG_ERR, "LCHKSUM error");
        return -1;
    }

    unsigned short chksumHigh = fromAscii(frame[length - 5], frame[length - 4]);
    unsigned short chksumLow = fromAscii(frame[length - 3], frame[length - 2]);
    unsigned short chksum = (chksumHigh << 8) | chksumLow;
    unsigned short calculatedCHKSUM = calculateCHKSUM(frame + 1, length - 6);
    if (chksum != calculatedCHKSUM) {
        proto_syslog(LOG_ERR, "CHKSUM error");
        return -1;
    }

    // 校验成功，提取DATA INFO，--注意：设置命令应答一般无INFO项，支持*infoLength = 0 
    uint8_t _header = 1 + 12; // SOI + 12 bytes header
    *infoLength = lenid/2;
    for (int i = 0; i < (*infoLength); i++) {
        info[i] = fromAscii(frame[_header + i * 2], frame[_header + 1 + i * 2]);
    }

    proto_syslog(LOG_INFO, "SOI: 0x%02x, VER: 0x%02x, ADR: 0x%02x, CID1: 0x%02x, CID2: 0x%02x, LENGTH: 0x%04x, INFO...  CHKSUM: 0x%04x, EOI: 0x%02x",
                    SOI, ver, adr, cid1, cid2, lengthValue, chksum, EOI);
    return 0;
}

static int extract_ydt1363_data(proto_ydt1363_poll *poll, uint8_t *rx_buff, int rx_len, long long cur_ms)
{
    int tag_sum = poll->tag_sum;
    tag_t *p_tag = NULL;
    ydt1363_parser_info_t *parser_info = NULL;
    uint8_t info_buff[2048] = {0}; // DATA INFO BUFF
    unsigned short info_len = 0;

    tag_value_t tag_value = {0};
    if (tag_sum > 0)
    {
        refresh_online_state((device_t *)(poll->p_tag[0]->dev_ptr));
    }

	if (ydt1363_parseFrame(rx_buff, rx_len, info_buff, &info_len) < 0){
		proto_syslog(LOG_ERR, "ydt1363 data Frame parse error!!!");
		return -1;
	}

	if (info_len == 0)  return 0; // 设置命令应答一般无INFO项，直接返回
    proto_syslog_hex(LOG_DEBUG, info_buff, info_len, "data info:: [%s]<-LEN(%d):", "data", info_len);

    for (int i = 0; i < tag_sum; i++)
    {
        p_tag = poll->p_tag[i];
        parser_info = (ydt1363_parser_info_t *)(p_tag->proto_ptr);
        switch (parser_info->data_type)
        {
        case TYPE_BOOL:
            tag_value.type = 0;
            if (parser_info->reversal != 0)
            {
                tag_value.value.to_int = (*(info_buff + parser_info->offset) != 0) ? 0 : 1;
            }
            else
            {
                tag_value.value.to_int = (*(info_buff + parser_info->offset) != 0) ? 1 : 0;
            }
            proto_set_tag(p_tag, tag_value, cur_ms);
            break;

        case TYPE_INT:
            GET_INT(tag_value, info_buff + parser_info->offset, parser_info->data_order);
            proto_set_tag(p_tag, tag_value, cur_ms);
            break;

        case TYPE_UINT:
            GET_UINT(tag_value, info_buff + parser_info->offset, parser_info->data_order);
            proto_set_tag(p_tag, tag_value, cur_ms);
            break;

        case TYPE_FLOAT:
            GET_FLOAT(tag_value, info_buff + parser_info->offset, parser_info->data_order);
            proto_set_tag(p_tag, tag_value, cur_ms);
            break;

        case TYPE_BITS:
            GET_BITS(tag_value, info_buff + parser_info->offset, parser_info->data_order);
            if (parser_info->reversal != 0)
            {
                tag_value.value.to_int = GET_BITS_IN_U16(~tag_value.value.to_int, parser_info->bit_st, parser_info->bit_len);
            }
            else
            {
                tag_value.value.to_int = GET_BITS_IN_U16(tag_value.value.to_int, parser_info->bit_st, parser_info->bit_len);
            }
            proto_set_tag(p_tag, tag_value, cur_ms);
            break;

        default:
            proto_syslog(LOG_ERR, "extract data error, date type isn't supported, data_type = %d, tag=%s, cid2:%d, offset:%d",
                         parser_info->data_type, p_tag->name, parser_info->cid2, parser_info->offset);
            break;
        }
    }

    return 0;
}

int ydt1363_proto_read(ydt1363_hd_t *hd, proto_ydt1363_poll *poll, void *rx_buff, int rx_buff_len)
{
	int rx_len = 0, info_len = 0;
    uint8_t info_buff[2048] = {0}; // COMMAND INFO BUFF

    if (poll->tag_sum <= 0) {
        proto_syslog(LOG_ERR, "ydt1363 period read no match read tag error!!!");
        return -1;
    }
    tag_t *r_tag = poll->p_tag[0];
    device_t *pd = (device_t *)r_tag->dev_ptr;
    template_t *template = pd->use_template;
    ydt1363_user_def_t *ydt1363_user_def = (ydt1363_user_def_t *)template->proto_ptr;    

    for (int i = GROUP; i < ITEM_MAX; i++) //>>读命令，只支持1字节在COMMAND INFO中赋值，此项可缺省
    {
        ydt1363_comman_info_item *comman_info_item = &poll->item[i];
        if (comman_info_item->exist == 1){
            comman_info_item->buff[0] = (uint8_t)poll->val;
            comman_info_item->len = 1;
            break;
        }
    }

    for (int i = 0; i < ITEM_MAX; i++)     // 按顺序拼接COMMAND_INFO
    {
        ydt1363_comman_info_item *comman_info_item = &poll->item[i];
        if (comman_info_item->exist == 1 && comman_info_item->len != 0){
            memcpy(info_buff + info_len, comman_info_item->buff, comman_info_item->len);
            proto_syslog_hex(LOG_DEBUG, comman_info_item->buff, comman_info_item->len, "[%s]->read. splice(%d)", comman_info_item->des, comman_info_item->len);
            info_len += comman_info_item->len;
        }
    }
    for (int i = GROUP; i < ITEM_MAX; i++)
    {
        ydt1363_comman_info_item *comman_info_item = &poll->item[i];
        memset(comman_info_item->buff, 0, sizeof(comman_info_item->buff));
        comman_info_item->len = 0;
    }
    proto_syslog_hex(LOG_DEBUG, info_buff, info_len, "period inquire command info::: [%s]->LEN(%d):", "command", info_len);
    
    uint8_t ver = stringToHex(ydt1363_user_def->ver);
    uint32_t adr = *(uint32_t *)pd->proto_ptr;
    ydt1363_createFrame(ver, (uint8_t)adr, ydt1363_user_def->cid1, poll->cid2, info_buff, info_len, poll->buff, &poll->len);
    	
	rx_len = ydt1363_read(hd, poll->buff, poll->len, rx_buff, rx_buff_len);
    return rx_len;
}

static int ydt1363_tag_write_dev(ydt1363_hd_t *hd, tag_t *w_tag, void *rx_buff, int rx_buff_len)
{
    int use_len = 0, rx_len = 0, info_len = 0;
    uint8_t info_buff[2048] = {0}; // COMMAND INFO BUFF

    channel_t *channel = (channel_t *)w_tag->chan_ptr;
    device_t *pd = (device_t *)w_tag->dev_ptr;
    template_t *template = pd->use_template;
    ydt1363_user_def_t *ydt1363_user_def = (ydt1363_user_def_t *)template->proto_ptr;

    proto_ydt1363_poll *poll = NULL;
    uint8_t *dest_buff = NULL;
    int *len = NULL;

    if (w_tag->poll_list_type == 0)
    {
        poll = (proto_ydt1363_poll *)channel->chunks_interv[w_tag->poll_list_num]->chunk_proto_ptr;
    }
    else
    {
        poll = (proto_ydt1363_poll *)channel->chunks_nonstop[w_tag->poll_list_num]->chunk_proto_ptr;
    }
    proto_syslog(LOG_DEBUG, "w_tag name %s ---> poll id: 0x%02x, tag_sum: %d ", w_tag->name, poll->id, poll->tag_sum);

	for (int i = 0; i < poll->tag_sum; i++){ 
		tag_t *write_tag = poll->p_tag[i];
		ydt1363_parser_info_t *parse_info = (ydt1363_parser_info_t *)write_tag->proto_ptr;
        ydt1363_comman_info_item *comman_info_item = NULL;
        int find = 0;
        for (int i = 0; i < ITEM_MAX; i++)
        {
            comman_info_item = &poll->item[i];
            if (comman_info_item->exist == 1 && strcmp(parse_info->item, comman_info_item->des) == 0){
                proto_syslog(LOG_DEBUG, "ydt1363 command info item match success: index=[%d]   %s", i, comman_info_item->des);
                find = 1;
                break;
            }
        }
        if (find == 0) continue;

	    dest_buff = &comman_info_item->buff[parse_info->offset];
        len = &comman_info_item->len;
	    switch (parse_info->data_type)
	    {
	    case TYPE_BOOL:
	        *dest_buff = (write_tag->write_cache.to_int != 0) ? 1 : 0;
			use_len = 1;
            *len += use_len;
	        break;

	    case TYPE_INT:
	        SET_INT(write_tag, dest_buff, use_len, parse_info->data_order);
            *len += use_len;
	        break;

	    case TYPE_UINT:
	        SET_UINT(write_tag, dest_buff, use_len, parse_info->data_order);
            *len += use_len;
	        break;

	    case TYPE_FLOAT:
	        SET_FLOAT(write_tag, dest_buff, use_len, parse_info->data_order);
            *len += use_len;
	        break;

	    default:
	        proto_syslog(LOG_ERR, "write reg type error, [id:%d, offset:%d, type=%d]", parse_info->id, parse_info->offset, parse_info->data_type);
	        break;
	    }
	}

    for (int i = 0; i < ITEM_MAX; i++)     // 按顺序拼接COMMAND_INFO
    {
        ydt1363_comman_info_item *comman_info_item = &poll->item[i];
        if (comman_info_item->exist == 1 && comman_info_item->len != 0){
            memcpy(info_buff + info_len, comman_info_item->buff, comman_info_item->len);
            proto_syslog_hex(LOG_DEBUG, comman_info_item->buff, comman_info_item->len, "[%s]->write. splice(%d)", comman_info_item->des, comman_info_item->len);
            info_len += comman_info_item->len;
        }
    }
    for (int i = GROUP; i < ITEM_MAX; i++)
    {
        ydt1363_comman_info_item *comman_info_item = &poll->item[i];
        memset(comman_info_item->buff, 0, sizeof(comman_info_item->buff));
        comman_info_item->len = 0;
    }
    proto_syslog_hex(LOG_DEBUG, info_buff, info_len, "set command info::: [%s]->LEN(%d):", "command", info_len);
    
    uint8_t ver = stringToHex(ydt1363_user_def->ver);
    uint32_t adr = *(uint32_t *)pd->proto_ptr;
    ydt1363_createFrame(ver, (uint8_t)adr, ydt1363_user_def->cid1, poll->cid2, info_buff, info_len, poll->buff, &poll->len);
	
	rx_len = ydt1363_write(hd, poll->buff, poll->len, rx_buff, rx_buff_len);
    return rx_len;
}

static int ydt1363_channel_check_write(ydt1363_hd_t *hd, channel_t *channel, int time_out)
{
	uint8_t rx_buff[4095 + 18] = {0}; 
    uint8_t info_buff[2048] = {0}; // DATA INFO BUFF
    unsigned short info_len = 0;
    
    int ret = sem_wait_time(&channel->write_sem, 0, time_out);
    if (ret == 0)
    {
        write_list_t *w_info = get_data_from_list(channel);
        if (w_info == NULL)
        {
            proto_syslog(LOG_ERR, "get data from list is null");
            return -1;
        }
        tag_t *write_tag = &w_info->w_tag;
        if (write_tag == NULL)
        {
            free(w_info);
            proto_syslog(LOG_ERR, "get data from list is null");
            return -2;
        }
		
        device_t *dev = (device_t *)write_tag->dev_ptr;
        proto_syslog(LOG_INFO, "get data from list is [%s.%s(type:%d, double=%f, long=%ld)]",
                     dev->no, write_tag->name, write_tag->data_type, write_tag->write_cache.to_float, write_tag->write_cache.to_int);

        ret = ydt1363_tag_write_dev(hd, write_tag, rx_buff, sizeof(rx_buff));
        if (ret < 0)
        {
            w_info->pw_tag->last_state_code = w_info->pw_tag->state_code;
            w_info->pw_tag->state_code = TAG_STATE_BAD;
            proto_syslog(LOG_ERR, "ydt1363 write tag error, tag:%s, ret = %d", w_info->pw_tag->name, ret);
            BUSINESS_LOG(BLOG_ERR, DAS_ID_YDT1363_3, write_tag->dev_tag_name, "[数采]: %s 通道写 tag: %s 失败", channel->channel, w_info->pw_tag->name);
        }
        else{
            BUSINESS_LOG(BLOG_NOTICE, DAS_ID_YDT1363_3, write_tag->dev_tag_name, "[数采]: %s 通道写 tag: %s 成功", channel->channel, w_info->pw_tag->name);
			int rx_len = ret;
            if (ydt1363_parseFrame(rx_buff, rx_len, info_buff, &info_len) < 0){  //ydt1363协议，设置命令无更新tag需求，只需确保设置命令应答正确即可
                proto_syslog(LOG_ERR, "ydt1363 set command answer error, set fail!!!");
                return -1;
            }
            else{
                proto_syslog(LOG_INFO, "ydt1363 set command answer normal, set success");
            }
			refresh_online_state(dev);

            w_info->pw_tag->last_state_code = w_info->pw_tag->state_code;
            w_info->pw_tag->state_code = TAG_STATE_GOOD;
            sync_wtag_read_cache(w_info);
        }
        free(w_info);
    }

    return 0;
}

static int ydt1363_channel(channel_t *channel)
{
	long long cur_ms = 0;
	int err_counter = 0;
	proto_ydt1363_poll *poll_info = NULL;
	uint8_t rx_buff[4095 + 18] = {0};
	device_t *dev = NULL;
	
	ydt1363_hd_t *hd = ydt1363_proto_init((ydt1363_proto_t *)channel->proto_ptr);
	if (hd == NULL)
	{
		proto_syslog(LOG_CRIT, "[数采]:ydt1363_proto_init error in the channel:[%s]", channel->channel);
        BUSINESS_LOG(BLOG_ERR, DAS_ID_YDT1363_1, channel->channel, "[数采]: %s 通道初始化失败", channel->channel);
		return -1;
	}
	channel->s = hd->fd;
	proto_syslog(LOG_INFO, "ydt1363_proto_init success, in the channel:[%s]", channel->channel);
    BUSINESS_LOG(BLOG_ERR, DAS_ID_YDT1363_1, channel->channel, "[数采]: %s 通道初始化成功", channel->channel);
	
	while (1)
	{
	    int contine_interv = 0;
		do
        {
            contine_interv = 0;
			check_tag_state(channel);
			if(ydt1363_channel_check_write(hd, channel, 100) < 0)
				err_counter++;
			
			cur_ms = clock_get_ms();
			for (int i = 0; i < channel->chunks_interv_size; i++)
			{
				if(ydt1363_channel_check_write(hd, channel, 0) < 0)
					err_counter++;
				
				dev = channel->chunks_interv[i]->dev_ptr;
				poll_info = (proto_ydt1363_poll *)channel->chunks_interv[i]->chunk_proto_ptr;
				if ((poll_info->id == 0) && (cur_ms - channel->chunks_interv[i]->last_poll) >= channel->chunks_interv[i]->poll_interv) // id=0，真正的周期性查询指令
				{
					channel->chunks_interv[i]->last_poll = cur_ms;
					ydt1363_set_time(hd, dev->timeout, dev->interval);
					channel_lock(channel);
					int ret = ydt1363_proto_read(hd, poll_info, rx_buff, sizeof(rx_buff));
					channel_unlock(channel);
					if (ret > 0)
					{
                        BUSINESS_LOG(BLOG_NOTICE, DAS_ID_CUSTOM_2, channel->channel, "[数采]: %s 通道回复正常", channel->channel);
					    int rx_len = ret;
					    cur_ms = clock_get_ms();
                        if (poll_info->fun != NULL) // 有自定义解析
                        {
                            poll_info->fun(poll_info);
                            refresh_online_state(channel->chunks_interv[i]->dev_ptr);
                        }
                        else{
                            extract_ydt1363_data(poll_info, rx_buff, rx_len, cur_ms);
                        }
						refresh_online_state(dev);
						err_counter = 0;
					}
					else if (ret == 0)
					{
						proto_syslog(LOG_ERR, "%s:ydt1363 read timeout", channel->channel);
                        BUSINESS_LOG(BLOG_ERR, DAS_ID_YDT1363_2, channel->channel, "[数采]: %s 通道回复超时", channel->channel);
						err_counter++;
						check_online_state(dev);
					}
					else
					{
						proto_syslog(LOG_ERR, "%s:ydt1363 read error(ret = %d), info:%s, errno:%d", channel->channel, ret, strerror(errno), errno);
                        BUSINESS_LOG(BLOG_ERR, DAS_ID_YDT1363_2, channel->channel, "[数采]: %s 通道回复失败", channel->channel);
						check_online_state(dev);
						if (errno > 0) // io错误
						{
                            BUSINESS_LOG(BLOG_ERR, DAS_ID_YDT1363_2, channel->channel, "[数采]: %s 通道出现系统级IO错误, 请检查 物理连接或设备 是否正常", channel->channel);
							goto ydt1363_channel_exit;
						}
						else // 协议错误
						{
                            BUSINESS_LOG(BLOG_ERR, DAS_ID_YDT1363_2, channel->channel, "[数采]: %s 通道出现协议错误, 请检查 回复数据 是否正确", channel->channel);
							err_counter++;
						}
					}
					contine_interv = 1;
				}

				if (err_counter > COMMUNIT_RETRY_MAX)
				{
					goto ydt1363_channel_exit;
				}
			}
	    } while ((contine_interv != 0));

		// interv空闲
		if (channel->chunks_nonstop_size > 0)
		{
			int nonstop_index = channel->chunks_nonstop_index;
			dev = channel->chunks_nonstop[nonstop_index]->dev_ptr;
			poll_info = (proto_ydt1363_poll *)channel->chunks_nonstop[nonstop_index]->chunk_proto_ptr;
			if (poll_info->id == 0)
			{
				channel_lock(channel);
				int ret = ydt1363_proto_read(hd, poll_info, rx_buff, sizeof(rx_buff)); // 读数据
				channel_unlock(channel);
				if (ret > 0)
				{
					int rx_len = ret;
					cur_ms = clock_get_ms();
                    if (poll_info->fun != NULL) // 有自定义解析
                    {
                        poll_info->fun(poll_info);
                        //refresh_online_state(channel->chunks_nonstop[i]->dev_ptr);
                    }
                    else{
                        extract_ydt1363_data(poll_info, rx_buff, rx_len, cur_ms);
                    }
					refresh_online_state(dev);
					err_counter = 0;
				}
				else if (ret == 0)
				{
					proto_syslog(LOG_ERR, "%s:ydt1363 read timeout", channel->channel);
					err_counter++;
					check_online_state(dev);
				}
				else
				{
					proto_syslog(LOG_ERR, "%s:ydt1363 read error(ret = %d), info:%s, errno:%d", channel->channel, ret, strerror(errno), errno);
					check_online_state(dev);
					if (errno > 0) // io错误
					{
						goto ydt1363_channel_exit;
					}
					else // 协议错误
					{
						err_counter++;
					}
				}
			
				if (err_counter > COMMUNIT_RETRY_MAX)
				{
					goto ydt1363_channel_exit;
				}
				
				nonstop_index++;
				if (nonstop_index >= channel->chunks_nonstop_size)
				{
					channel->chunks_nonstop_index = 0;
				}
				else
				{
					channel->chunks_nonstop_index = nonstop_index;
				}
			}
		}  
	}
	
ydt1363_channel_exit:
	if (hd != NULL)
	{
		ydt1363_proto_deinit(hd);
	}
	hd = NULL;

	return -1;
}

#ifdef EN_EMS
static const char *dev_no_map(const device_t *dev)
{
    if (strstr(dev->user_def, "is_lc=1;"))
    {
        return "LC";
    }
    else if (strstr(dev->user_def, "is_ems=1;"))
    {
        return "EMS";
    }
    else
    {
        const char *dot = strchr(dev->no, '.');
        if (dot != NULL && strlen(dot + 1) > 0)
        {
            return dot + 1;
        }
    }
    return dev->no;
}

int lcua_write_node(UA_Client *client, tag_t *tag)
{
    UA_Int32 int_value = 0;
    UA_Double double_value = 0;
    UA_Variant variant;
    char node_buff[TAG_NAME_LEN + 10] = {0};
    UA_Variant_init(&variant);
    if (tag->data_type == 0)
    {
        int_value = tag->write_cache.to_int;
        UA_Variant_setScalar(&variant, &int_value, &UA_TYPES[UA_TYPES_INT32]);
    }
    else
    {
        double_value = tag->write_cache.to_float;
        UA_Variant_setScalar(&variant, &double_value, &UA_TYPES[UA_TYPES_DOUBLE]);
    }

    if (is_dev_tag(tag) != 0)
    {
        const char *dot = strchr(tag->dev_tag_name, '.');
        if (dot != NULL)
        {
            snprintf(node_buff, sizeof(node_buff), "%s", dot + 1);
        }
    }
    else
    {
        const char *p = dev_no_map((device_t *)tag->dev_ptr);
        snprintf(node_buff, sizeof(node_buff), "%s.%s", p, tag->name);
    }
    UA_StatusCode statusCode = UA_Client_writeValueAttribute(client, UA_NODEID_STRING(1, node_buff), &variant);
    return (statusCode == UA_STATUSCODE_GOOD) ? 0 : -2;
}

int lcua_report_heartbeat(UA_Client *client, lc_it* lc_ptr)//1：被禁用，0：启用
{
    int is_disable = 0;
    if(lc_ptr->disabled == 1) is_disable = 255;
    else is_disable = 1;

    //只考虑了opcua通道接emu从机
    UA_Int32 int_value = is_disable;
    UA_Variant variant;
    UA_Variant_init(&variant);
    UA_Variant_setScalar(&variant, &int_value, &UA_TYPES[UA_TYPES_INT32]);
    UA_StatusCode statusCode = UA_Client_writeValueAttribute(client, UA_NODEID_STRING(1, EMS_HEARTBEAT), &variant);
    return (statusCode == UA_STATUSCODE_GOOD) ? 0 : -2;
}


static int lcua_channel_check_write(UA_Client *client, channel_t *channel, int time_out)
{
    int ret = sem_wait_time(&channel->write_sem, 0, time_out);
    while (ret == 0)
    {
        write_list_t *w_info = get_data_from_list(channel);
        if (w_info == NULL)
        {
            proto_syslog(LOG_ERR, "get data from list is null");
            return -1;
        }
        tag_t *write_tag = &w_info->w_tag;
        if (write_tag == NULL)
        {
            proto_syslog(LOG_ERR, "get data from list is null");
            free(w_info);
            return -2;
        }
        device_t *dev = (device_t *)write_tag->dev_ptr;
        proto_syslog(LOG_INFO, "get data from list is [%s.%s(type:%d, double=%f, int=%ld)]",
                     dev->no, write_tag->name, write_tag->data_type, write_tag->write_cache.to_float, write_tag->write_cache.to_int);

        ret = lcua_write_node(client, write_tag);
        if (0 > ret)
        {
            w_info->pw_tag->last_state_code = w_info->pw_tag->state_code;
            w_info->pw_tag->state_code = TAG_STATE_BAD;
            proto_syslog(LOG_ERR, "can write tag error, tag:%s, ret = %d", write_tag->name, ret);
            free(w_info);
            return -3;
        }
        else
        {
            w_info->pw_tag->last_state_code = w_info->pw_tag->state_code;
            w_info->pw_tag->state_code = TAG_STATE_GOOD;
            sync_wtag_read_cache(w_info);
        }
        free(w_info);
        ret = sem_wait_time(&channel->write_sem, 0, 0);

    }

    return 0;
}

typedef struct channel_dev_ua_req
{
    device_t *dev;
    UA_ReadRequest request;
    tag_t **tags;
    int *tag_subscription_status;// 订阅成功标志
    UA_UInt32 subscriptionId; 
    //
    UA_NodeId online_node_id;
} channel_dev_ua_req;

static void handler_data_changed(UA_Client *client, UA_UInt32 subId, void *subContext, UA_UInt32 monId, void *monContext, UA_DataValue *value)
{
    tag_t *tag = (tag_t *)monContext;
    tag_value_t tag_value = {};
    UA_Variant val = value->value;

    if (val.type == &UA_TYPES[UA_TYPES_INT32])
    {
        tag_value.type = 0;
        tag_value.value.to_int = *(UA_Int32 *)val.data;
        proto_syslog(LOG_NOTICE, "%s.%s.%s changed, value:%ld", ((channel_t *)subContext)->channel, ((device_t *)tag->dev_ptr)->no, tag->name, tag_value.value.to_int);
    }
    else if (val.type == &UA_TYPES[UA_TYPES_DOUBLE])
    {
        tag_value.type = 1;
        tag_value.value.to_float = *(UA_Double *)val.data;
        proto_syslog(LOG_NOTICE, "%s.%s.%s changed, value:%f", ((channel_t *)subContext)->channel, ((device_t *)tag->dev_ptr)->no, tag->name, tag_value.value.to_float);
    }
    else
    {
        proto_syslog(LOG_ERR, "data type error");
        return;
    }
    proto_set_tag_no_scale_offset(tag, tag_value, clock_get_ms());
    if (g_usercfg_variant.EnDataCollect) //使能数据汇聚功能
    {
        proto_syslog(LOG_NOTICE,"g_usercfg_variant.EnDataCollect:%d\n",g_usercfg_variant.EnDataCollect);
    }
}

int ua_node_already_exist(UA_Client *client, const UA_NodeId *node_id)
{
    UA_NodeClass nodeClass;
    int ok = 0;
    UA_StatusCode retval = UA_Client_readNodeClassAttribute(client, *node_id, &nodeClass);
    if (UA_STATUSCODE_GOOD == retval)
    {
        ok = 1;
    }
    return ok;
}

int creat_data_change_sub(UA_Client *client, channel_t *channel, channel_dev_ua_req *req, int req_len)
{
    int ret = 0;
    UA_CreateSubscriptionRequest request = UA_CreateSubscriptionRequest_default();
    UA_CreateSubscriptionResponse response = UA_Client_Subscriptions_create(client, request,
                                                                            channel, NULL, NULL);

    if (response.responseHeader.serviceResult == UA_STATUSCODE_GOOD)
    {
        proto_syslog(LOG_NOTICE, "Create %s subscription succeeded", channel->channel);
    }
    else
    {
        proto_syslog(LOG_ERR, "Create %s subscription error", channel->channel);
        return -1;
    }

    time_t start_time = time(NULL);
    time_t cur_time = start_time;
    for (int d = 0; d < req_len; d++)
    {
        req[d].subscriptionId = response.subscriptionId; 
        for (int r = 0; r < req[d].request.nodesToReadSize; r++)
        {   // 下面这个do while是为了防止服务器能够连接了但是点位订阅还没有就绪导致没有订阅数据问题
            while (cur_time - start_time < 20) // ua连上到服务器点位完全就绪设置为30s，20s后还读不到代表点位的确不存在
            {
                cur_time = time(NULL);
                if (ua_node_already_exist(client, &req[d].request.nodesToRead[r].nodeId) != 0) // 点位存在
                {
                    break;
                }
                else
                {
                    proto_syslog(LOG_WARNING, "waiting for node[%s], %ld, %ld", req[d].tags[r]->dev_tag_name, cur_time, start_time);
                    sleep(1);
                }
            }

            UA_MonitoredItemCreateRequest monRequest =
                UA_MonitoredItemCreateRequest_default(req[d].request.nodesToRead[r].nodeId);

            UA_MonitoredItemCreateResult monResponse =
                UA_Client_MonitoredItems_createDataChange(client, response.subscriptionId,
                                                          UA_TIMESTAMPSTORETURN_BOTH,
                                                          monRequest, req[d].tags[r], handler_data_changed, NULL);
            if (monResponse.statusCode == UA_STATUSCODE_GOOD)
            {
                proto_syslog(LOG_NOTICE, "Monitoring %s DataChange succeeded", req[d].tags[r]->dev_tag_name);
                req[d].tag_subscription_status[r] = 1;
                // 订阅成功
            }
            else
            {
                proto_syslog(LOG_ERR, "Monitoring %s DataChange error", req[d].tags[r]->dev_tag_name);
            }
        }
    }

    UA_Client_run_iterate(client, 1000);

    return ret;
}



extern tag_t **get_ext_tag(void *ext_par, int *out_sum);
static channel_dev_ua_req *new_dev_re_req(channel_t *channel, int *ret)
{
    if (channel->devs_size < 1)
    {
        *ret = -1;
        return NULL;
    }
    channel_dev_ua_req *req = calloc(channel->devs_size, sizeof(channel_dev_ua_req));
    if (req == NULL)
    {
        *ret = -2;
        return NULL;
    }

    //
    device_t *dev = NULL;
    int tag_sum = 0;
    UA_ReadValueId *id_arr = NULL;
    const char *dev_no = NULL;
    char buff[TAG_NAME_LEN + DEV_NAME_LEN] = {0};
    tag_t **ext_tags = NULL;
    int ext_tag_sum = 0;
    for (int i = 0; i < channel->devs_size; i++)
    {
        dev = channel->devs[i];
        //
        dev_no = dev_no_map(dev); // 这儿主要区分是否需要映射
        if (dev->ext_par != NULL)
        {
            ext_tags = get_ext_tag(dev->ext_par, &ext_tag_sum);
        }
        else
        {
            ext_tags = NULL; 
        }
        //
        tag_sum = ext_tag_sum + dev->tags.ro_tags_size + dev->tags.rw_tags_size + dev->tags.sys_info_tags_size;
        id_arr = calloc(tag_sum, sizeof(UA_ReadValueId));
        req[i].tags = calloc(tag_sum, sizeof(tag_t **));
        req[i].tag_subscription_status = calloc(tag_sum, sizeof(int));
        //
        int index = 0;
        if (ext_tags != NULL)
        {
            for (int t = 0; t < ext_tag_sum; t++)
            {
                UA_ReadValueId_init(&id_arr[index]);
                id_arr[index].attributeId = UA_ATTRIBUTEID_VALUE;
                id_arr[index].nodeId = UA_NODEID_STRING_ALLOC(1, strstr(ext_tags[t]->dev_tag_name, ".") + 1);
                req[i].tags[index] = ext_tags[t];
                proto_syslog(LOG_NOTICE, "====%s:ext_tags[%d]=%s", dev->no, t, ext_tags[t]->dev_tag_name);
                req[i].tag_subscription_status[index] = 0;
                index++;
            }
        }
        for (int t = 0; t < dev->tags.ro_tags_size; t++)
        {
            UA_ReadValueId_init(&id_arr[index]);
            id_arr[index].attributeId = UA_ATTRIBUTEID_VALUE;
            snprintf(buff, sizeof(buff), "%s.%s", dev_no, dev->tags.ro_tags[t]->name);
            id_arr[index].nodeId = UA_NODEID_STRING_ALLOC(1, buff);
            req[i].tags[index] = dev->tags.ro_tags[t];
            proto_syslog(LOG_NOTICE, "====%s:ro_tags[%d]=%s", dev->no, t, buff);
            req[i].tag_subscription_status[index] = 0;
            index++;
        }

        for (int t = 0; t < dev->tags.rw_tags_size; t++)
        {
            UA_ReadValueId_init(&id_arr[index]);
            id_arr[index].attributeId = UA_ATTRIBUTEID_VALUE;
            snprintf(buff, sizeof(buff), "%s.%s", dev_no, dev->tags.rw_tags[t]->name);
            id_arr[index].nodeId = UA_NODEID_STRING_ALLOC(1, buff);
            req[i].tags[index] = dev->tags.rw_tags[t];
            proto_syslog(LOG_NOTICE, "====%s:rw_tags[%d]=%s", dev->no, t, buff);
            req[i].tag_subscription_status[index] = 0;
            index++;
        }

        for (int t = 0; t < dev->tags.sys_info_tags_size; t++)
        {
            UA_ReadValueId_init(&id_arr[index]);
            id_arr[index].attributeId = UA_ATTRIBUTEID_VALUE;
            snprintf(buff, sizeof(buff), "%s.%s", dev_no, dev->tags.sys_info_tags[t]->name);
            id_arr[index].nodeId = UA_NODEID_STRING_ALLOC(1, buff);
            req[i].tags[index] = dev->tags.sys_info_tags[t];
            if (strcmp(dev->tags.sys_info_tags[t]->name, "Online") == 0) //特殊处理
            {
                req[i].online_node_id = UA_NODEID_STRING_ALLOC(1, buff);
            }
            proto_syslog(LOG_NOTICE, "====%s:sys_info_tags[%d]=%s", dev->no, t, buff);
            req[i].tag_subscription_status[index] = 0;
            index++;
        }

        req[i].dev = dev;
        UA_ReadRequest_init(&req[i].request);
        req[i].request.nodesToRead = id_arr;
        req[i].request.nodesToReadSize = tag_sum;
    }

    *ret = channel->devs_size;
    return req;
}

static int check_lc_device_online(UA_Client *ua_client, channel_dev_ua_req *req, int req_len)
{
    UA_Variant val;
    int online = 0;
    for (int i = 0; i < req_len; i++)
    {
        online = 0;
        UA_Variant_init(&val);
        UA_StatusCode retval = UA_Client_readValueAttribute(ua_client, req[i].online_node_id, &val);
        if (retval == UA_STATUSCODE_GOOD)
        {
            if ((val.type == &UA_TYPES[UA_TYPES_INT32]) && (*(UA_Int32 *)val.data != 0))
            {
                online = 1;
            }
        }
        UA_Variant_clear(&val);

        if (online != 0)
        {
            refresh_online_state(req[i].dev);
        }
        else
        {
            check_online_state(req[i].dev);
        }
    }

    return 0;
}


static void retry_failed_subscriptions(UA_Client *client, channel_t *channel, channel_dev_ua_req *req, int req_len) {

    for (int i = 0; i < req_len; i++) {
        for (int j = 0; j < req[i].request.nodesToReadSize; j++) {
            if (req[i].tag_subscription_status[j] == 0) {
                UA_MonitoredItemCreateRequest monRequest = UA_MonitoredItemCreateRequest_default(req[i].request.nodesToRead[j].nodeId);

                proto_syslog(LOG_NOTICE, "Retrying Monitoring item for NodeId: %s, subscriptionId: %u",
                             req[i].tags[j]->dev_tag_name, req[i].subscriptionId);

                UA_MonitoredItemCreateResult monResponse =
                    UA_Client_MonitoredItems_createDataChange(client, req[i].subscriptionId,
                                                              UA_TIMESTAMPSTORETURN_BOTH,
                                                              monRequest, req[i].tags[j], handler_data_changed, NULL);

                if (monResponse.statusCode == UA_STATUSCODE_GOOD) {
                    proto_syslog(LOG_NOTICE, "Retry subscription Monitoring succeeded for %s", req[i].tags[j]->dev_tag_name);
                    req[i].tag_subscription_status[j] = 1;
                } else {
                    proto_syslog(LOG_ERR, "Retry subscription Monitoring failed for %s, statusCode: 0x%08X",
                                 req[i].tags[j]->dev_tag_name, monResponse.statusCode);
                }
            }
        }
    }
}

lc_it* channel_contain_lc(channel_t *channel)
{
    // 从机列表
    pcs_ctrl_var_t *var = get_pcs_ctrl_var();
    int lc_cnt = var->lc_list_len;
    lc_it *lc_list_tmp = var->lc_list;

    for (int i = 0; i < channel->devs_size; i++)
    {
        if (strstr(channel->devs[i]->dev_type, "LC") != NULL)
        {
            for (int j = 0; j < lc_cnt; j++)
            {
                if (strcmp(channel->devs[i]->no, lc_list_tmp[j].no) == 0)
                {
                    return &lc_list_tmp[j];
                }
            }
        }
    }
    return NULL;
}

static int get_ps_by_chan_name(const char *cfg, const char *url, char *name, int name_max, char *pass, int pass_max)
{
/*/app/config/opcua_chan_cfg.json指定opcua的用户名密码，格式：
{
    "info" : [
        {
            "url" : "frps.lnxall.com:38930",
            "username" : "admin",
            "password" : "lnxall123"
        }
    ]
}
*/
	char *cfgdat = read_file_data(cfg);
	if (cfgdat == NULL)
		return -1;

    // 解析 JSON 字符串
    cJSON *root = cJSON_Parse(cfgdat);
    free(cfgdat);
    if (root == NULL)
    {
        return -2;
    }

    // 获取 info 数组
    cJSON *info_array = cJSON_GetObjectItem(root, "info");
    if (info_array == NULL || !cJSON_IsArray(info_array))
    {
        cJSON_Delete(root);
        return -3;
    }

    // 遍历数组查找匹配的 url
    int array_size = cJSON_GetArraySize(info_array);
    int found = 0;

    for (int i = 0; i < array_size; i++)
    {
        cJSON *item = cJSON_GetArrayItem(info_array, i);

        cJSON *json_url = cJSON_GetObjectItem(item, "url");
        cJSON *json_username = cJSON_GetObjectItem(item, "username");
        cJSON *json_password = cJSON_GetObjectItem(item, "password");

        // 检查 url 是否匹配
        if (cJSON_IsString(json_url) && json_url->valuestring != NULL &&
            strcmp(json_url->valuestring, url) == 0)
        {

            // 提取 username
            if (cJSON_IsString(json_username) && json_username->valuestring != NULL)
            {
                strncpy(name, json_username->valuestring, name_max - 1);
                name[name_max - 1] = '\0';
            }
            else
            {
                cJSON_Delete(root);
                return -4;
            }

            // 提取 password
            if (cJSON_IsString(json_password) && json_password->valuestring != NULL)
            {
                strncpy(pass, json_password->valuestring, pass_max - 1);
                pass[pass_max - 1] = '\0';
            }
            else
            {
                cJSON_Delete(root);
                return -5;
            }

            found = 1;
            break;
        }
    }

    cJSON_Delete(root);

    // 未找到匹配的 url
    if (!found)
    {
        return -6;
    }

    return 0;
}

static int opcua_channel(channel_t *channel)
{

    // pthread_setcanceltype(PTHREAD_CANCEL_ENABLE, NULL);
    int dev_sum = 0;
    char name[64] = {0};
    char pass[64] = {0};
    channel_dev_ua_req *req = new_dev_re_req(channel, &dev_sum);
    if (req == NULL)
    {
        proto_syslog(LOG_ERR, "new_dev_re_req error:%d", dev_sum);
        return -1;
    }
    //
    UA_Client *ua_client = NULL;
    lc_it* lc_ptr = channel_contain_lc(channel);
    while (1)
    {
        char *url = calloc(1, strlen(channel->channel) + 20);
        if (url == NULL)
        {
            proto_syslog(LOG_ERR, "calloc error");
            sleep(1);
            continue;
        }
        
        sprintf(url, "opc.tcp://%s", channel->channel);
        if (get_ps_by_chan_name("/app/config/opcua_chan_cfg.json", channel->channel, name, sizeof(name), pass, sizeof(pass)) == 0)
        {
            ua_client = ua_connect_opcua_ps(url, name, pass);
        }
        else
        {
            ua_client = ua_connect_opcua(url);
        }
        free(url);
        if (ua_client == NULL)
        {
            proto_syslog(LOG_CRIT, "UA通道[%s]连接不上", channel->channel);
            for (int i = 0; i < dev_sum; i++)
            {
                check_online_state(req[i].dev);
            }
            sleep(1);
            continue;
        }
        creat_data_change_sub(ua_client, channel, req, dev_sum);
        check_lc_device_online(ua_client, req, dev_sum);
        int err = 0;
        uint64_t last_time = 0;
        uint64_t retry_time = 0;
        uint64_t second_t = 0;
        while (1)
        {
            lcua_channel_check_write(ua_client, channel, 50);
            if (UA_STATUSCODE_GOOD != UA_Client_run_iterate(ua_client, 50))
            {
                err++;
                if (err >= 5)
                {
                    ua_client_delete(ua_client);
                    break;
                }
            }
            else
            {
                err = 0;
            }

            uint64_t cur_ms = clock_get_ms();
            if (cur_ms - second_t > 999)
            {
                second_t = cur_ms;
                if (lc_ptr != NULL)
                {
                    lcua_report_heartbeat(ua_client, lc_ptr);
                }
                check_lc_device_online(ua_client, req, dev_sum);
            }

            if (cur_ms - retry_time > 10000) // 每 10 秒重试一次
            {
                retry_time = cur_ms;
                retry_failed_subscriptions(ua_client, channel, req ,dev_sum);
            }
            if (cur_ms - last_time < (60 * 1000))
            {
                continue;
            }
            last_time = cur_ms;
            //
            UA_Variant val;
            tag_value_t tag_value = {.type = 0};
            for (int i = 0; i < dev_sum; i++)
            {
                if (dev_is_online(req[i].dev) == 0) // 离线了不读取了
                {
                    continue;
                }
                for (int j = 0; j < req[i].request.nodesToReadSize; j++)
                {
                    UA_Variant_init(&val);
                    UA_StatusCode retval = UA_Client_readValueAttribute(ua_client, req[i].request.nodesToRead[j].nodeId, &val);
                    if (retval == UA_STATUSCODE_GOOD)
                    {
                        if (val.type == &UA_TYPES[UA_TYPES_INT32])
                        {
                            tag_value.type = 0;
                            tag_value.value.to_int = *(UA_Int32 *)val.data;
                        }
                        else if (val.type == &UA_TYPES[UA_TYPES_DOUBLE])
                        {
                            tag_value.type = 1;
                            tag_value.value.to_float = *(UA_Double *)val.data;
                        }
                        else
                        {
                            continue;
                        }
                        proto_set_tag_no_scale_offset(req[i].tags[j], tag_value, cur_ms);
                    }
                    else
                    {
                        // check_online_state(req[i].dev);
                    }
                    UA_Variant_clear(&val);
                }
            }
            // check_tag_state(channel); //不用检查
        }
    }
    return -2;
}
#endif

void *prot_transf_loop(void *param)
{
    // pthread_setcanceltype(PTHREAD_CANCEL_ENABLE, NULL);
    channel_t *channel = (channel_t *)param;

    int ret = -1;
    
    while (1)
    {
        proto_syslog(LOG_NOTICE, "the channel:[%s] start, proto_type = %d", channel->channel, channel->proto_type);
        BUSINESS_LOG(BLOG_NOTICE, DAS_ID_0002, NULL, "[数采]: %s 通道启动, 通道类型 %d (2:MODBUS_TCP; 3:MODBUS_RTU; 4:DLT645_2007; 5:DLT645_1997; 7:UPS; 8:BMSER_BIN; 9:CAN; 13:IEC104; 14:CUSTOM; 15:YDT1363_2014;)", channel->channel, channel->proto_type);
        channel->s = -1; //初始无效值
        switch (channel->proto_type)
        {
        case PROTOCOL_TYPE_NONE:
            break;

        case PROTOCOL_TYPE_SCRIPT:
            break;

        case PROTOCOL_TYPE_MODBUS_TCP:
        case PROTOCOL_TYPE_MODBUS_RTU:
        case PROTOCOL_TYPE_MODBUS_RTU_RCRC:
#ifdef EN_EMS
            if (contain_meter(channel) != 0)
            {
                ret = net_get_modbus_task(channel);
            }
            ret = modbus_channel(channel);
#endif
            break;

        case PROTOCOL_TYPE_DLT645_2007:
        case PROTOCOL_TYPE_DLT645_1997:
            ret = dlt645_channel(channel); //dlt不判断返回值 由接收广播通道判断设备是否在线
#ifdef EN_EMS
            if(channel->some_sign == 0)
            {
                channel->some_sign = 1;
                pthread_t tid;
                if (0 == pthread_create(&tid,NULL,net_get_dlt_task,channel))
                {
                    pthread_setname_np(tid, "NetGetDltTask");
                }
            }
            ret = 0;
#endif
            break;
        case PROTOCOL_TYPE_EMS:
            //set_priority();
            ret = ems_channel(channel);
            break;
#ifdef EN_EMS
        case PROTOCOL_TYPE_OPCUA:
            ret = opcua_channel(channel);
            break;
#endif
        case PROTOCOL_TYPE_LC:
            ret = lc_channel(channel);
            break;

        case PROTOCOL_TYPE_UPS:
            ret = ups_channel(channel);
            break;
        case PROTOCOL_TYPE_BMSER_BIN:
            if (NULL != strstr(((bmser_proto_s *)channel->proto_ptr)->ip, "bmser_lib"))
            {
                ret = bmser_bin_lib_channel(channel);
            }
            else
            {
                ret = bmser_bin_channel(channel);
            }
            break;
        case PROTOCOL_TYPE_CAN:
            ret = can_channel(channel);
            break;

        case PROTOCOL_TYPE_IEC104_CLIENT:
            iec104_client_channel(channel);
            break;
		
        case PROTOCOL_TYPE_CUSTOM:
            custom_channel(channel);
            break;

        case PROTOCOL_TYPE_YDT1363_2014:
            ydt1363_channel(channel);
            break;
            
        default:
            proto_syslog(LOG_ERR, "%s, the proto_type(%d) isn't supported", channel->channel, channel->proto_type);
            BUSINESS_LOG(BLOG_ERR, DAS_ID_0002, NULL, "[数采]: %s 通道,不支持的通道类型 %d", channel->channel, channel->proto_type);       
            break;
        }

        if (ret < 0)
        {
            for (size_t i = 0; i < channel->devs_size; i++)
            {
                check_online_state(channel->devs[i]);
            }
            //check_tag_state(channel);
        }
        //check_write_sleep(channel);
        sleep(3);
    }
    return NULL;
}
