Files
esp-mqtt/mqtt.c

540 lines
19 KiB
C
Raw Normal View History

2016-09-11 10:10:24 +07:00
/*
* @Author: Tuan PM
* @Date: 2016-09-10 09:33:06
* @Last Modified by: Tuan PM
2017-02-15 13:16:44 +07:00
* @Last Modified time: 2017-02-15 13:11:53
2016-09-11 10:10:24 +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"
#include "ringbuf.h"
#include "mqtt.h"
2016-09-11 21:42:34 +07:00
static TaskHandle_t xMqttTask = NULL;
static TaskHandle_t xMqttSendingTask = NULL;
2016-09-11 21:42:34 +07:00
2016-09-23 22:07:46 +07:00
static int resolve_dns(const char *host, struct sockaddr_in *ip) {
2016-09-11 21:42:34 +07:00
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;
}
2016-09-12 21:15:28 +07:00
static void mqtt_queue(mqtt_client *client)
{
2017-02-15 13:16:44 +07:00
int msg_len;
2017-02-15 13:03:07 +07:00
while (rb_available(&client->send_rb) < client->mqtt_state.outbound_message->length) {
xQueueReceive(client->xSendingQueue, &msg_len, 1000 / portTICK_RATE_MS);
rb_read(&client->send_rb, client->mqtt_state.out_buffer, msg_len);
}
2016-09-12 21:15:28 +07:00
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);
}
static bool client_connect(mqtt_client *client)
2016-09-11 21:42:34 +07:00
{
int ret;
2016-09-12 13:06:48 +07:00
struct sockaddr_in remote_ip;
2017-02-15 13:16:44 +07:00
while (1) {
2016-09-11 21:42:34 +07:00
bzero(&remote_ip, sizeof(struct sockaddr_in));
2016-09-12 13:06:48 +07:00
remote_ip.sin_family = AF_INET;
remote_ip.sin_port = htons(client->settings->port);
//if host is not ip address, resolve it
if (inet_aton( client->settings->host, &(remote_ip.sin_addr)) == 0) {
mqtt_info("Resolve dns for domain: %s", client->settings->host);
if (!resolve_dns(client->settings->host, &remote_ip)) {
2016-09-11 21:42:34 +07:00
vTaskDelay(1000 / portTICK_RATE_MS);
continue;
}
}
#if defined(CONFIG_MQTT_SECURITY_ON) // ENABLE MQTT OVER SSL
client->ctx = NULL;
client->ssl = NULL;
client->ctx = SSL_CTX_new(TLSv1_2_client_method());
if (!client->ctx) {
mqtt_error("Failed to create SSL CTX");
goto failed1;
}
#endif
2017-02-15 13:16:44 +07:00
client->socket = socket(PF_INET, SOCK_STREAM, 0);
if (client->socket == -1) {
mqtt_error("Failed to create socket");
goto failed2;
2016-09-11 21:42:34 +07:00
}
2016-09-12 13:06:48 +07:00
mqtt_info("Connecting to server %s:%d,%d",
inet_ntoa((remote_ip.sin_addr)),
client->settings->port,
2016-09-12 13:06:48 +07:00
remote_ip.sin_port);
2017-02-15 13:16:44 +07:00
if (connect(client->socket, (struct sockaddr *)(&remote_ip), sizeof(struct sockaddr)) != 00) {
mqtt_error("Connect failed");
goto failed3;
2016-09-11 21:42:34 +07:00
}
#if defined(CONFIG_MQTT_SECURITY_ON) // ENABLE MQTT OVER SSL
mqtt_info("Creating SSL object...");
client->ssl = SSL_new(client->ctx);
if (!client->ssl) {
mqtt_error("Unable to creat new SSL");
goto failed3;
}
if (!SSL_set_fd(client->ssl, client->socket)) {
mqtt_error("SSL set_fd failed");
goto failed3;
}
mqtt_info("Start SSL connect..");
ret = SSL_connect(client->ssl);
if (!ret) {
mqtt_error("SSL Connect FAILED");
goto failed4;
}
#endif
mqtt_info("Connected!");
return true;
2017-02-15 13:16:44 +07:00
//failed5:
// SSL_shutdown(client->ssl);
2017-02-15 13:16:44 +07:00
#if defined(CONFIG_MQTT_SECURITY_ON)
failed4:
SSL_free(client->ssl);
client->ssl = NULL;
#endif
failed3:
close(client->socket);
client->socket = -1;
failed2:
2017-02-15 13:16:44 +07:00
#if defined(CONFIG_MQTT_SECURITY_ON)
SSL_CTX_free(client->ctx);
failed1:
client->ctx = NULL;
2017-02-15 13:16:44 +07:00
#endif
vTaskDelay(1000 / portTICK_RATE_MS);
}
}
// Close client socket
// including SSL objects if CNFIG_MQTT_SECURITY_ON is enabled
void closeclient(mqtt_client *client)
{
2017-02-15 13:16:44 +07:00
#if defined(CONFIG_MQTT_SECURITY_ON)
if (client->ssl != NULL)
{
SSL_shutdown(client->ssl);
SSL_free(client->ssl);
client->ssl = NULL;
}
#endif
if (client->socket != -1)
{
close(client->socket);
client->socket = -1;
}
2017-02-15 13:16:44 +07:00
#if defined(CONFIG_MQTT_SECURITY_ON)
if (client->ctx != NULL)
{
SSL_CTX_free(client->ctx);
client->ctx = NULL;
}
2017-02-15 13:16:44 +07:00
#endif
2016-09-11 21:42:34 +07:00
}
/*
* 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 = ClientWrite(
2016-09-11 21:42:34 +07:00
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");
read_len = ClientRead(client->mqtt_state.in_buffer, CONFIG_MQTT_BUFFER_SIZE_BYTE);
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
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);
ClientWrite(client->mqtt_state.out_buffer, send_len);
2016-09-12 21:15:28 +07:00
//TODO: Check sending type, to callback publish message
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");
ClientWrite(
client->mqtt_state.outbound_message->data,
client->mqtt_state.outbound_message->length);
}
}
}
vTaskDelete(NULL);
}
2016-09-12 21:15:28 +07:00
void deliver_publish(mqtt_client *client, uint8_t *message, int length)
{
mqtt_event_data_t event_data;
int len_read, total_mqtt_len = 0, mqtt_len = 0, mqtt_offset = 0;
do
{
event_data.topic_length = length;
event_data.topic = mqtt_get_publish_topic(message, &event_data.topic_length);
event_data.data_length = length;
event_data.data = mqtt_get_publish_data(message, &event_data.data_length);
if(total_mqtt_len == 0){
total_mqtt_len = client->mqtt_state.message_length - client->mqtt_state.message_length_read + event_data.data_length;
mqtt_len = event_data.data_length;
} else {
mqtt_len = len_read;
}
event_data.data_total_length = total_mqtt_len;
event_data.data_offset = mqtt_offset;
event_data.data_length = mqtt_len;
mqtt_info("Data received: %d/%d bytes ", mqtt_len, total_mqtt_len);
if(client->settings->data_cb) {
client->settings->data_cb(client, &event_data);
}
mqtt_offset += mqtt_len;
if (client->mqtt_state.message_length_read >= client->mqtt_state.message_length)
break;
len_read = ClientRead(client->mqtt_state.in_buffer, CONFIG_MQTT_BUFFER_SIZE_BYTE);
2016-09-12 21:15:28 +07:00
client->mqtt_state.message_length_read += len_read;
} while (1);
}
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 = ClientRead(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:
2016-09-12 21:15:28 +07:00
if (client->mqtt_state.pending_msg_type == MQTT_MSG_TYPE_SUBSCRIBE && client->mqtt_state.pending_msg_id == msg_id) {
mqtt_info("Subscribe successful");
2016-09-12 21:15:28 +07:00
if (client->settings->subscribe_cb) {
client->settings->subscribe_cb(client, NULL);
}
}
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);
2016-09-12 21:15:28 +07:00
if (msg_qos == 1 || msg_qos == 2) {
mqtt_info("Queue response QoS: %d", msg_qos);
2016-09-12 21:15:28 +07:00
mqtt_queue(client);
// if (QUEUE_Puts(&client->msgQueue, client->mqtt_state.outbound_message->data, client->mqtt_state.outbound_message->length) == -1) {
// mqtt_info("MQTT: Queue full");
// }
}
2016-09-12 21:15:28 +07:00
client->mqtt_state.message_length_read = read_len;
client->mqtt_state.message_length = mqtt_get_total_length(client->mqtt_state.in_buffer, client->mqtt_state.message_length_read);
mqtt_info("deliver_publish");
deliver_publish(client, client->mqtt_state.in_buffer, client->mqtt_state.message_length_read);
// 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);
2016-09-12 21:15:28 +07:00
mqtt_queue(client);
break;
case MQTT_MSG_TYPE_PUBREL:
client->mqtt_state.outbound_message = mqtt_msg_pubcomp(&client->mqtt_state.mqtt_connection, msg_id);
2016-09-12 21:15:28 +07:00
mqtt_queue(client);
break;
case MQTT_MSG_TYPE_PUBCOMP:
2017-04-30 22:24:09 +08:00
if (client->mqtt_state.pending_msg_type == MQTT_MSG_TYPE_PUBREL && client->mqtt_state.pending_msg_id == msg_id) {
2016-09-12 21:15:28 +07:00
mqtt_info("Receive 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);
2016-09-12 21:15:28 +07:00
mqtt_queue(client);
break;
case MQTT_MSG_TYPE_PINGRESP:
2016-09-12 21:15:28 +07:00
mqtt_info("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);
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;
while (1) {
client_connect(client);
2017-02-15 13:16:44 +07:00
mqtt_info("Connected to server %s:%d", client->settings->host, client->settings->port);
if (!mqtt_connect(client)) {
closeclient(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) {
2016-09-12 21:15:28 +07:00
client->settings->connected_cb(client, NULL);
}
mqtt_info("mqtt_start_receive_schedule");
mqtt_start_receive_schedule(client);
2017-02-15 13:16:44 +07:00
closeclient(client);
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-11 21:42:34 +07:00
}
2016-09-12 21:15:28 +07:00
mqtt_client *mqtt_start(mqtt_settings *settings)
2016-09-11 21:42:34 +07:00
{
int stackSize = 2048;
uint8_t *rb_buf;
2016-09-11 21:42:34 +07:00
if (xMqttTask != NULL)
2016-09-12 21:15:28 +07:00
return NULL;
2016-09-11 21:42:34 +07:00
mqtt_client *client = malloc(sizeof(mqtt_client));
2016-09-11 21:42:34 +07:00
if (client == NULL) {
mqtt_error("Memory not enough");
2016-09-12 21:15:28 +07:00
return NULL;
2016-09-11 21:42:34 +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-13 00:19:00 +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;
client->socket = -1;
#if defined(CONFIG_MQTT_SECURITY_ON) // ENABLE MQTT OVER SSL
client->ctx = NULL;
client->ssl = NULL;
stackSize = 10240; // Need more stack to handle SSL handshake
#endif
/* 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");
2016-09-12 21:15:28 +07:00
return NULL;
}
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);
xTaskCreate(&mqtt_task, "mqtt_task", stackSize, client, CONFIG_MQTT_PRIORITY, &xMqttTask);
2016-09-12 21:15:28 +07:00
return client;
2016-09-11 21:42:34 +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);
2016-09-12 21:15:28 +07:00
mqtt_info("Queue subscribe, topic\"%s\", id: %d", topic, client->mqtt_state.pending_msg_id);
mqtt_queue(client);
}
2016-09-12 21:15:28 +07:00
void mqtt_publish(mqtt_client* client, char *topic, char *data, int len, int qos, int retain)
{
client->mqtt_state.outbound_message = mqtt_msg_publish(&client->mqtt_state.mqtt_connection,
topic, data, len,
qos, retain,
&client->mqtt_state.pending_msg_id);
mqtt_queue(client);
mqtt_info("Queuing publish, length: %d, queue size(%d/%d)\r\n",
client->mqtt_state.outbound_message->length,
client->send_rb.fill_cnt,
client->send_rb.size);
}
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