369 lines
6.5 KiB
C
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;
|
|
} |