#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 #include /* * 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; }