2016-09-11 10:10:24 +07:00
|
|
|
/*
|
|
|
|
|
* @Author: Tuan PM
|
|
|
|
|
* @Date: 2016-09-10 09:33:06
|
|
|
|
|
* @Last Modified by: Tuan PM
|
2016-09-12 17:05:28 +07:00
|
|
|
* @Last Modified time: 2016-09-12 17:03:56
|
2016-09-11 10:10:24 +07:00
|
|
|
*/
|
2016-09-12 17:05:28 +07:00
|
|
|
#include <stdio.h>
|
2016-09-11 21:42:34 +07:00
|
|
|
#include "freertos/FreeRTOS.h"
|
|
|
|
|
#include "freertos/task.h"
|
|
|
|
|
#include "freertos/semphr.h"
|
|
|
|
|
#include "freertos/queue.h"
|
2016-09-11 10:10:24 +07:00
|
|
|
|
2016-09-11 21:42:34 +07:00
|
|
|
#include "lwip/sockets.h"
|
|
|
|
|
#include "lwip/dns.h"
|
|
|
|
|
#include "lwip/netdb.h"
|
2016-09-12 17:05:28 +07:00
|
|
|
#include "ringbuf.h"
|
|
|
|
|
#include "mqtt.h"
|
2016-09-11 21:42:34 +07:00
|
|
|
|
|
|
|
|
static TaskHandle_t xMqttTask = NULL;
|
2016-09-12 17:05:28 +07:00
|
|
|
static TaskHandle_t xMqttSendingTask = NULL;
|
2016-09-11 21:42:34 +07:00
|
|
|
|
|
|
|
|
|
|
|
|
|
static int resolev_dns(const char *host, struct sockaddr_in *ip) {
|
|
|
|
|
struct hostent *he;
|
|
|
|
|
struct in_addr **addr_list;
|
|
|
|
|
he = gethostbyname(host);
|
|
|
|
|
if (he == NULL) return 0;
|
|
|
|
|
addr_list = (struct in_addr **)he->h_addr_list;
|
|
|
|
|
if (addr_list[0] == NULL) return 0;
|
|
|
|
|
ip->sin_family = AF_INET;
|
|
|
|
|
memcpy(&ip->sin_addr, addr_list[0], sizeof(ip->sin_addr));
|
|
|
|
|
return 1;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
static int client_connect(const char *stream_host, int stream_port)
|
|
|
|
|
{
|
2016-09-12 13:06:48 +07:00
|
|
|
int sock;
|
|
|
|
|
struct sockaddr_in remote_ip;
|
2016-09-11 21:42:34 +07:00
|
|
|
while (1) {
|
|
|
|
|
bzero(&remote_ip, sizeof(struct sockaddr_in));
|
2016-09-12 13:06:48 +07:00
|
|
|
remote_ip.sin_family = AF_INET;
|
2016-09-11 21:42:34 +07:00
|
|
|
//if stream_host is not ip address, resolve it
|
2016-09-12 13:06:48 +07:00
|
|
|
if (inet_aton(stream_host, &(remote_ip.sin_addr)) == 0) {
|
|
|
|
|
mqtt_info("Resolve dns for domain: %s", stream_host);
|
2016-09-11 21:42:34 +07:00
|
|
|
if (!resolev_dns(stream_host, &remote_ip)) {
|
|
|
|
|
vTaskDelay(1000 / portTICK_RATE_MS);
|
|
|
|
|
continue;
|
|
|
|
|
}
|
|
|
|
|
}
|
2016-09-12 13:06:48 +07:00
|
|
|
sock = socket(PF_INET, SOCK_STREAM, 0);
|
2016-09-11 21:42:34 +07:00
|
|
|
if (sock == -1) {
|
|
|
|
|
continue;
|
|
|
|
|
}
|
|
|
|
|
remote_ip.sin_port = htons(stream_port);
|
2016-09-12 13:06:48 +07:00
|
|
|
mqtt_info("Connecting to server %s:%d,%d",
|
|
|
|
|
inet_ntoa((remote_ip.sin_addr)),
|
|
|
|
|
stream_port,
|
|
|
|
|
remote_ip.sin_port);
|
|
|
|
|
|
2016-09-11 21:42:34 +07:00
|
|
|
if (connect(sock, (struct sockaddr *)(&remote_ip), sizeof(struct sockaddr)) != 00) {
|
|
|
|
|
close(sock);
|
2016-09-12 13:06:48 +07:00
|
|
|
mqtt_error("Conn err.");
|
2016-09-11 21:42:34 +07:00
|
|
|
vTaskDelay(1000 / portTICK_RATE_MS);
|
|
|
|
|
continue;
|
|
|
|
|
}
|
|
|
|
|
return sock;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
/*
|
|
|
|
|
* mqtt_connect
|
|
|
|
|
* input - client
|
|
|
|
|
* return 1: success, 0: fail
|
|
|
|
|
*/
|
|
|
|
|
static bool mqtt_connect(mqtt_client *client)
|
|
|
|
|
{
|
|
|
|
|
int write_len, read_len, connect_rsp_code;
|
2016-09-12 13:06:48 +07:00
|
|
|
struct timeval tv;
|
|
|
|
|
|
|
|
|
|
tv.tv_sec = 10; /* 30 Secs Timeout */
|
|
|
|
|
tv.tv_usec = 0; // Not init'ing this can cause strange errors
|
|
|
|
|
|
|
|
|
|
setsockopt(client->socket, SOL_SOCKET, SO_RCVTIMEO, (char *)&tv, sizeof(struct timeval));
|
|
|
|
|
|
2016-09-11 21:42:34 +07:00
|
|
|
mqtt_msg_init(&client->mqtt_state.mqtt_connection,
|
|
|
|
|
client->mqtt_state.out_buffer,
|
|
|
|
|
client->mqtt_state.out_buffer_length);
|
|
|
|
|
client->mqtt_state.outbound_message = mqtt_msg_connect(&client->mqtt_state.mqtt_connection,
|
|
|
|
|
client->mqtt_state.connect_info);
|
|
|
|
|
client->mqtt_state.pending_msg_type = mqtt_get_type(client->mqtt_state.outbound_message->data);
|
|
|
|
|
client->mqtt_state.pending_msg_id = mqtt_get_id(client->mqtt_state.outbound_message->data,
|
|
|
|
|
client->mqtt_state.outbound_message->length);
|
2016-09-12 13:06:48 +07:00
|
|
|
mqtt_info("Sending MQTT CONNECT message, type: %d, id: %04X",
|
2016-09-11 21:42:34 +07:00
|
|
|
client->mqtt_state.pending_msg_type,
|
|
|
|
|
client->mqtt_state.pending_msg_id);
|
|
|
|
|
write_len = write(client->socket,
|
|
|
|
|
client->mqtt_state.outbound_message->data,
|
|
|
|
|
client->mqtt_state.outbound_message->length);
|
2016-09-12 13:06:48 +07:00
|
|
|
mqtt_info("Reading MQTT CONNECT response message");
|
2016-09-11 21:42:34 +07:00
|
|
|
read_len = read(client->socket, client->mqtt_state.in_buffer, CONFIG_MQTT_BUFFER_SIZE_BYTE);
|
2016-09-12 17:05:28 +07:00
|
|
|
|
|
|
|
|
tv.tv_sec = 0; /* No timeout */
|
|
|
|
|
setsockopt(client->socket, SOL_SOCKET, SO_RCVTIMEO, (char *)&tv, sizeof(struct timeval));
|
|
|
|
|
|
2016-09-12 13:06:48 +07:00
|
|
|
if (read_len < 0) {
|
|
|
|
|
mqtt_error("Error network response");
|
2016-09-11 21:42:34 +07:00
|
|
|
return false;
|
|
|
|
|
}
|
|
|
|
|
if (mqtt_get_type(client->mqtt_state.in_buffer) != MQTT_MSG_TYPE_CONNACK) {
|
2016-09-12 13:06:48 +07:00
|
|
|
mqtt_error("Invalid MSG_TYPE response: %d, read_len: %d", mqtt_get_type(client->mqtt_state.in_buffer), read_len);
|
2016-09-11 21:42:34 +07:00
|
|
|
return false;
|
|
|
|
|
}
|
|
|
|
|
connect_rsp_code = mqtt_get_connect_return_code(client->mqtt_state.in_buffer);
|
|
|
|
|
switch (connect_rsp_code) {
|
|
|
|
|
case CONNECTION_ACCEPTED:
|
2016-09-12 13:06:48 +07:00
|
|
|
mqtt_info("Connected");
|
2016-09-11 21:42:34 +07:00
|
|
|
return true;
|
|
|
|
|
case CONNECTION_REFUSE_PROTOCOL:
|
|
|
|
|
case CONNECTION_REFUSE_SERVER_UNAVAILABLE:
|
|
|
|
|
case CONNECTION_REFUSE_BAD_USERNAME:
|
|
|
|
|
case CONNECTION_REFUSE_NOT_AUTHORIZED:
|
2016-09-12 13:06:48 +07:00
|
|
|
mqtt_warn("Connection refuse, reason code: %d", connect_rsp_code);
|
2016-09-11 21:42:34 +07:00
|
|
|
return false;
|
|
|
|
|
default:
|
2016-09-12 13:06:48 +07:00
|
|
|
mqtt_warn("Connection refuse, Unknow reason");
|
2016-09-11 21:42:34 +07:00
|
|
|
return false;
|
|
|
|
|
}
|
|
|
|
|
return false;
|
|
|
|
|
}
|
2016-09-11 10:10:24 +07:00
|
|
|
|
2016-09-12 17:05:28 +07:00
|
|
|
void mqtt_sending_task(void *pvParameters)
|
|
|
|
|
{
|
|
|
|
|
mqtt_client *client = (mqtt_client *)pvParameters;
|
|
|
|
|
uint32_t msg_len, send_len;
|
|
|
|
|
mqtt_info("mqtt_sending_task");
|
|
|
|
|
|
|
|
|
|
while (1) {
|
|
|
|
|
if (xQueueReceive(client->xSendingQueue, &msg_len, 1000 / portTICK_RATE_MS)) {
|
|
|
|
|
//queue available
|
|
|
|
|
while (msg_len > 0) {
|
|
|
|
|
send_len = msg_len;
|
|
|
|
|
if (send_len > CONFIG_MQTT_BUFFER_SIZE_BYTE)
|
|
|
|
|
send_len = CONFIG_MQTT_BUFFER_SIZE_BYTE;
|
|
|
|
|
mqtt_info("Sending...%d bytes", send_len);
|
|
|
|
|
|
|
|
|
|
rb_read(&client->send_rb, client->mqtt_state.out_buffer, send_len);
|
|
|
|
|
client->mqtt_state.pending_msg_type = mqtt_get_type(client->mqtt_state.out_buffer);
|
|
|
|
|
client->mqtt_state.pending_msg_id = mqtt_get_id(client->mqtt_state.out_buffer, send_len);
|
|
|
|
|
write(client->socket, client->mqtt_state.out_buffer, send_len);
|
|
|
|
|
msg_len -= send_len;
|
|
|
|
|
}
|
|
|
|
|
//invalidate keepalive timer
|
|
|
|
|
client->keepalive_tick = client->settings->keepalive / 2;
|
|
|
|
|
}
|
|
|
|
|
else {
|
|
|
|
|
if (client->keepalive_tick > 0) client->keepalive_tick --;
|
|
|
|
|
else {
|
|
|
|
|
client->keepalive_tick = client->settings->keepalive / 2;
|
|
|
|
|
client->mqtt_state.outbound_message = mqtt_msg_pingreq(&client->mqtt_state.mqtt_connection);
|
|
|
|
|
client->mqtt_state.pending_msg_type = mqtt_get_type(client->mqtt_state.outbound_message->data);
|
|
|
|
|
client->mqtt_state.pending_msg_id = mqtt_get_id(client->mqtt_state.outbound_message->data,
|
|
|
|
|
client->mqtt_state.outbound_message->length);
|
|
|
|
|
mqtt_info("Sending pingreq");
|
|
|
|
|
write(client->socket,
|
|
|
|
|
client->mqtt_state.outbound_message->data,
|
|
|
|
|
client->mqtt_state.outbound_message->length);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
vTaskDelete(NULL);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void mqtt_start_receive_schedule(mqtt_client *client)
|
|
|
|
|
{
|
|
|
|
|
int read_len;
|
|
|
|
|
uint8_t msg_type;
|
|
|
|
|
uint8_t msg_qos;
|
|
|
|
|
uint16_t msg_id;
|
|
|
|
|
|
|
|
|
|
while (1) {
|
|
|
|
|
read_len = read(client->socket, client->mqtt_state.in_buffer, CONFIG_MQTT_BUFFER_SIZE_BYTE);
|
|
|
|
|
mqtt_info("Read len %d", read_len);
|
|
|
|
|
if (read_len == 0)
|
|
|
|
|
break;
|
|
|
|
|
|
|
|
|
|
msg_type = mqtt_get_type(client->mqtt_state.in_buffer);
|
|
|
|
|
msg_qos = mqtt_get_qos(client->mqtt_state.in_buffer);
|
|
|
|
|
msg_id = mqtt_get_id(client->mqtt_state.in_buffer, client->mqtt_state.in_buffer_length);
|
|
|
|
|
// mqtt_info("msg_type %d, msg_id: %d, pending_id: %d", msg_type, msg_id, client->mqtt_state.pending_msg_type);
|
|
|
|
|
switch (msg_type)
|
|
|
|
|
{
|
|
|
|
|
case MQTT_MSG_TYPE_SUBACK:
|
|
|
|
|
if (client->mqtt_state.pending_msg_type == MQTT_MSG_TYPE_SUBSCRIBE && client->mqtt_state.pending_msg_id == msg_id)
|
|
|
|
|
mqtt_info("Subscribe successful");
|
|
|
|
|
break;
|
|
|
|
|
case MQTT_MSG_TYPE_UNSUBACK:
|
|
|
|
|
if (client->mqtt_state.pending_msg_type == MQTT_MSG_TYPE_UNSUBSCRIBE && client->mqtt_state.pending_msg_id == msg_id)
|
|
|
|
|
mqtt_info("UnSubscribe successful");
|
|
|
|
|
break;
|
|
|
|
|
case MQTT_MSG_TYPE_PUBLISH:
|
|
|
|
|
if (msg_qos == 1)
|
|
|
|
|
client->mqtt_state.outbound_message = mqtt_msg_puback(&client->mqtt_state.mqtt_connection, msg_id);
|
|
|
|
|
else if (msg_qos == 2)
|
|
|
|
|
client->mqtt_state.outbound_message = mqtt_msg_pubrec(&client->mqtt_state.mqtt_connection, msg_id);
|
|
|
|
|
if (msg_qos == 1 || msg_qos == 2) {
|
|
|
|
|
mqtt_info("Queue response QoS: %d", msg_qos);
|
|
|
|
|
// if (QUEUE_Puts(&client->msgQueue, client->mqtt_state.outbound_message->data, client->mqtt_state.outbound_message->length) == -1) {
|
|
|
|
|
// mqtt_info("MQTT: Queue full");
|
|
|
|
|
// }
|
|
|
|
|
}
|
|
|
|
|
// deliver_publish(client, client->mqtt_state.in_buffer, client->mqtt_state.message_length_read);
|
|
|
|
|
break;
|
|
|
|
|
case MQTT_MSG_TYPE_PUBACK:
|
|
|
|
|
if (client->mqtt_state.pending_msg_type == MQTT_MSG_TYPE_PUBLISH && client->mqtt_state.pending_msg_id == msg_id) {
|
|
|
|
|
mqtt_info("received MQTT_MSG_TYPE_PUBACK, finish QoS1 publish");
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
break;
|
|
|
|
|
case MQTT_MSG_TYPE_PUBREC:
|
|
|
|
|
client->mqtt_state.outbound_message = mqtt_msg_pubrel(&client->mqtt_state.mqtt_connection, msg_id);
|
|
|
|
|
// if (QUEUE_Puts(&client->msgQueue, client->mqtt_state.outbound_message->data, client->mqtt_state.outbound_message->length) == -1) {
|
|
|
|
|
// mqtt_info("MQTT: Queue full");
|
|
|
|
|
// }
|
|
|
|
|
break;
|
|
|
|
|
case MQTT_MSG_TYPE_PUBREL:
|
|
|
|
|
client->mqtt_state.outbound_message = mqtt_msg_pubcomp(&client->mqtt_state.mqtt_connection, msg_id);
|
|
|
|
|
// if (QUEUE_Puts(&client->msgQueue, client->mqtt_state.outbound_message->data, client->mqtt_state.outbound_message->length) == -1) {
|
|
|
|
|
// mqtt_info("MQTT: Queue full");
|
|
|
|
|
// }
|
|
|
|
|
break;
|
|
|
|
|
case MQTT_MSG_TYPE_PUBCOMP:
|
|
|
|
|
if (client->mqtt_state.pending_msg_type == MQTT_MSG_TYPE_PUBLISH && client->mqtt_state.pending_msg_id == msg_id) {
|
|
|
|
|
mqtt_info("eceive MQTT_MSG_TYPE_PUBCOMP, finish QoS2 publish");
|
|
|
|
|
}
|
|
|
|
|
break;
|
|
|
|
|
case MQTT_MSG_TYPE_PINGREQ:
|
|
|
|
|
client->mqtt_state.outbound_message = mqtt_msg_pingresp(&client->mqtt_state.mqtt_connection);
|
|
|
|
|
// if (QUEUE_Puts(&client->msgQueue, client->mqtt_state.outbound_message->data, client->mqtt_state.outbound_message->length) == -1) {
|
|
|
|
|
// mqtt_info("MQTT: Queue full");
|
|
|
|
|
// }
|
|
|
|
|
break;
|
|
|
|
|
case MQTT_MSG_TYPE_PINGRESP:
|
|
|
|
|
// Ignore
|
|
|
|
|
break;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
mqtt_info("network disconnected");
|
|
|
|
|
}
|
2016-09-12 13:06:48 +07:00
|
|
|
|
|
|
|
|
void mqtt_destroy(mqtt_client *client)
|
|
|
|
|
{
|
|
|
|
|
free(client->mqtt_state.in_buffer);
|
|
|
|
|
free(client->mqtt_state.out_buffer);
|
|
|
|
|
free(client);
|
2016-09-12 17:05:28 +07:00
|
|
|
vTaskDelete(xMqttTask);
|
2016-09-12 13:06:48 +07:00
|
|
|
}
|
|
|
|
|
|
2016-09-11 10:10:24 +07:00
|
|
|
void mqtt_task(void *pvParameters)
|
|
|
|
|
{
|
2016-09-11 21:42:34 +07:00
|
|
|
mqtt_client *client = (mqtt_client *)pvParameters;
|
|
|
|
|
|
2016-09-12 17:05:28 +07:00
|
|
|
while (1) {
|
|
|
|
|
client->socket = client_connect(client->settings->host, client->settings->port);
|
|
|
|
|
mqtt_info("Connected to server %s:%d", client->settings->host, client->settings->port);
|
|
|
|
|
if (!mqtt_connect(client)) {
|
|
|
|
|
continue;
|
|
|
|
|
//return;
|
|
|
|
|
}
|
|
|
|
|
mqtt_info("Connected to MQTT broker, create sending thread before call connected callback");
|
|
|
|
|
xTaskCreate(&mqtt_sending_task, "mqtt_sending_task", 2048, client, CONFIG_MQTT_PRIORITY + 1, &xMqttSendingTask);
|
|
|
|
|
if (client->settings->connected_cb) {
|
|
|
|
|
client->settings->connected_cb(client);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
mqtt_info("mqtt_start_receive_schedule");
|
|
|
|
|
mqtt_start_receive_schedule(client);
|
|
|
|
|
|
2016-09-11 21:42:34 +07:00
|
|
|
close(client->socket);
|
2016-09-12 17:05:28 +07:00
|
|
|
vTaskDelete(xMqttSendingTask);
|
|
|
|
|
vTaskDelay(1000 / portTICK_RATE_MS);
|
|
|
|
|
|
2016-09-11 21:42:34 +07:00
|
|
|
}
|
2016-09-12 13:06:48 +07:00
|
|
|
mqtt_destroy(client);
|
2016-09-12 17:05:28 +07:00
|
|
|
|
2016-09-11 21:42:34 +07:00
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
2016-09-12 13:06:48 +07:00
|
|
|
void mqtt_start(mqtt_settings *settings)
|
2016-09-11 21:42:34 +07:00
|
|
|
{
|
2016-09-12 17:05:28 +07:00
|
|
|
uint8_t *rb_buf;
|
2016-09-11 21:42:34 +07:00
|
|
|
if (xMqttTask != NULL)
|
|
|
|
|
return;
|
|
|
|
|
mqtt_client *client = malloc(sizeof(mqtt_client));
|
2016-09-12 17:05:28 +07:00
|
|
|
|
2016-09-11 21:42:34 +07:00
|
|
|
if (client == NULL) {
|
2016-09-12 17:05:28 +07:00
|
|
|
mqtt_error("Memory not enough");
|
2016-09-11 21:42:34 +07:00
|
|
|
return;
|
|
|
|
|
}
|
2016-09-12 17:05:28 +07:00
|
|
|
memset(client, 0, sizeof(mqtt_client));
|
|
|
|
|
|
2016-09-12 13:06:48 +07:00
|
|
|
client->settings = settings;
|
|
|
|
|
client->connect_info.client_id = settings->client_id;
|
|
|
|
|
client->connect_info.username = settings->username;
|
|
|
|
|
client->connect_info.password = settings->password;
|
|
|
|
|
client->connect_info.will_topic = settings->lwt_topic;
|
|
|
|
|
client->connect_info.will_message = settings->lwt_msg;
|
|
|
|
|
client->connect_info.will_qos = settings->lwt_qos;
|
|
|
|
|
client->connect_info.will_retain = settings->lwt_retain;
|
2016-09-12 17:05:28 +07:00
|
|
|
client->keepalive_tick = settings->keepalive / 2;
|
2016-09-12 13:06:48 +07:00
|
|
|
|
|
|
|
|
client->connect_info.keepalive = settings->keepalive;
|
|
|
|
|
client->connect_info.clean_session = settings->clean_session;
|
2016-09-11 21:42:34 +07:00
|
|
|
|
|
|
|
|
client->mqtt_state.in_buffer = (uint8_t *)malloc(CONFIG_MQTT_BUFFER_SIZE_BYTE);
|
|
|
|
|
client->mqtt_state.in_buffer_length = CONFIG_MQTT_BUFFER_SIZE_BYTE;
|
|
|
|
|
client->mqtt_state.out_buffer = (uint8_t *)malloc(CONFIG_MQTT_BUFFER_SIZE_BYTE);
|
|
|
|
|
client->mqtt_state.out_buffer_length = CONFIG_MQTT_BUFFER_SIZE_BYTE;
|
|
|
|
|
client->mqtt_state.connect_info = &client->connect_info;
|
|
|
|
|
|
2016-09-12 17:05:28 +07:00
|
|
|
|
|
|
|
|
|
|
|
|
|
/* Create a queue capable of containing 64 unsigned long values. */
|
|
|
|
|
client->xSendingQueue = xQueueCreate(64, sizeof( uint32_t ));
|
|
|
|
|
rb_buf = (uint8_t*) malloc(CONFIG_MQTT_QUEUE_BUFFER_SIZE_WORD * 4);
|
|
|
|
|
|
|
|
|
|
if (rb_buf == NULL) {
|
|
|
|
|
mqtt_error("Memory not enough");
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
rb_init(&client->send_rb, rb_buf, CONFIG_MQTT_QUEUE_BUFFER_SIZE_WORD * 4, 1);
|
|
|
|
|
|
2016-09-11 21:42:34 +07:00
|
|
|
mqtt_msg_init(&client->mqtt_state.mqtt_connection,
|
|
|
|
|
client->mqtt_state.out_buffer,
|
|
|
|
|
client->mqtt_state.out_buffer_length);
|
|
|
|
|
|
2016-09-12 13:06:48 +07:00
|
|
|
xTaskCreate(&mqtt_task, "mqtt_task", 2048, client, CONFIG_MQTT_PRIORITY, &xMqttTask);
|
2016-09-11 21:42:34 +07:00
|
|
|
}
|
|
|
|
|
|
2016-09-12 17:05:28 +07:00
|
|
|
void mqtt_subscribe(mqtt_client *client, char *topic, uint8_t qos)
|
|
|
|
|
{
|
|
|
|
|
|
|
|
|
|
client->mqtt_state.outbound_message = mqtt_msg_subscribe(&client->mqtt_state.mqtt_connection,
|
|
|
|
|
topic, qos,
|
|
|
|
|
&client->mqtt_state.pending_msg_id);
|
|
|
|
|
mqtt_info("MQTT: queue subscribe, topic\"%s\", id: %d", topic, client->mqtt_state.pending_msg_id);
|
|
|
|
|
rb_write(&client->send_rb,
|
|
|
|
|
client->mqtt_state.outbound_message->data,
|
|
|
|
|
client->mqtt_state.outbound_message->length);
|
|
|
|
|
xQueueSend(client->xSendingQueue, &client->mqtt_state.outbound_message->length, 0);
|
|
|
|
|
|
|
|
|
|
}
|
2016-09-11 21:42:34 +07:00
|
|
|
void mqtt_stop()
|
|
|
|
|
{
|
|
|
|
|
|
2016-09-11 10:10:24 +07:00
|
|
|
}
|
2016-09-11 21:42:34 +07:00
|
|
|
|