#ifndef __CLOUD_MQTT_COMMON_H
#define __CLOUD_MQTT_COMMON_H

/***************HEADER FILES****************/
#include <string.h>
#include <syslog.h>
#include "mqtt_session.h"
#include "ipc_session.h"
#include "lnxall_list.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"

#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"

#define CLOUD_NANDU_CFG "/app/config/nandu/uplink_channels.json"

#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 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 关闭断点续传

    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 *produce_topic_conf;
    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 列表，重连时需要重新订阅
} 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]; //上行通道
} cloud_mqtt_var_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 cloud_kafka_channel_create(cloud_mqtt_var_t *var);
extern int kafka_message_produce(channel_t *channel, char *topic, char *payload, int payloadLen);

#endif /* __CLOUD_MQTT_COMMON_H */
