#include <stdio.h>
#include <stdlib.h>
#include <sys/stat.h>
#include <sys/time.h>
#include <sys/types.h>
#include <unistd.h>
#include <signal.h>

#include "dy_utils/dy_common.h"
#include "cloud_mqtt_common.h"
#include "dy_utils/cJSON.h"
#include "dy_utils/protocol.h"
#include "rdkafka.h"

int recv_cloud_kafka_msg(channel_t *channel, mqtt_message_t *mqtt_msg)
{
	char topic[TOPIC_MAX_LEN] = {0};
	bool matched;
	int j = 0, i = 0;
	int ret;
	cloud_mqtt_var_t *var = channel->var;

	dbg_syslog(LOG_INFO, "External MQTT client received MQTT topic:%s payload:%s", mqtt_msg->topic, mqtt_msg->payload);

	for (j = 0; j < channel->topic_maps.downlink_cnt; j++)
	{
		ret = mosquitto_topic_matches_sub(channel->topic_maps.downlink[j].external_topic, mqtt_msg->topic, &matched);
		if (ret == 0 && matched)
		{

			if (strstr(channel->topic_maps.downlink[j].internal_topic, "[PORT]") != NULL ||
				strstr(channel->topic_maps.downlink[j].internal_topic, "[DEV_SN]") != NULL ||
				strstr(channel->topic_maps.downlink[j].internal_topic, "[SID]") != NULL)
			{
				char port[64] = {0};
				char sn[SN_MAX_LEN] = {0};
				char identifier[64] = {0};
				char topic_2[TOPIC_MAX_LEN] = {0};
				cJSON *node = cJSON_Parse(mqtt_msg->payload);
				if (!node)
				{
					dbg_syslog(LOG_WARNING, "json rsData %s error!!!", mqtt_msg->payload);
					continue;
				}
				GET_JSON_VALUE_STRING(node, "sn", sn);
				GET_JSON_VALUE_STRING(node, "identifier", identifier);
				replace_sub_str(channel->topic_maps.downlink[j].internal_topic, "[DEV_SN]", strlen(sn) == 0 ? "UNKNOWN_SN" : sn, topic);

				GET_JSON_VALUE_STRING(node, "port", port);
				if (strlen(sn) != 0 && strlen(port) == 0)
				{
					// find port according SN
					//遍历设备列表
					for (i = 0; i < var->nodes_cfg_table->node_cnt; i++)
					{
						if (strcmp(sn, var->nodes_cfg_table->node[i].sn) == 0)
						{
							if (strlen(var->nodes_cfg_table->node[i].app_key) != 0)
							{
								strcpy(port, var->nodes_cfg_table->node[i].app_key);
							}
							else
							{
								strcpy(port, port_enum2char(var->nodes_cfg_table->node[i].port));
							}
							break;
						}
					}
				}

				replace_sub_str(topic, "[PORT]", strlen(port) == 0 ? "UNKNOWN_PORT" : port, topic_2);
				replace_sub_str(topic_2, "[SID]", strlen(identifier) == 0 ? "UNKNOWN_SID" : identifier, topic);
				dbg_syslog(LOG_INFO, "forward to internal topic:%s", topic);
				ipc_session_publish(var->session, topic, (unsigned char *) mqtt_msg->payload, mqtt_msg->payloadLen);
				cJSON_Delete(node);
			}
			else
			{
				dbg_syslog(LOG_INFO, "forward to internal topic:%s", channel->topic_maps.downlink[j].internal_topic);
				ipc_session_publish(var->session, channel->topic_maps.downlink[j].internal_topic,
									(unsigned char *) mqtt_msg->payload, mqtt_msg->payloadLen);
			}
		}
	}

	return 0;
}

/*Kafka logger callback (optional)*/
static void logger(const rd_kafka_t *rk, int level, const char *fac, const char *buf)
{
	struct timeval tv;
	gettimeofday(&tv, NULL);
	dbg_syslog(LOG_DEBUG, "%u.%03u RDKAFKA-%i-%s: %s: %s\n", (int)tv.tv_sec, (int)(tv.tv_usec / 1000), level, fac, rk ? rd_kafka_name(rk) : NULL, buf);
}

/**
 * Message delivery report callback.
 * Called once for each message.
 * See rdkafka.h for more information.
 */
static void msg_delivered(rd_kafka_t *rk, const rd_kafka_message_t *rkmessage, void *opaque)
{
	p_msg_node_t *msg_node = (p_msg_node_t *)rkmessage->_private;
	time_t now = time(NULL);
	tag_table_t *tag;
	tag_list_node_t *tag_list = NULL;

	if (rkmessage->err)
	{
		printf("%% Message delivery failed: %s msg_node->topic:%s\n", rd_kafka_err2str(rkmessage->err), msg_node->topic);
		dbg_syslog(LOG_DEBUG, "%% Message delivery failed: DB:%d %s msg_node->topic:%s\n", msg_node->src, rd_kafka_err2str(rkmessage->err), msg_node->topic);
		if (msg_node->src != FROM_DB && msg_node->channel->db_session)
		{
			// 不是从DB来的，失败时需要插入DB
			tag_table_t tag = {0};

			tag.report = NOT_REPORT;
			tag.tag_node = malloc(msg_node->payloadLen + 1);
			memcpy(tag.tag_node, msg_node->payload, msg_node->payloadLen);
			tag.tag_node[msg_node->payloadLen] = '\0';
			tag.tag_len = msg_node->payloadLen;
			tag.time = now;
			tag.ext = strdup(msg_node->topic);
			dy_db_session_insert_tag(msg_node->channel->db_session, &tag);
		}
		free(msg_node);
	}
	else
	{
		printf("%% Message delivered (%zd bytes): msg_node:%s\n", rkmessage->len, msg_node->topic);
		dbg_syslog(LOG_DEBUG, "%% Message delivered DB:%d (%zd bytes): %.*s msg_node:%s\n", msg_node->src, rkmessage->len, (int)rkmessage->len, (const char *)rkmessage->payload, msg_node->topic);

		if (msg_node->src == FROM_DB)
		{
			// 从DB读上来的数据需要更新DB里面的report标记
			tag_list = calloc(1, sizeof(tag_list_node_t));
			if (tag_list != NULL)
			{
				tag = &tag_list->tag;
				tag->id = msg_node->db_id;
				tag->report = REPORTED;
				list_add_tail(&tag_list->list, &msg_node->channel->p_update_list);
			}
		}

		free(msg_node);
	}
}

static void msg_consume(channel_t *channel, rd_kafka_message_t *rkmessage,
						void *opaque)
{
	const char *topic = rd_kafka_topic_name(rkmessage->rkt);
	if (rkmessage->err)
	{
		if (rkmessage->err == RD_KAFKA_RESP_ERR__PARTITION_EOF)
		{
			dbg_syslog(LOG_WARNING, "%% Consumer reached end of %s [%" PRId32 "] message queue at offset %" PRId64 "\n", topic, rkmessage->partition, rkmessage->offset);
			return;
		}

		dbg_syslog(LOG_ERR, "%% Consume error for topic \"%s\" [%" PRId32 "] offset %" PRId64 ": %s\n", topic, rkmessage->partition, rkmessage->offset, rd_kafka_message_errstr(rkmessage));
		return;
	}

	rd_kafka_timestamp_type_t tstype;
	int64_t timestamp;
	dbg_syslog(LOG_DEBUG, "%% Message (offset %" PRId64 ", %zd bytes):\n", rkmessage->offset, rkmessage->len);

	timestamp = rd_kafka_message_timestamp(rkmessage, &tstype);
	if (tstype != RD_KAFKA_TIMESTAMP_NOT_AVAILABLE)
	{
		const char *tsname = "?";
		if (tstype == RD_KAFKA_TIMESTAMP_CREATE_TIME)
			tsname = "create time";
		else if (tstype == RD_KAFKA_TIMESTAMP_LOG_APPEND_TIME)
			tsname = "log append time";

		dbg_syslog(LOG_DEBUG, "%% topic %s Message timestamp: %s %" PRId64 " (%ds ago)\n", topic, tsname, timestamp, !timestamp ? 0 : (int)time(NULL) - (int)(timestamp / 1000));
	}

	if (rkmessage->key_len)
	{
		dbg_syslog_hex(LOG_DEBUG, rkmessage->key, rkmessage->key_len, "topic %s	Message key_len %lu			网关<-KAFKA		",
			topic, (unsigned long) rkmessage->key_len);
	}

	// dbg_syslog_hex(LOG_DEBUG, rkmessage->payload, rkmessage->len, "topic %s	len %d			网关<-KAFKA		", topic, rkmessage->len);
	mqtt_message_t *mqtt_msg = calloc(1, sizeof(mqtt_message_t) + rkmessage->len);
	strncpy(mqtt_msg->topic, topic, sizeof(mqtt_msg->topic) - 1);
	mqtt_msg->payloadLen = rkmessage->len;
	memcpy(mqtt_msg->payload, rkmessage->payload, mqtt_msg->payloadLen);
	recv_cloud_kafka_msg(channel, mqtt_msg);
	free(mqtt_msg);
}

int kafka_massage_produce(channel_t *channel, p_msg_node_t *msg_node)
{
	int ret = -1;

	dbg_syslog(LOG_INFO, "topic:%s payloadLen %d produce_rk:%p", msg_node->topic, msg_node->payloadLen, channel->produce_rk);

	kafka_topic_node_t *node = NULL;
	list_for_each_entry(node, &channel->produce_topic_list, list)
	{
		if (strcmp(node->topic, msg_node->topic) == 0)
		{
			/* Create topic */
			if (node->rkt_topic == NULL)
			{
				/* Topic configuration */
				rd_kafka_topic_conf_t *produce_topic_conf;
				char errstr[448];

				produce_topic_conf = rd_kafka_topic_conf_new();

				if (rd_kafka_topic_conf_set(produce_topic_conf, "message.timeout.ms", "10000", errstr, sizeof(errstr)) != RD_KAFKA_CONF_OK)
				{
					dbg_syslog(LOG_INFO, "%% message.timeout.ms configuration failed: %s\n", errstr);
				}
				if (rd_kafka_topic_conf_set(produce_topic_conf, "request.required.acks", "1", errstr, sizeof(errstr)) != RD_KAFKA_CONF_OK)
				{
					dbg_syslog(LOG_INFO, "%% request.required.acks configuration failed: %s\n", errstr);
				}
				node->rkt_topic = rd_kafka_topic_new(channel->produce_rk, msg_node->topic, produce_topic_conf);
			}

			/* Send/Produce message. */
			ret = rd_kafka_produce(node->rkt_topic, 0, RD_KAFKA_MSG_F_COPY, msg_node->payload, msg_node->payloadLen, NULL, 0, msg_node);
			if (ret == -1)
			{
				dbg_syslog(LOG_ERR, "%% Failed to produce to topic %s %s\n", rd_kafka_topic_name(node->rkt_topic), rd_kafka_err2str(rd_kafka_last_error()));
			}
			else
			{
				ret = 0;
			}

			/* Destroy topic */
			// rd_kafka_topic_destroy(node->rkt_topic);

			break;
		}
	}

	return ret;
}

static void rebalance_cb(rd_kafka_t *rk, rd_kafka_resp_err_t err, rd_kafka_topic_partition_list_t *partitions, void *opaque)
{
	switch (err)
	{
	case RD_KAFKA_RESP_ERR__ASSIGN_PARTITIONS:
		dbg_syslog(LOG_DEBUG, "%% Group rebalanced: %d partition(s) assigned\n", partitions->cnt);
		rd_kafka_assign(rk, partitions);
		break;
	case RD_KAFKA_RESP_ERR__REVOKE_PARTITIONS:
		dbg_syslog(LOG_DEBUG, "%% Group rebalanced: %d partition(s) revoked\n", partitions->cnt);
		rd_kafka_assign(rk, NULL);
		break;
	default:
		break;
	}
}

static void err_cb(rd_kafka_t *rk, int err, const char *reason, void *opaque)
{
	dbg_syslog(LOG_DEBUG, "%% ERROR CALLBACK: %s: %s: %s\n", rd_kafka_name(rk), rd_kafka_err2str(err), reason);
	printf("%% ERROR CALLBACK: %s: %s: %s\n", rd_kafka_name(rk), rd_kafka_err2str(err), reason);
}

static void throttle_cb(rd_kafka_t *rk, const char *broker_name, int32_t broker_id, int throttle_time_ms, void *opaque)
{
	dbg_syslog(LOG_DEBUG, "%% THROTTLED %dms by %s (%" PRId32 ")\n", throttle_time_ms, broker_name, broker_id);
}

static void offset_commit_cb(rd_kafka_t *rk, rd_kafka_resp_err_t err, rd_kafka_topic_partition_list_t *offsets, void *opaque)
{
	int i;

	if (err)
		dbg_syslog(LOG_DEBUG, "%% Offset commit of %d partition(s): %s\n", offsets->cnt, rd_kafka_err2str(err));

	for (i = 0; i < offsets->cnt; i++)
	{
		rd_kafka_topic_partition_t *rktpar = &offsets->elems[i];
		if (rktpar->err)
			dbg_syslog(LOG_DEBUG, "%%  %s [%" PRId32 "] @ %" PRId64 ": %s\n", rktpar->topic, rktpar->partition, rktpar->offset, rd_kafka_err2str(err));
	}
}

static int stats_cb(rd_kafka_t *rk, char *json, size_t json_len, void *opaque)
{
	/* Extract values for our own stats */
	// json_parse_stats(json);
	// dbg_syslog(LOG_DEBUG, "json:%s\n", json);
	return 0;
}

static int _do_produce(channel_t *channel, struct list_head *recv_list)
{
	p_msg_node_t *msg = NULL;
	p_msg_node_t *tmp = NULL;
	int rc = 0;
	time_t now = time(NULL);

	if (list_empty(recv_list))
	{
		return 0;
	}
	list_for_each_entry_safe(msg, tmp, recv_list, list)
	{
		rc = kafka_massage_produce(channel, msg);
		if (rc == -1)
		{
			int last_err = rd_kafka_last_error();
			if (msg->src != FROM_DB && msg->channel->db_session)
			{
				// 不是从DB来的，失败时需要插入DB
				tag_table_t tag = {0};

				tag.report = NOT_REPORT;
				tag.tag_node = malloc(msg->payloadLen + 1);
				memcpy(tag.tag_node, msg->payload, msg->payloadLen);
				tag.tag_node[msg->payloadLen] = '\0';
				tag.tag_len = msg->payloadLen;
				tag.time = now;
				tag.ext = strdup(msg->topic);
				dy_db_session_insert_tag(msg->channel->db_session, &tag);
			}
			free(msg);
			dbg_syslog(LOG_ERR, "produce err:%s", rd_kafka_err2str(last_err));
		}
		else
		{
			rd_kafka_poll(channel->produce_rk, 0);
		}
	}
	return 0;
}

int kafka_produce_in_DB(channel_t *channel)
{
	tag_list_node_t *tmp = NULL;
	tag_list_node_t *tag_list = NULL;
	dy_db_session_t *session = channel->db_session;

	if (!list_empty(&channel->p_update_list)) {
		dy_db_session_add_update_list(session, &channel->p_update_list);
		dy_db_session_update_reported(session);
	}
	dy_db_session_select_not_repoted(session, 100);

	if (list_empty(&session->query_list))
	{
		return 0;
	}
	pthread_mutex_lock(&channel->produce_lock);
	list_for_each_entry_safe(tag_list, tmp, &session->query_list, list)
	{
		tag_table_t *p = &tag_list->tag;

		if (p->ext != NULL && p->tag_node != NULL && p->tag_len > 0)
		{
			p_msg_node_t *node = calloc(1, sizeof(p_msg_node_t) + p->tag_len);
			strcpy(node->topic, p->ext);
			node->payloadLen = p->tag_len;
			node->channel = channel;
			node->src = FROM_DB;
			node->db_id = p->id;
			memcpy(node->payload, p->tag_node, p->tag_len);

			list_add_tail(&node->list, &channel->p_ready_list);
			channel->received++;
		}

		list_del(&tag_list->list);
		free(tag_list->tag.tag_node);
		free(tag_list->tag.ext);
		free(tag_list);
	}
	pthread_mutex_unlock(&channel->produce_lock);
	return 0;
}

void *kafka_produce_loop(void *param)
{
	channel_t *channel = (channel_t *)param;
	while (1)
	{
		struct list_head tmp_recv_list = LIST_HEAD_INIT(tmp_recv_list);

		pthread_mutex_lock(&channel->produce_lock);
		while (channel->received == 0)
		{
			int ret;
			ret = lnxall_condvar_timedwait(&channel->produce_cond, &channel->produce_lock, 3000);
			if (ret == ETIMEDOUT)
			{
				break;
			}
		}

		if (channel->received)
		{
			// real received packets
			if (!list_empty(&channel->p_ready_list))
			{
				list_move_tail_list(&channel->p_ready_list, &tmp_recv_list);
			}
			channel->received = 0;
			pthread_mutex_unlock(&channel->produce_lock);

			_do_produce(channel, &tmp_recv_list);
		}
		else
		{
			pthread_mutex_unlock(&channel->produce_lock);
			rd_kafka_poll(channel->produce_rk, 100);
			if (channel->resume && channel->db_session) {
				kafka_produce_in_DB(channel);
				dy_db_session_delete_reported(channel->db_session);
			}
		}
	}
}

int kafka_produce_init(channel_t *channel, char *brokers, int partition)
{
	rd_kafka_conf_t *conf;
	char errstr[448];
	char tmp[16];
	int i = 0;

	/* Kafka configuration */
	conf = rd_kafka_conf_new();

	/* Set logger */
	rd_kafka_conf_set_log_cb(conf, logger);

	/* Quick termination */
	snprintf(tmp, sizeof(tmp), "%i", SIGIO);
	rd_kafka_conf_set(conf, "internal.termination.signal", tmp, NULL, 0);
	if (strlen(channel->user) && strlen(channel->pass))
	{
		dbg_syslog(LOG_DEBUG, "P== user:(%s %s)", channel->user, channel->pass);
		if (rd_kafka_conf_set(conf, "security.protocol", "SASL_PLAINTEXT", errstr, sizeof(errstr)) ||
			rd_kafka_conf_set(conf, "sasl.mechanism", "SCRAM-SHA-256", errstr, sizeof(errstr)) ||
			rd_kafka_conf_set(conf, "sasl.username", channel->user, errstr, sizeof(errstr)) ||
			rd_kafka_conf_set(conf, "sasl.password", channel->pass, errstr, sizeof(errstr)) ||
			rd_kafka_conf_set(conf, "debug", "security", errstr, sizeof(errstr)))
		{
			dbg_syslog(LOG_ERR, "conf_set user info failed: %s\n", errstr);
			goto __out;
		}
	}

	if (rd_kafka_conf_set(conf, "debug", "broker", errstr, sizeof(errstr)) != RD_KAFKA_CONF_OK)
	{
		dbg_syslog(LOG_DEBUG, "%% Debug configuration failed: %s\n", errstr);
	}

	for (i = 0; i < channel->property_cnt; i++)
	{
		if (rd_kafka_conf_set(conf, channel->properties[i].property, channel->properties[i].value,
							  errstr, sizeof(errstr)) != RD_KAFKA_CONF_OK)
		{
			dbg_syslog(LOG_ERR, "conf_set %s failed: %s\n", channel->properties[i].property, errstr);
		}
		else
		{
			dbg_syslog(LOG_DEBUG, "conf_set %s=%s", channel->properties[i].property, channel->properties[i].value);
		}
	}

	rd_kafka_conf_set_dr_msg_cb(conf, msg_delivered);
	/* Create Kafka handle */
	if (!(channel->produce_rk = rd_kafka_new(RD_KAFKA_PRODUCER, conf, errstr, sizeof(errstr))))
	{
		dbg_syslog(LOG_ERR, "%% Failed to create new producer: %s\n", errstr);
		goto __out;
	}

	/* Add brokers */
	if (rd_kafka_brokers_add(channel->produce_rk, brokers) == 0)
	{
		dbg_syslog(LOG_ERR, "%% No valid brokers specified\n");
		goto __out;
	}

	rd_kafka_dump(stdout, channel->produce_rk);

	pthread_t thread_produce;
	pthread_create(&thread_produce, NULL, kafka_produce_loop, (void *)channel);

	return 0;

__out:
	channel->topic_maps.uplink_cnt = 0;
	return 1;
}

int kafka_consume_init(channel_t *channel, char *brokers, int partition)
{
	char tmp[128];
	char errstr[448];
	rd_kafka_conf_t *conf;
	int i;

	/* Kafka configuration */
	conf = rd_kafka_conf_new();
	rd_kafka_conf_set_error_cb(conf, err_cb);
	rd_kafka_conf_set_throttle_cb(conf, throttle_cb);
	rd_kafka_conf_set_offset_commit_cb(conf, offset_commit_cb);
	/* Quick termination */
	snprintf(tmp, sizeof(tmp), "%i", SIGIO);
	rd_kafka_conf_set(conf, "internal.termination.signal", tmp, NULL, 0);

	/* config */
	rd_kafka_conf_set(conf, "queue.buffering.max.messages", "500000", NULL, 0);
	rd_kafka_conf_set(conf, "message.send.max.retries", "3", NULL, 0);
	rd_kafka_conf_set(conf, "retry.backoff.ms", "500", NULL, 0);
	rd_kafka_conf_set(conf, "queued.min.messages", "1000000", NULL, 0);
	rd_kafka_conf_set(conf, "session.timeout.ms", "6000", NULL, 0);

	/* Kafka topic configuration */
	channel->consume_topic_conf = rd_kafka_topic_conf_new();
	rd_kafka_topic_conf_set(channel->consume_topic_conf, "auto.offset.reset", "earliest", NULL, 0);
	if (strlen(channel->user) && strlen(channel->pass))
	{
		dbg_syslog(LOG_DEBUG, "C== user:(%s %s)", channel->user, channel->pass);
		if (rd_kafka_conf_set(conf, "security.protocol", "SASL_PLAINTEXT", errstr, sizeof(errstr)) ||
			rd_kafka_conf_set(conf, "sasl.mechanism", "SCRAM-SHA-256", errstr, sizeof(errstr)) ||
			rd_kafka_conf_set(conf, "sasl.username", channel->user, errstr, sizeof(errstr)) ||
			rd_kafka_conf_set(conf, "sasl.password", channel->pass, errstr, sizeof(errstr)) ||
			rd_kafka_conf_set(conf, "debug", "security", errstr, sizeof(errstr)))
		{
			dbg_syslog(LOG_ERR, "conf_set user info failed: %s\n", errstr);
			goto __out;
		}
	}

	for (i = 0; i < channel->property_cnt; i++)
	{
		if (rd_kafka_conf_set(conf, channel->properties[i].property, channel->properties[i].value,
							  errstr, sizeof(errstr)) != RD_KAFKA_CONF_OK)
		{
			dbg_syslog(LOG_ERR, "conf_set %s failed: %s\n", channel->properties[i].property, errstr);
		}
		else
		{
			dbg_syslog(LOG_DEBUG, "conf_set %s=%s", channel->properties[i].property, channel->properties[i].value);
		}
	}

	channel->topics = rd_kafka_topic_partition_list_new(channel->consume_topic_cnt);
	topic_node_t *node = NULL;
	list_for_each_entry(node, &channel->consume_topic_list, list)
	{
		dbg_syslog(LOG_INFO, "consume_topic_cnt %d node->topic %s", channel->consume_topic_cnt, node->topic);
		if (rd_kafka_conf_set(conf, "group.id", node->topic, errstr, sizeof(errstr)) != RD_KAFKA_CONF_OK)
		{
			dbg_syslog(LOG_ERR, "%% conf_set topic group.id %s\n", errstr);
			goto __out;
		}
		rd_kafka_topic_partition_list_add(channel->topics, node->topic, partition);
	}

	/* Always enable stats (for RTT extraction), and if user supplied
	 * the -T <intvl> option we let her take part of the stats aswell. */
	rd_kafka_conf_set_stats_cb(conf, stats_cb);

	/*
	 * High-level balanced Consumer
	 */
	rd_kafka_resp_err_t err;

	rd_kafka_conf_set_rebalance_cb(conf, rebalance_cb);
	rd_kafka_conf_set_default_topic_conf(conf, channel->consume_topic_conf);
	/* Create Kafka handle */
	if (!(channel->consume_rk = rd_kafka_new(RD_KAFKA_CONSUMER, conf, errstr, sizeof(errstr))))
	{
		dbg_syslog(LOG_DEBUG, "%% Failed to create Kafka consumer: %s\n", errstr);
	}

	/* Forward all events to consumer queue */
	rd_kafka_poll_set_consumer(channel->consume_rk);
	/* Add broker(s) */
	if (brokers && rd_kafka_brokers_add(channel->consume_rk, brokers) < 1)
	{
		dbg_syslog(LOG_DEBUG, "%% No valid brokers specified\n");
		goto __out;
	}

	err = rd_kafka_subscribe(channel->consume_rk, channel->topics);
	if (err)
	{
		dbg_syslog(LOG_DEBUG, "%% Subscribe failed: %s\n", rd_kafka_err2str(err));
	}

	dbg_syslog(LOG_DEBUG, "%% Waiting for group rebalance..\n");
	return 0;
__out:
	channel->topic_maps.downlink_cnt = 0;
	return 1;
}

void *kafka_consume_loop(void *param)
{
	channel_t *channel = (channel_t *)param;
	rd_kafka_resp_err_t err;

	while (1)
	{
		rd_kafka_message_t *rkmessage;

		rkmessage = rd_kafka_consumer_poll(channel->consume_rk, 1000);
		if (rkmessage)
		{
			msg_consume(channel, rkmessage, NULL);
			rd_kafka_message_destroy(rkmessage);
		}
		// print_stats(rk, mode, otype, compression);
	}

	err = rd_kafka_consumer_close(channel->consume_rk);
	if (err)
		dbg_syslog(LOG_DEBUG, "%% Failed to close consumer: %s\n", rd_kafka_err2str(err));

	rd_kafka_destroy(channel->consume_rk);

	return 0;
}

int kafka_session_subscribe(channel_t *channel, char *topic)
{
	size_t tlen;
	topic_node_t *node = NULL;
	list_for_each_entry(node, &channel->consume_topic_list, list)
	{
		if (strcmp(node->topic, topic) == 0)
		{
			dbg_syslog(LOG_DEBUG, "find node->topic:%s", node->topic);
			return 0;
		}
	}

	tlen = strlen(topic);
	node = malloc(sizeof(topic_node_t) + tlen + 1);
	if (node == NULL)
	{
		dbg_syslog(LOG_ERR, "malloc failed, strlen(topic):%d", (int) tlen);
		return -1;
	}

	strcpy(node->topic, topic);
	dbg_syslog(LOG_DEBUG, "node->topic:%s", node->topic);
	channel->consume_topic_cnt++;
	list_add_tail(&node->list, &channel->consume_topic_list);
	return 0;
}

int kafka_add_session_produce(channel_t *channel, char *topic)
{
	size_t tlen;
	kafka_topic_node_t *node = NULL;
	list_for_each_entry(node, &channel->produce_topic_list, list)
	{
		if (strcmp(node->topic, topic) == 0)
		{
			dbg_syslog(LOG_DEBUG, "produce find node->topic:%s", node->topic);
			return 0;
		}
	}

	tlen = strlen(topic);
	node = calloc(sizeof(kafka_topic_node_t) + tlen + 1, 1);
	if (node == NULL)
	{
		dbg_syslog(LOG_ERR, "malloc failed, strlen(topic):%d", (int) tlen);
		return -1;
	}

	strcpy(node->topic, topic);
	dbg_syslog(LOG_INFO, "produce node->topic:%s", node->topic);
	channel->consume_topic_cnt++;
	list_add_tail(&node->list, &channel->produce_topic_list);
	return 0;
}

// 建立与broker之间的MQTT连接
int cloud_kafka_channel_create(cloud_mqtt_var_t *var)
{
	int i = 0, j = 0;
	const char *board_name = NULL;
	board_name = get_board_name();

	if (var->kafka_channel_cnt <= 0)
	{
		dbg_syslog(LOG_DEBUG, "kafka_channel_cnt:%d", var->kafka_channel_cnt);
		return 0;
	}

	for (i = 0; i < var->kafka_channel_cnt; i++)
	{
		channel_t *channel = &var->kafka_channel[i];
		char db_name[32] = {0};

		dbg_syslog(LOG_DEBUG, "channel:%p board_name:%s", channel, board_name);
		dbg_syslog(LOG_DEBUG, "resume:%d", channel->resume);

		channel->var = var;
		INIT_LIST_HEAD(&channel->consume_topic_list);
		INIT_LIST_HEAD(&channel->produce_topic_list);
		INIT_LIST_HEAD(&channel->p_ready_list);
		INIT_LIST_HEAD(&channel->p_update_list);
		pthread_mutex_init(&channel->produce_lock, NULL);
		lnxall_condvar_init(&channel->produce_cond);

		if (channel->resume)
		{
			dbg_syslog(LOG_DEBUG, "resume:%d", channel->resume);
			if (strcmp(board_name, BOARD_WOOLINK_MT7628) == 0)
			{
				sprintf(db_name, "/tmp/KAFKA_REPORT_%d.db", i);
			}
			else
			{
				sprintf(db_name, "/app/KAFKA_REPORT_%d.db", i);
			}
			dbg_syslog(LOG_DEBUG, "db_name:%s", db_name);
			dy_db_session_init(&channel->db_session, db_name, (void *)channel);
		}
		dbg_syslog(LOG_DEBUG, "resume:%d", channel->resume);

		char broker[64] = {0};

		snprintf(broker, 64, "%s:%d", channel->addr, channel->port);
		dbg_syslog(LOG_DEBUG, "broker:%s", broker);

		if (channel->topic_maps.uplink_cnt > 0)
		{
			kafka_produce_init(channel, broker, 0);
		}

		if (channel->topic_maps.downlink_cnt > 0)
		{
			if (kafka_consume_init(channel, broker, 0) == 0)
			{
				pthread_t thread_consum;
				pthread_create(&thread_consum, NULL, kafka_consume_loop, (void *)channel);
			}
		}
		for (j = 0; j < channel->topic_maps.downlink_cnt; j++)
		{
			kafka_session_subscribe(channel, channel->topic_maps.downlink[j].external_topic);
		}

		for (j = 0; j < channel->topic_maps.uplink_cnt; j++)
		{
			kafka_add_session_produce(channel, channel->topic_maps.uplink[j].external_topic);
		}
	}
	return 0;
}
