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

#include "mqtt_session.h"
#include "dy_utils/dy_common.h"
#include "cloud_nandu_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};
	unsigned char matched;
	int j = 0, i = 0;
	int ret;
	cloud_mqtt_var_t *var = channel->var;

	dy_syslog(LOG_DEBUG, "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 = mqtt_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)
				{
					dy_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);
				dy_syslog(LOG_DEBUG, "forward to internal topic:%s", topic);
				ipc_session_publish(var->session, topic, mqtt_msg->payload, mqtt_msg->payloadLen);
				cJSON_Delete(node);
			}
			else
			{
				dy_syslog(LOG_DEBUG, "forward to internal topic:%s", channel->topic_maps.downlink[j].internal_topic);
				ipc_session_publish(var->session, channel->topic_maps.downlink[j].internal_topic,
										mqtt_msg->payload, mqtt_msg->payloadLen);
			}
		}
	}

out:
	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);
	dy_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, void *payload, size_t len, int error_code, void *opaque, void *msg_opaque)
{
	if (error_code)
		dy_syslog(LOG_WARNING, "%% Message delivery failed: %s\n", rd_kafka_err2str(error_code));
	else
		dy_syslog(LOG_DEBUG, "%% Message delivered (%zd bytes): %.*s\n", len, (int)len, (const char *)payload);
}

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)
		{
			dy_syslog(LOG_WARNING, "%% Consumer reached end of %s [%" PRId32 "] message queue at offset %" PRId64 "\n", topic, rkmessage->partition, rkmessage->offset);
			return;
		}

		dy_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;
	dy_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";

		dy_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)
	{
		dy_syslog_hex(LOG_DEBUG, rkmessage->key, rkmessage->key_len, "topic %s	Message key_len %d			网关<-KAFKA		", topic, rkmessage->key_len);
	}

	//dy_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));
	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 nandu_message_produce_handle(channel_t *channel, char *topic, char *payload, int payloadLen)
{
	dy_syslog(LOG_DEBUG, "topic:%s payloadLen %d payload %s produce_rk:%p produce_topic_conf %p", topic, payloadLen, payload, channel->produce_rk, channel->produce_topic_conf);

	kafka_topic_node_t *node = NULL;
	list_for_each_entry(node, &channel->produce_topic_list, list)
	{
		if (strcmp(node->topic, topic) == 0)
		{
			rd_kafka_resp_err_t err = RD_KAFKA_RESP_ERR_NO_ERROR;
			int retry_cnt = 0;
			/* Create topic */
			if (node->rkt_topic == NULL)
				node->rkt_topic = rd_kafka_topic_new(channel->produce_rk, topic, channel->produce_topic_conf);

		retry:
			/* Send/Produce message. */
			if (rd_kafka_produce(node->rkt_topic, RD_KAFKA_PARTITION_UA, RD_KAFKA_MSG_F_COPY, payload, payloadLen, NULL, 0, NULL) == -1)
			{
				err = rd_kafka_last_error();
				if (err == RD_KAFKA_RESP_ERR__QUEUE_FULL)
				{
					/* If the internal queue is full, wait for
					* messages to be delivered and then retry.
					* The internal queue represents both
					* messages to be sent and messages that have
					* been sent or failed, awaiting their
					* delivery report callback to be called.
					*
					* The internal queue is limited by the
					* configuration property
					* queue.buffering.max.messages */
					rd_kafka_poll(channel->produce_rk, 1000 /*block for max 1000ms*/);
					if (retry_cnt == 0)
					{
						retry_cnt++;
						goto retry;
					}
				}
				dy_syslog(LOG_ERR, "%% Failed to produce to topic %s %s\n", rd_kafka_topic_name(node->rkt_topic), rd_kafka_err2str(err));
			}
			/* Poll to handle delivery reports */
			rd_kafka_poll(channel->produce_rk, 0);

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

			break;
		}
	}
}

int kafka_message_produce(channel_t *channel, char *topic, char *payload, int payloadLen)
{
	dy_syslog(LOG_DEBUG, "topic:%s payloadLen %d payload %s produce_rk:%p produce_topic_conf %p", topic, payloadLen, payload, channel->produce_rk, channel->produce_topic_conf);
	int ret = 0;
	char nandu_topic[64] = {0};
	char *nandu_payload = calloc(payloadLen * 2, 1);
	cJSON *root = cJSON_Parse(payload);
	if (root)
	{
		char time_str[16] = {0};
		char sn[64] = {0};
		time_t time = 0;

		GET_JSON_VALUE_STRING(root, "sn", sn);
		GET_JSON_VALUE_STRING(root, "siteId", nandu_topic);
		if (strlen(nandu_topic) == 0)
		{
			dy_syslog(LOG_WARNING, "payload con't find siteId!!!, use default topic");
			strcpy(nandu_topic, topic);
		}
		GET_JSON_VALUE_INT(root, "time", time);
		struct tm *timenow = localtime(&time);
		sprintf(time_str, "%04d%02d%02d%02d%02d000", (timenow->tm_year + 1900), timenow->tm_mon + 1, timenow->tm_mday, timenow->tm_hour, timenow->tm_min);
		strcat(nandu_payload, time_str);
		cJSON *tags = cJSON_GetObjectItem(root, "tags");
		if (tags)
		{
			cJSON *child = NULL;
			for (child = tags->child; child; child = child->next)
			{
				if (child->string)
				{
					char tag[64] = {0};
					char tagv[64] = {0};
					char *tag_s = strstr(child->string, "TAG_");
					if (tag_s)
					{
						char sn_pref[64] = {0};
						sscanf(sn,"%[0-9a-fA-F]",sn_pref);
						sprintf(tag, "|%s%05X", sn_pref, strtol(tag_s + 4, NULL, 16));
						if (child->valuestring)
						{
							sprintf(tagv, "|%s", child->valuestring);
						}
						else
						{
							sprintf(tagv, "|%1.15g", child->valuedouble);
						}
						strcat(nandu_payload, tag);
						strcat(nandu_payload, tagv);
					}
				}
			}
		}
		cJSON_Delete(root);
	}

	nandu_message_produce_handle(channel, nandu_topic, nandu_payload, strlen(nandu_payload));
_out:
	if (nandu_payload)
		free(nandu_payload);
	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:
		dy_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:
		dy_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)
{
	dy_syslog(LOG_DEBUG, "%% 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)
{
	dy_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)
		dy_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)
			dy_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);
	//dy_syslog(LOG_DEBUG, "json:%s\n", json);
	return 0;
}

int kafka_produce_init(channel_t *channel, char *brokers, int partition)
{
	rd_kafka_topic_t *rkt;
	rd_kafka_conf_t *conf;
	char errstr[512];
	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))
	{
		dy_syslog(LOG_DEBUG, "== 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)))
		{
			dy_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)
		{
			dy_syslog(LOG_ERR, "conf_set %s failed: %s\n", channel->properties[i].property, errstr);
		}
		else
		{
			dy_syslog(LOG_DEBUG, "conf_set %s=%s", channel->properties[i].property, channel->properties[i].value);
		}
	}

	/* Topic configuration */
	channel->produce_topic_conf = rd_kafka_topic_conf_new();
	if (rd_kafka_conf_set(conf, "debug", "broker", errstr, sizeof(errstr)) != RD_KAFKA_CONF_OK)
	{
		dy_syslog(LOG_DEBUG, "%% Debug configuration failed: %s\n", errstr);
	}

	rd_kafka_conf_set_dr_cb(conf, msg_delivered);

	/* Create Kafka handle */
	if (!(channel->produce_rk = rd_kafka_new(RD_KAFKA_PRODUCER, conf, errstr, sizeof(errstr))))
	{
		dy_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)
	{
		dy_syslog(LOG_ERR, "%% No valid brokers specified\n");
		goto __out;
	}

	rd_kafka_dump(stdout, channel->produce_rk);

	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[512];
	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))
	{
		dy_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)))
		{
			dy_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)
		{
			dy_syslog(LOG_ERR, "conf_set %s failed: %s\n", channel->properties[i].property, errstr);
		}
		else
		{
			dy_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)
	{
		dy_syslog(LOG_DEBUG, "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)
		{
			dy_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))))
	{
		dy_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)
	{
		dy_syslog(LOG_DEBUG, "%% No valid brokers specified\n");
		goto __out;
	}

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

	dy_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;
		uint64_t fetch_latency;

		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)
		dy_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)
{
	topic_node_t *node = NULL;
	list_for_each_entry(node, &channel->consume_topic_list, list)
	{
		if (strcmp(node->topic, topic) == 0)
		{
			dy_syslog(LOG_DEBUG, "find node->topic:%s", node->topic);
			return 0;
		}
	}

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

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

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

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

	strcpy(node->topic, topic);
	dy_syslog(LOG_DEBUG, "produce node->topic:%s", node->topic);
	channel->consume_topic_cnt++;
	list_add_tail(&node->list, &channel->produce_topic_list);
}

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

	if (var->kafka_channel_cnt <= 0)
	{
		dy_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};

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

		channel->var = var;
		INIT_LIST_HEAD(&channel->consume_topic_list);
		INIT_LIST_HEAD(&channel->produce_topic_list);

		if (channel->resume)
		{
			dy_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);
			}
			dy_syslog(LOG_DEBUG, "db_name:%s", db_name);
			dy_db_session_init(&channel->db_session, db_name, (void *)channel);
		}
		dy_syslog(LOG_DEBUG, "resume:%d", channel->resume);

		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);
		}

		char broker[64] = {0};

		snprintf(broker, 64, "%s:%d", channel->addr, channel->port);
		dy_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);
			}
		}
	}
}
