#ifndef __CLOUD_MQTT_COMMON_H
#define __CLOUD_MQTT_COMMON_H

/***************HEADER FILES****************/
#include <string.h>
#include <syslog.h>
#include "ipc_session.h"
#include "lnxall_list.h"
#include "mqtt_session.h"
#include "dy_utils/dy_common.h"
#include "dy_utils/dy_db.h"
#include "dy_utils/dy_ipc.h"
#include "dy_utils/led.h"
#include "rdkafka.h"
#include "lnxall_ubuslog.h"
#define TOPIC_DEFAULT_EVT_HEARTBEAT "Evt_HeartBeat"
#define TOPIC_DEFAULT_RSP_HEARTBEAT "Rsp_HeartBeat"
#define TOPIC_DEFAULT_LOGIN "LogIn"
#define TOPIC_DEFAULT_RSP_LOGIN "Rsp_LogIn"

/* cloud mqtt force write database */
#define CLOUD_MQTT_FORCE_WRITEDB "/tmp/CMFORCE_WRITEDB"

#define MAX_UPLINK_CHANNEL_CNT 8
struct _cloud_mqtt_var_s;

typedef enum _COMM_STATUS_E
{
    COMM_STATUS_IDLE = 0,   //未连接
    COMM_STATUS_CONNECTING, //连接中
    COMM_STATUS_CONNECTED,  //已连接
    COMM_STATUS_LOGIN,      //已经登录
} COMM_STATUS_E;

typedef enum
{
    LOGIN_PASS = 0,
    LOGIN_REFUSE = 501,
} login_result_code;

typedef enum
{
    LINK_TYPE_MQTT = 0,
    LINK_TYPE_KAFKA,
} link_type_e;

typedef struct
{
    int serial_num;
    time_t timestamp;
    login_result_code error_code;
} login_info_t;

typedef struct
{
    char *pub_topic;
    char *sub_topic;
    int interval; //如果不需要周期发送则设置interval为0
    char *payload;
    int last_ts;
    int mi;
} cloud_mqtt_hb_login_msg_t;

typedef struct
{
    char *internal_topic;
    char *external_topic;
} topic_pair_t;

typedef struct
{
    int uplink_cnt;
    topic_pair_t *uplink;
    int downlink_cnt;
    topic_pair_t *downlink;
} topic_maps_t;

typedef struct kafka_topic_node_s
{
    struct list_head list;
    int tid;
    rd_kafka_topic_t *rkt_topic;
    char topic[0];
} kafka_topic_node_t;

typedef struct
{
    char *property;
    char *value;
} property_t;

typedef struct
{
    mqtt_session_t *session;       //外部broker的MQTT client，用于连接云平台
    dy_db_session_t *db_session;   //用于操作DB的会话
    struct _cloud_mqtt_var_s *var; //上层数据
    cloud_mqtt_hb_login_msg_t heart_beat;
    cloud_mqtt_hb_login_msg_t login;
    topic_maps_t topic_maps;
    COMM_STATUS_E comm_status;
    unsigned int heart_beat_lost_cnt;
    int resume; ///<断点续传，1 开启断点续传  0 关闭断点续传
    link_type_e type;

    char addr[MAX_HOST_LEN]; ///<地址
    unsigned short port;     ///<端口
    char user[MAX_USER_LEN]; ///<用户名
    char pass[MAX_PASS_LEN]; ///<密码地址

    rd_kafka_t *produce_rk;
    rd_kafka_t *consume_rk;
    int consume_topic_cnt;
    rd_kafka_topic_conf_t *consume_topic_conf;
    rd_kafka_topic_partition_list_t *topics;
    int property_cnt;
    property_t *properties;
    struct list_head produce_topic_list; // 已经订阅的topic 列表，重连时需要重新订阅
    struct list_head consume_topic_list; // 已经订阅的topic 列表，重连时需要重新订阅

    pthread_mutex_t produce_lock;
    pthread_cond_t produce_cond;
    int received;
    struct list_head p_ready_list; //待发送的消息列表
    struct list_head p_update_list; //待发送的消息列表
} channel_t;

typedef struct _cloud_mqtt_var_s
{
    ipc_session_t *session;             //内部broker的MQTT client，用于进程间通信
    nodes_cfg_table_t *nodes_cfg_table; //节点信息
    char sn_str[SN_MAX_LEN];
    char proc_name[PROGRAM_NAME_LEN];
    int status_poll_timer;
    int report_timer;
    int channel_cnt;
    channel_t channel[MAX_UPLINK_CHANNEL_CNT]; //上行通道
    int kafka_channel_cnt;
    channel_t kafka_channel[MAX_UPLINK_CHANNEL_CNT]; //上行通道
    char product_name[PRODUCT_NAME_LEN];
} cloud_mqtt_var_t;

typedef enum
{
    FROM_IPC = 0,
    FROM_DB,
} source_e;

typedef struct _prod_msg_node_s
{
    struct list_head list;
    /* data */
    channel_t *channel;
    source_e src;
    int db_id;
    char topic[TOPIC_MAX_LEN];
    int payloadLen;
    unsigned char payload[0];
} p_msg_node_t;

extern int cloud_mqtt_load_cfg(cloud_mqtt_var_t *var, const char *cfg);
extern int cloud_mqtt_internal_mqtt_init(cloud_mqtt_var_t *var);
extern int recv_cloud_mqtt_msg(void *obj, mqtt_message_t *mqtt_msg);
extern int status_poll(cloud_mqtt_var_t *var);
extern int cloud_mqtt_channel_create(cloud_mqtt_var_t *var);
extern int cloud_kafka_channel_create(cloud_mqtt_var_t *var);
extern int kafka_massage_produce(channel_t *channel, p_msg_node_t *msg_node);
extern int replace_heart_beat_interval(const char *cfg_file, unsigned int interval);
extern int get_heart_beat_interval(const char *cfg_file);

#endif /* __CLOUD_MQTT_COMMON_H */
