Files

369 lines
6.5 KiB
C

#include "mqtt_manager.h"
#include "mqtt_config.h"
#include "mqtt_protocol.h"
#include "wifi_manager.h"
#include "app_log.h"
#include "mqtt/MQTTClient.h"
#include "os/os_api.h"
#include <string.h>
#include <stdio.h>
/*
* MQTT对象
*
* 按AC79官方demo改成静态生命周期
*/
static Client mqtt_client;
static Network mqtt_network;
/*
* MQTT状态
*/
static mqtt_state_t mqtt_state = MQTT_STATE_INIT;
static u32 mqtt_reconnect_delay = 0;
/*
* MQTT消息回调
*/
static mqtt_message_callback_t message_callback = NULL;
/*
* MQTT缓存
*
* AC79 Paho需要
*/
static unsigned char mqtt_send_buffer[2048];
static unsigned char mqtt_read_buffer[2048];
/*
* 收到MQTT消息
*
* 对应ESP32:
*
* mqtt event callback
*/
static void message_arrived
(
MessageData *data
)
{
/*
* 先交给协议层
*
* 已识别的小智消息由协议层处理
*/
if(
mqtt_protocol_handle_message(
data->topicName->lenstring.data,
(uint8_t *)data->message->payload,
data->message->payloadlen
) == 0
)
{
return;
}
if(message_callback == NULL)
{
return;
}
APP_LOG("mqtt message arrived\n");
APP_LOG(
"topic:%s\n",
data->topicName->lenstring.data
);
APP_LOG(
"payload:%.*s\n",
data->message->payloadlen,
data->message->payload
);
mqtt_message_callback_t cb =
message_callback;
cb(
data->topicName->lenstring.data,
(uint8_t *)data->message->payload,
data->message->payloadlen
);
}
/*
* MQTT连接
*/
static int mqtt_connect_server(void)
{
mqtt_config_t *cfg = mqtt_config_get();
MQTTPacket_connectData connect_data = MQTTPacket_connectData_initializer;
/*
* 创建Network
*/
NewNetwork(&mqtt_network);
SetNetworkRecvTimeout(
&mqtt_network,
1000
);
/*
* TCP连接
*/
if(ConnectNetwork(&mqtt_network,cfg->host,cfg->port))
{
APP_LOG("mqtt tcp connect fail\n");
mqtt_state =
MQTT_STATE_ERROR;
return -1;
}
APP_LOG("mqtt tcp connect ok\n");
/*
* 创建MQTT Client
*/
MQTTClient(
&mqtt_client, // 客户端对象-->我的设备
&mqtt_network, // 已经建立的 TCP 网络
1000, // 命令超时 ms
mqtt_send_buffer, // 发送缓冲区
sizeof(mqtt_send_buffer),
mqtt_read_buffer, // 接收缓冲区
sizeof(mqtt_read_buffer)
);
/*
* MQTT身份
*
* 暂时使用config中的固定值
*/
connect_data.clientID.cstring = cfg->client_id;
if (cfg->username[0])
{
connect_data.username.cstring =
cfg->username;
}
if (cfg->password[0])
{
connect_data.password.cstring =
cfg->password;
}
connect_data.keepAliveInterval =
cfg->keepalive;
/*
* 建立MQTT连接
*/
int rc = MQTTConnect( &mqtt_client, &connect_data);
APP_LOG("MQTTConnect rc=%d\n", rc);
if(rc)
{
APP_LOG("mqtt connect failed\n");
mqtt_network.disconnect(
&mqtt_network
);
mqtt_state = MQTT_STATE_ERROR;
return -1;
}
APP_LOG("mqtt connect ok\n");
mqtt_state = MQTT_STATE_CONNECTED;
return 0;
}
/*
* MQTT任务
*/
static void mqtt_task(void *arg)
{
while(1)
{
if(
!wifi_manager_is_network_ready()
)
{
os_time_dly(200);
continue;
}
if(
mqtt_state != MQTT_STATE_CONNECTED
)
{
mqtt_state = MQTT_STATE_CONNECTING;
if(mqtt_connect_server())
{
if (mqtt_reconnect_delay < 16000)
{
mqtt_reconnect_delay =
mqtt_reconnect_delay ?
mqtt_reconnect_delay * 2 : 500;
}
os_time_dly(mqtt_reconnect_delay);
continue;
}
mqtt_reconnect_delay = 0;
/*
* 连接成功后订阅
*/
if(MQTTSubscribe
(
&mqtt_client,
mqtt_config_get()->subscribe_topic,
QOS0,
message_arrived
)
)
{
APP_LOG("mqtt subscribe failed\n");
MQTTDisconnect(
&mqtt_client
);
mqtt_network.disconnect(
&mqtt_network
);
mqtt_state =
MQTT_STATE_ERROR;
continue;
}
APP_LOG("mqtt subscribe ok\n");
char online_payload[256];
snprintf(
online_payload,
sizeof(online_payload),
"{\"id\":1,\"method\":\"thing.event.property.post\","
"\"params\":{\"deviceId\":\"%s\",\"reportType\":\"heartbeat\","
"\"personDetected\":0,\"heartbeat\":1}}",
mqtt_config_get()->username
);
int pub_ret = mqtt_manager_publish(
mqtt_config_get()->publish_topic,
(const uint8_t *)online_payload,
(uint32_t)strlen(online_payload)
);
APP_LOG("mqtt online publish ret=%d\n", pub_ret);
}
/*
* 保持MQTT通信
*
* 收消息依靠yield
*/
MQTTYield
(
&mqtt_client,
1000
);
os_time_dly(10);
}
}
int mqtt_manager_init(void)
{
mqtt_config_init();
mqtt_state = MQTT_STATE_INIT;
return 0;
}
int mqtt_manager_start(void)
{
return thread_fork(
"mqtt_task",
10,
2048,
0,
NULL,
mqtt_task,
NULL
);
}
int mqtt_manager_publish(const char *topic,const uint8_t *data,uint32_t len)
{
MQTTMessage message;
memset(
&message,
0,
sizeof(message)
);
message.qos = QOS0;
message.payload = (void *)data;
message.payloadlen = len;
return MQTTPublish
(
&mqtt_client,
topic,
&message
);
}
void mqtt_manager_set_callback(mqtt_message_callback_t callback)
{
message_callback = callback;
}
mqtt_state_t mqtt_manager_get_state(void)
{
return mqtt_state;
}