46#include "contiki-net.h"
47#include "contiki-lib.h"
56#include "tcp-socket.h"
58#include "lib/assert.h"
69#define PRINTF(...) printf(__VA_ARGS__)
75 MQTT_FHDR_DUP_FLAG = 0x08,
77 MQTT_FHDR_QOS_LEVEL_0 = 0x00,
78 MQTT_FHDR_QOS_LEVEL_1 = 0x02,
79 MQTT_FHDR_QOS_LEVEL_2 = 0x04,
81 MQTT_FHDR_RETAIN_FLAG = 0x01,
85 MQTT_VHDR_USERNAME_FLAG = 0x80,
86 MQTT_VHDR_PASSWORD_FLAG = 0x40,
88 MQTT_VHDR_WILL_RETAIN_FLAG = 0x20,
89 MQTT_VHDR_WILL_QOS_LEVEL_0 = 0x00,
90 MQTT_VHDR_WILL_QOS_LEVEL_1 = 0x08,
91 MQTT_VHDR_WILL_QOS_LEVEL_2 = 0x10,
93 MQTT_VHDR_WILL_FLAG = 0x04,
94 MQTT_VHDR_CLEAN_SESSION_FLAG = 0x02,
95} mqtt_vhdr_conn_fields_t;
98 MQTT_VHDR_CONN_ACCEPTED,
99 MQTT_VHDR_CONN_REJECTED_PROTOCOL,
100 MQTT_VHDR_CONN_REJECTED_IDENTIFIER,
101 MQTT_VHDR_CONN_REJECTED_UNAVAILABLE,
102 MQTT_VHDR_CONN_REJECTED_BAD_USER_PASS,
103 MQTT_VHDR_CONN_REJECTED_UNAUTHORIZED,
104} mqtt_vhdr_connack_ret_code_t;
107 MQTT_VHDR_CONNACK_SESSION_PRESENT = 0x1
108} mqtt_vhdr_connack_flags_t;
112 MQTT_SUBACK_RET_QOS_0 = 0x00,
113 MQTT_SUBACK_RET_QOS_1 = 0x01,
114 MQTT_SUBACK_RET_QOS_2 = 0x02,
115 MQTT_SUBACK_RET_FAIL = 0x08,
116} mqtt_suback_ret_code_t;
120 MQTT_VHDR_RC_SUCCES_OR_NORMAL = 0x00,
121 MQTT_VHDR_RC_QOS_0 = 0x01,
122 MQTT_VHDR_RC_QOS_1 = 0x02,
123 MQTT_VHDR_RC_QOS_2 = 0x03,
124 MQTT_VHDR_RC_DISC_WITH_WILL = 0x04,
125 MQTT_VHDR_RC_NO_MATCH_SUB = 0x10,
126 MQTT_VHDR_RC_NO_SUB_EXISTED = 0x11,
127 MQTT_VHDR_RC_CONTINUE_AUTH = 0x18,
128 MQTT_VHDR_RC_REAUTH = 0x19,
129 MQTT_VHDR_RC_UNSPEC_ERR = 0x80,
130 MQTT_VHDR_RC_MALFORMED_PKT = 0x81,
131 MQTT_VHDR_RC_PROTOCOL_ERR = 0x82,
132 MQTT_VHDR_RC_IMPL_SPEC_ERR = 0x83,
133 MQTT_VHDR_RC_PROT_VER_UNUSUPPORTED = 0x84,
134 MQTT_VHDR_RC_CLIENT_ID_INVALID = 0x85,
135 MQTT_VHDR_RC_BAD_USER_PASS = 0x86,
136 MQTT_VHDR_RC_NOT_AUTH = 0x87,
137 MQTT_VHDR_RC_SRV_UNAVAIL = 0x88,
138 MQTT_VHDR_RC_SRV_BUSY = 0x89,
139 MQTT_VHDR_RC_BANNED = 0x8A,
140 MQTT_VHDR_RC_SRV_SHUTDOWN = 0x8B,
141 MQTT_VHDR_RC_BAD_AUTH_METHOD = 0x8C,
142 MQTT_VHDR_RC_KEEP_ALIVE_TIMEOUT = 0x8D,
143 MQTT_VHDR_RC_SESS_TAKEN_OVER = 0x8E,
144 MQTT_VHDR_RC_TOPIC_FILT_INVAL = 0x8F,
145 MQTT_VHDR_RC_TOPIC_NAME_INVAL = 0x90,
146 MQTT_VHDR_RC_PKT_ID_IN_USE = 0x91,
147 MQTT_VHDR_RC_PKT_ID_NOT_FOUND = 0x92,
148 MQTT_VHDR_RC_RECV_MAX_EXCEEDED = 0x93,
149 MQTT_VHDR_RC_TOPIC_ALIAS_INVAL = 0x94,
150 MQTT_VHDR_RC_PKT_TOO_LARGE = 0x95,
151 MQTT_VHDR_RC_MSG_RATE_TOO_HIGH = 0x96,
152 MQTT_VHDR_RC_QUOTA_EXCEEDED = 0x97,
153 MQTT_VHDR_RC_ADMIN_ACTION = 0x98,
154 MQTT_VHDR_RC_PAYLD_FMT_INVAL = 0x99,
155 MQTT_VHDR_RC_RETAIN_UNSUPPORTED = 0x9A,
156 MQTT_VHDR_RC_QOS_UNSUPPORTED = 0x9B,
157 MQTT_VHDR_RC_USE_ANOTHER_SRV = 0x9C,
158 MQTT_VHDR_RC_SRV_MOVED = 0x9D,
159 MQTT_VHDR_RC_SHARED_SUB_UNSUPPORTED = 0x9E,
160 MQTT_VHDR_RC_CONN_RATE_EXCEEDED = 0x9F,
161 MQTT_VHDR_RC_MAX_CONN_TIME = 0xA0,
162 MQTT_VHDR_RC_SUB_ID_UNSUPPORTED = 0xA1,
163 MQTT_VHDR_RC_WILD_SUB_UNSUPPORTED = 0xA2,
166#define RESPONSE_WAIT_TIMEOUT (CLOCK_SECOND * 10)
168#define INCREMENT_MID(conn) (conn)->mid_counter += 2
169#define MQTT_STRING_LENGTH(s) (((s)->length) == 0 ? 0 : (MQTT_STRING_LEN_SIZE + (s)->length))
172#define PT_MQTT_WRITE_BYTES(conn, data, len) \
173 conn->out_write_pos = 0; \
174 while(write_bytes(conn, data, len)) { \
175 PT_WAIT_UNTIL(pt, (conn)->out_buffer_sent); \
178#define PT_MQTT_WRITE_BYTE(conn, data) \
179 while(write_byte(conn, data)) { \
180 PT_WAIT_UNTIL(pt, (conn)->out_buffer_sent); \
189#define PT_MQTT_WAIT_SEND() \
191 if (PROCESS_ERR_OK == \
192 process_post(PROCESS_CURRENT(), mqtt_continue_send_event, NULL)) { \
194 PROCESS_WAIT_EVENT(); \
195 if(ev == mqtt_abort_now_event) { \
196 conn->state = MQTT_CONN_STATE_ABORT_IMMEDIATE; \
197 PT_INIT(&conn->out_proto_thread); \
198 process_post(PROCESS_CURRENT(), ev, data); \
199 } else if(ev >= mqtt_event_min && ev <= mqtt_event_max) { \
200 process_post(PROCESS_CURRENT(), ev, data); \
202 } while (ev != mqtt_continue_send_event); \
206static process_event_t mqtt_do_connect_tcp_event;
207static process_event_t mqtt_do_connect_mqtt_event;
208static process_event_t mqtt_do_disconnect_mqtt_event;
209static process_event_t mqtt_do_subscribe_event;
210static process_event_t mqtt_do_unsubscribe_event;
211static process_event_t mqtt_do_publish_event;
212static process_event_t mqtt_do_pingreq_event;
213static process_event_t mqtt_continue_send_event;
214static process_event_t mqtt_abort_now_event;
215static process_event_t mqtt_do_auth_event;
216process_event_t mqtt_update_event;
223static process_event_t mqtt_event_min;
224static process_event_t mqtt_event_max;
228tcp_input(
struct tcp_socket *s,
void *ptr,
const uint8_t *input_data_ptr,
231static void tcp_event(
struct tcp_socket *s,
void *ptr,
232 tcp_socket_event_t event);
234static void reset_packet(
struct mqtt_in_packet *packet);
238PROCESS(mqtt_process,
"MQTT process");
241call_event(
struct mqtt_connection *conn,
245 conn->event_callback(conn, event, data);
246 process_post(conn->app_process, mqtt_update_event, NULL);
250reset_defaults(
struct mqtt_connection *conn)
252 conn->mid_counter = 1;
253 PT_INIT(&conn->out_proto_thread);
254 conn->waiting_for_pingresp = 0;
256 reset_packet(&conn->in_packet);
257 conn->out_buffer_sent = 0;
261abort_connection(
struct mqtt_connection *conn)
263 conn->out_buffer_ptr = conn->out_buffer;
264 conn->out_queue_full = 0;
267 memset(&conn->out_packet, 0,
sizeof(conn->out_packet));
269 tcp_socket_close(&conn->socket);
270 tcp_socket_unregister(&conn->socket);
272 memset(&conn->socket, 0,
sizeof(conn->socket));
274 conn->state = MQTT_CONN_STATE_NOT_CONNECTED;
278connect_tcp(
struct mqtt_connection *conn)
280 conn->state = MQTT_CONN_STATE_TCP_CONNECTING;
282 reset_defaults(conn);
283 tcp_socket_register(&(conn->socket),
286 MQTT_TCP_INPUT_BUFF_SIZE,
288 MQTT_TCP_OUTPUT_BUFF_SIZE,
291 tcp_socket_connect(&(conn->socket), &(conn->server_ip), conn->server_port);
295disconnect_tcp(
struct mqtt_connection *conn)
297 conn->state = MQTT_CONN_STATE_DISCONNECTING;
298 tcp_socket_close(&(conn->socket));
299 tcp_socket_unregister(&conn->socket);
301 memset(&conn->socket, 0,
sizeof(conn->socket));
305send_out_buffer(
struct mqtt_connection *conn)
307 if(conn->out_buffer_ptr - conn->out_buffer == 0) {
308 conn->out_buffer_sent = 1;
311 conn->out_buffer_sent = 0;
313 DBG(
"MQTT - (send_out_buffer) Space used in buffer: %i\n",
314 conn->out_buffer_ptr - conn->out_buffer);
316 tcp_socket_send(&conn->socket, conn->out_buffer,
317 conn->out_buffer_ptr - conn->out_buffer);
321string_to_mqtt_string(
struct mqtt_string *mqtt_string,
char *
string)
323 if(mqtt_string == NULL) {
326 mqtt_string->string = string;
329 mqtt_string->length = strlen(
string);
331 mqtt_string->length = 0;
336write_byte(
struct mqtt_connection *conn, uint8_t data)
338 DBG(
"MQTT - (write_byte) buff_size: %i write: '%02X'\n",
339 &conn->out_buffer[MQTT_TCP_OUTPUT_BUFF_SIZE] - conn->out_buffer_ptr,
342 if(&conn->out_buffer[MQTT_TCP_OUTPUT_BUFF_SIZE] - conn->out_buffer_ptr == 0) {
343 send_out_buffer(conn);
347 *conn->out_buffer_ptr = data;
348 conn->out_buffer_ptr++;
353write_bytes(
struct mqtt_connection *conn, uint8_t *data, uint16_t len)
355 uint16_t write_bytes;
357 MIN(&conn->out_buffer[MQTT_TCP_OUTPUT_BUFF_SIZE] - conn->out_buffer_ptr,
358 len - conn->out_write_pos);
360 memcpy(conn->out_buffer_ptr, &data[conn->out_write_pos], write_bytes);
361 conn->out_write_pos += write_bytes;
362 conn->out_buffer_ptr += write_bytes;
364 DBG(
"MQTT - (write_bytes) len: %u write_pos: %i\n", len,
365 conn->out_write_pos);
367 if(len - conn->out_write_pos == 0) {
368 conn->out_write_pos = 0;
371 send_out_buffer(conn);
372 return len - conn->out_write_pos;
377mqtt_decode_var_byte_int(
const uint8_t *input_data_ptr,
380 uint32_t *pkt_byte_count,
383 uint8_t read_bytes = 0;
385 uint32_t multiplier = 1;
386 uint32_t input_pos_0 = 0;
388 if(input_pos == NULL) {
389 input_pos = &input_pos_0;
395 if(*input_pos >= input_data_len) {
399 byte_in = input_data_ptr[*input_pos];
405 DBG(
"MQTT - Read Variable Byte Integer byte %i\n", byte_in);
408 DBG(
"Received more than 4 byte 'Variable Byte Integer'.");
412 *dest += (byte_in & 127) * multiplier;
414 }
while((byte_in & 128) != 0);
420mqtt_encode_var_byte_int(uint8_t *vbi_out,
426 DBG(
"MQTT - Encoding Variable Byte Integer %u\n", val);
433 digit = digit | 0x80;
436 vbi_out[*vbi_bytes] = digit;
438 DBG(
"MQTT - Encode VBI digit '%u' length '%i'\n", digit, val);
439 }
while(val > 0 && *vbi_bytes < 5);
440 DBG(
"MQTT - var_byte_int bytes %u\n", *vbi_bytes);
444keep_alive_callback(
void *ptr)
446 struct mqtt_connection *conn = ptr;
448 DBG(
"MQTT - (keep_alive_callback) Called!\n");
451 if(conn->waiting_for_pingresp) {
452 PRINTF(
"MQTT - Disconnect due to no PINGRESP from broker.\n");
453 disconnect_tcp(conn);
457 process_post(&mqtt_process, mqtt_do_pingreq_event, conn);
461reset_packet(
struct mqtt_in_packet *packet)
463 memset(packet, 0,
sizeof(
struct mqtt_in_packet));
468PT_THREAD(write_out_props(
struct pt *pt,
struct mqtt_connection *conn,
469 struct mqtt_prop_list *prop_list))
473 static struct mqtt_prop_out_property *prop;
476 DBG(
"MQTT - Writing %i property bytes\n", prop_list->properties_len + prop_list->properties_len_enc_bytes);
478 PT_MQTT_WRITE_BYTES(conn,
479 prop_list->properties_len_enc,
480 prop_list->properties_len_enc_bytes);
482 prop = (
struct mqtt_prop_out_property *)
list_head(prop_list->props);
485 DBG(
"MQTT - Property ID %i len %i\n", prop->id, prop->property_len);
486 PT_MQTT_WRITE_BYTE(conn, prop->id);
487 PT_MQTT_WRITE_BYTES(conn,
492 }
while(prop != NULL);
495 DBG(
"MQTT - No properties to write\n");
496 PT_MQTT_WRITE_BYTE(conn, 0);
504PT_THREAD(connect_pt(
struct pt *pt,
struct mqtt_connection *conn))
509 static struct mqtt_prop_list *will_props = MQTT_PROP_LIST_NONE;
510 if(conn->will.properties) {
511 will_props = (
struct mqtt_prop_list *)
list_head(conn->will.properties);
515 DBG(
"MQTT - Sending CONNECT message...\n");
518 conn->out_packet.fhdr = MQTT_FHDR_MSG_TYPE_CONNECT;
519 conn->out_packet.remaining_length = 0;
520 conn->out_packet.remaining_length += MQTT_CONNECT_VHDR_SIZE;
521 conn->out_packet.remaining_length += MQTT_STRING_LENGTH(&conn->client_id);
522#if (MQTT_PROTOCOL_VERSION > MQTT_PROTOCOL_VERSION_3_1) && MQTT_SRV_SUPPORTS_EMPTY_CLIENT_ID
524 if(MQTT_STRING_LENGTH(&conn->client_id) == 0) {
525 conn->out_packet.remaining_length += 2;
528 conn->out_packet.remaining_length += MQTT_STRING_LENGTH(&conn->credentials.username);
529 conn->out_packet.remaining_length += MQTT_STRING_LENGTH(&conn->credentials.password);
530 conn->out_packet.remaining_length += MQTT_STRING_LENGTH(&conn->will.topic);
531 conn->out_packet.remaining_length += MQTT_STRING_LENGTH(&conn->will.message);
535 conn->out_packet.remaining_length +=
536 conn->out_props ? (conn->out_props->properties_len + conn->out_props->properties_len_enc_bytes)
540 if(conn->connect_vhdr_flags & MQTT_VHDR_WILL_FLAG) {
541 conn->out_packet.remaining_length +=
542 will_props ? will_props->properties_len + will_props->properties_len_enc_bytes
547 mqtt_encode_var_byte_int(conn->out_packet.remaining_length_enc,
548 &conn->out_packet.remaining_length_enc_bytes,
549 conn->out_packet.remaining_length);
550 if(conn->out_packet.remaining_length_enc_bytes > 4) {
551 call_event(conn, MQTT_EVENT_PROTOCOL_ERROR, NULL);
552 PRINTF(
"MQTT - Error, remaining length > 4 bytes\n");
557 PT_MQTT_WRITE_BYTE(conn, conn->out_packet.fhdr);
558 PT_MQTT_WRITE_BYTES(conn,
559 conn->out_packet.remaining_length_enc,
560 conn->out_packet.remaining_length_enc_bytes);
561 PT_MQTT_WRITE_BYTE(conn, 0);
562 PT_MQTT_WRITE_BYTE(conn, strlen(MQTT_PROTOCOL_NAME));
563 PT_MQTT_WRITE_BYTES(conn, (uint8_t *)MQTT_PROTOCOL_NAME, strlen(MQTT_PROTOCOL_NAME));
564 PT_MQTT_WRITE_BYTE(conn, MQTT_PROTOCOL_VERSION);
565 PT_MQTT_WRITE_BYTE(conn, conn->connect_vhdr_flags);
566 PT_MQTT_WRITE_BYTE(conn, (conn->keep_alive >> 8));
567 PT_MQTT_WRITE_BYTE(conn, (conn->keep_alive & 0x00FF));
571 write_out_props(pt, conn, conn->out_props);
575 PT_MQTT_WRITE_BYTE(conn, conn->client_id.length >> 8);
576 PT_MQTT_WRITE_BYTE(conn, conn->client_id.length & 0x00FF);
577 PT_MQTT_WRITE_BYTES(conn, (uint8_t *)conn->client_id.string,
578 conn->client_id.length);
580 if(conn->connect_vhdr_flags & MQTT_VHDR_WILL_FLAG) {
583 DBG(
"MQTT - Writing will properties\n");
584 write_out_props(pt, conn, will_props);
586 PT_MQTT_WRITE_BYTE(conn, conn->will.topic.length >> 8);
587 PT_MQTT_WRITE_BYTE(conn, conn->will.topic.length & 0x00FF);
588 PT_MQTT_WRITE_BYTES(conn, (uint8_t *)conn->will.topic.string,
589 conn->will.topic.length);
590 PT_MQTT_WRITE_BYTE(conn, conn->will.message.length >> 8);
591 PT_MQTT_WRITE_BYTE(conn, conn->will.message.length & 0x00FF);
592 PT_MQTT_WRITE_BYTES(conn, (uint8_t *)conn->will.message.string,
593 conn->will.message.length);
594 DBG(
"MQTT - Setting will topic to '%s' %u bytes and message to '%s' %u bytes\n",
595 conn->will.topic.string,
596 conn->will.topic.length,
597 conn->will.message.string,
598 conn->will.message.length);
600 if(conn->connect_vhdr_flags & MQTT_VHDR_USERNAME_FLAG) {
601 PT_MQTT_WRITE_BYTE(conn, conn->credentials.username.length >> 8);
602 PT_MQTT_WRITE_BYTE(conn, conn->credentials.username.length & 0x00FF);
603 PT_MQTT_WRITE_BYTES(conn,
604 (uint8_t *)conn->credentials.username.string,
605 conn->credentials.username.length);
607 if(conn->connect_vhdr_flags & MQTT_VHDR_PASSWORD_FLAG) {
608 PT_MQTT_WRITE_BYTE(conn, conn->credentials.password.length >> 8);
609 PT_MQTT_WRITE_BYTE(conn, conn->credentials.password.length & 0x00FF);
610 PT_MQTT_WRITE_BYTES(conn,
611 (uint8_t *)conn->credentials.password.string,
612 conn->credentials.password.length);
616 send_out_buffer(conn);
617 conn->state = MQTT_CONN_STATE_CONNECTING_TO_BROKER;
619 timer_set(&conn->t, RESPONSE_WAIT_TIMEOUT);
622 reset_packet(&conn->in_packet);
623 PT_WAIT_UNTIL(pt, conn->out_packet.qos_state == MQTT_QOS_STATE_GOT_ACK ||
626 DBG(
"Timeout waiting for CONNACK\n");
634 reset_packet(&conn->in_packet);
636 DBG(
"MQTT - Done sending CONNECT\n");
639 DBG(
"MQTT - CONNECT message sent: \n");
641 for(i = 0; i < (conn->out_buffer_ptr - conn->out_buffer); i++) {
642 DBG(
"%02X ", conn->out_buffer[i]);
651PT_THREAD(disconnect_pt(
struct pt *pt,
struct mqtt_connection *conn))
655 PT_MQTT_WRITE_BYTE(conn, MQTT_FHDR_MSG_TYPE_DISCONNECT);
656 PT_MQTT_WRITE_BYTE(conn, 0);
660 write_out_props(pt, conn, conn->out_props);
663 send_out_buffer(conn);
677PT_THREAD(subscribe_pt(
struct pt *pt,
struct mqtt_connection *conn))
681 DBG(
"MQTT - Sending subscribe message! topic %s topic_length %i\n",
682 conn->out_packet.topic,
683 conn->out_packet.topic_length);
684 DBG(
"MQTT - Buffer space is %i \n",
685 &conn->out_buffer[MQTT_TCP_OUTPUT_BUFF_SIZE] - conn->out_buffer_ptr);
688 conn->out_packet.fhdr = MQTT_FHDR_MSG_TYPE_SUBSCRIBE | MQTT_FHDR_QOS_LEVEL_1;
689 conn->out_packet.remaining_length = MQTT_MID_SIZE +
690 MQTT_STRING_LEN_SIZE +
691 conn->out_packet.topic_length +
695 conn->out_packet.remaining_length +=
696 conn->out_props ? (conn->out_props->properties_len + conn->out_props->properties_len_enc_bytes)
700 mqtt_encode_var_byte_int(conn->out_packet.remaining_length_enc,
701 &conn->out_packet.remaining_length_enc_bytes,
702 conn->out_packet.remaining_length);
703 if(conn->out_packet.remaining_length_enc_bytes > 4) {
704 call_event(conn, MQTT_EVENT_PROTOCOL_ERROR, NULL);
705 PRINTF(
"MQTT - Error, remaining length > 4 bytes\n");
710 PT_MQTT_WRITE_BYTE(conn, conn->out_packet.fhdr);
711 PT_MQTT_WRITE_BYTES(conn,
712 conn->out_packet.remaining_length_enc,
713 conn->out_packet.remaining_length_enc_bytes);
715 PT_MQTT_WRITE_BYTE(conn, (conn->out_packet.mid >> 8));
716 PT_MQTT_WRITE_BYTE(conn, (conn->out_packet.mid & 0x00FF));
720 write_out_props(pt, conn, conn->out_props);
724 PT_MQTT_WRITE_BYTE(conn, (conn->out_packet.topic_length >> 8));
725 PT_MQTT_WRITE_BYTE(conn, (conn->out_packet.topic_length & 0x00FF));
726 PT_MQTT_WRITE_BYTES(conn, (uint8_t *)conn->out_packet.topic,
727 conn->out_packet.topic_length);
730 PT_MQTT_WRITE_BYTE(conn, conn->out_packet.sub_options);
732 PT_MQTT_WRITE_BYTE(conn, conn->out_packet.qos);
736 send_out_buffer(conn);
737 timer_set(&conn->t, RESPONSE_WAIT_TIMEOUT);
740 reset_packet(&conn->in_packet);
741 PT_WAIT_UNTIL(pt, conn->out_packet.qos_state == MQTT_QOS_STATE_GOT_ACK ||
745 DBG(
"Timeout waiting for SUBACK\n");
747 reset_packet(&conn->in_packet);
750 conn->out_queue_full = 0;
752 DBG(
"MQTT - Done in send_subscribe!\n");
758PT_THREAD(unsubscribe_pt(
struct pt *pt,
struct mqtt_connection *conn))
762 DBG(
"MQTT - Sending unsubscribe message on topic %s topic_length %i\n",
763 conn->out_packet.topic,
764 conn->out_packet.topic_length);
765 DBG(
"MQTT - Buffer space is %i \n",
766 &conn->out_buffer[MQTT_TCP_OUTPUT_BUFF_SIZE] - conn->out_buffer_ptr);
769 conn->out_packet.fhdr = MQTT_FHDR_MSG_TYPE_UNSUBSCRIBE |
770 MQTT_FHDR_QOS_LEVEL_1;
771 conn->out_packet.remaining_length = MQTT_MID_SIZE +
772 MQTT_STRING_LEN_SIZE +
773 conn->out_packet.topic_length;
776 conn->out_packet.remaining_length +=
777 conn->out_props ? (conn->out_props->properties_len + conn->out_props->properties_len_enc_bytes)
781 mqtt_encode_var_byte_int(conn->out_packet.remaining_length_enc,
782 &conn->out_packet.remaining_length_enc_bytes,
783 conn->out_packet.remaining_length);
784 if(conn->out_packet.remaining_length_enc_bytes > 4) {
785 call_event(conn, MQTT_EVENT_PROTOCOL_ERROR, NULL);
786 PRINTF(
"MQTT - Error, remaining length > 4 bytes\n");
791 PT_MQTT_WRITE_BYTE(conn, conn->out_packet.fhdr);
792 PT_MQTT_WRITE_BYTES(conn, (uint8_t *)conn->out_packet.remaining_length_enc,
793 conn->out_packet.remaining_length_enc_bytes);
796 PT_MQTT_WRITE_BYTE(conn, (conn->out_packet.mid >> 8));
797 PT_MQTT_WRITE_BYTE(conn, (conn->out_packet.mid & 0x00FF));
800 write_out_props(pt, conn, conn->out_props);
804 PT_MQTT_WRITE_BYTE(conn, (conn->out_packet.topic_length >> 8));
805 PT_MQTT_WRITE_BYTE(conn, (conn->out_packet.topic_length & 0x00FF));
806 PT_MQTT_WRITE_BYTES(conn, (uint8_t *)conn->out_packet.topic,
807 conn->out_packet.topic_length);
810 send_out_buffer(conn);
811 timer_set(&conn->t, RESPONSE_WAIT_TIMEOUT);
814 reset_packet(&conn->in_packet);
815 PT_WAIT_UNTIL(pt, conn->out_packet.qos_state == MQTT_QOS_STATE_GOT_ACK ||
819 DBG(
"Timeout waiting for UNSUBACK\n");
822 reset_packet(&conn->in_packet);
825 conn->out_queue_full = 0;
827 DBG(
"MQTT - Done writing subscribe message to out buffer!\n");
833PT_THREAD(publish_pt(
struct pt *pt,
struct mqtt_connection *conn))
837 DBG(
"MQTT - Sending publish message! topic %s topic_length %i\n",
838 conn->out_packet.topic,
839 conn->out_packet.topic_length);
840 DBG(
"MQTT - Buffer space is %i \n",
841 &conn->out_buffer[MQTT_TCP_OUTPUT_BUFF_SIZE] - conn->out_buffer_ptr);
844 conn->out_packet.fhdr = MQTT_FHDR_MSG_TYPE_PUBLISH |
845 conn->out_packet.qos << 1;
846 if(conn->out_packet.retain == MQTT_RETAIN_ON) {
847 conn->out_packet.fhdr |= MQTT_FHDR_RETAIN_FLAG;
849 conn->out_packet.remaining_length = MQTT_STRING_LEN_SIZE +
850 conn->out_packet.topic_length +
851 conn->out_packet.payload_size;
852 if(conn->out_packet.qos > MQTT_QOS_LEVEL_0) {
853 conn->out_packet.remaining_length += MQTT_MID_SIZE;
857 conn->out_packet.remaining_length +=
858 conn->out_props ? (conn->out_props->properties_len + conn->out_props->properties_len_enc_bytes)
862 mqtt_encode_var_byte_int(conn->out_packet.remaining_length_enc,
863 &conn->out_packet.remaining_length_enc_bytes,
864 conn->out_packet.remaining_length);
865 if(conn->out_packet.remaining_length_enc_bytes > 4) {
866 call_event(conn, MQTT_EVENT_PROTOCOL_ERROR, NULL);
867 PRINTF(
"MQTT - Error, remaining length > 4 bytes\n");
872 if(conn->out_packet.qos == MQTT_QOS_LEVEL_0) {
873 conn->out_packet.fhdr &= ~MQTT_FHDR_DUP_FLAG;
877 PT_MQTT_WRITE_BYTE(conn, conn->out_packet.fhdr);
878 PT_MQTT_WRITE_BYTES(conn, (uint8_t *)conn->out_packet.remaining_length_enc,
879 conn->out_packet.remaining_length_enc_bytes);
881 PT_MQTT_WRITE_BYTE(conn, (conn->out_packet.topic_length >> 8));
882 PT_MQTT_WRITE_BYTE(conn, (conn->out_packet.topic_length & 0x00FF));
883 PT_MQTT_WRITE_BYTES(conn, (uint8_t *)conn->out_packet.topic,
884 conn->out_packet.topic_length);
885 if(conn->out_packet.qos > MQTT_QOS_LEVEL_0) {
886 PT_MQTT_WRITE_BYTE(conn, (conn->out_packet.mid >> 8));
887 PT_MQTT_WRITE_BYTE(conn, (conn->out_packet.mid & 0x00FF));
892 write_out_props(pt, conn, conn->out_props);
896 PT_MQTT_WRITE_BYTES(conn,
897 conn->out_packet.payload,
898 conn->out_packet.payload_size);
900 send_out_buffer(conn);
901 timer_set(&conn->t, RESPONSE_WAIT_TIMEOUT);
909 if(conn->out_packet.qos == 0) {
910 process_post(conn->app_process, mqtt_update_event, NULL);
911 }
else if(conn->out_packet.qos == 1) {
913 reset_packet(&conn->in_packet);
914 PT_WAIT_UNTIL(pt, conn->out_packet.qos_state == MQTT_QOS_STATE_GOT_ACK ||
917 DBG(
"Timeout waiting for PUBACK\n");
919 if(conn->in_packet.mid != conn->out_packet.mid) {
920 DBG(
"MQTT - Warning, got PUBACK with none matching MID. Currently there "
921 "is no support for several concurrent PUBLISH messages.\n");
923 }
else if(conn->out_packet.qos == 2) {
924 DBG(
"MQTT - QoS not implemented yet.\n");
928 reset_packet(&conn->in_packet);
931 conn->out_queue_full = 0;
933 DBG(
"MQTT - Publish Enqueued\n");
939PT_THREAD(pingreq_pt(
struct pt *pt,
struct mqtt_connection *conn))
943 DBG(
"MQTT - Sending PINGREQ\n");
946 PT_MQTT_WRITE_BYTE(conn, MQTT_FHDR_MSG_TYPE_PINGREQ);
947 PT_MQTT_WRITE_BYTE(conn, 0);
949 send_out_buffer(conn);
952 conn->waiting_for_pingresp = 1;
955 reset_packet(&conn->in_packet);
956 timer_set(&conn->t, RESPONSE_WAIT_TIMEOUT);
960 reset_packet(&conn->in_packet);
962 conn->waiting_for_pingresp = 0;
969PT_THREAD(auth_pt(
struct pt *pt,
struct mqtt_connection *conn))
973 conn->out_packet.remaining_length +=
974 conn->out_props ? (conn->out_props->properties_len + conn->out_props->properties_len_enc_bytes)
977 mqtt_encode_var_byte_int(conn->out_packet.remaining_length_enc,
978 &conn->out_packet.remaining_length_enc_bytes,
979 conn->out_packet.remaining_length);
981 if(conn->out_packet.remaining_length_enc_bytes > 4) {
982 call_event(conn, MQTT_EVENT_PROTOCOL_ERROR, NULL);
983 PRINTF(
"MQTT - Error, remaining length > 4 bytes\n");
988 PT_MQTT_WRITE_BYTE(conn, conn->out_packet.fhdr);
989 PT_MQTT_WRITE_BYTES(conn, (uint8_t *)conn->out_packet.remaining_length_enc,
990 conn->out_packet.remaining_length_enc_bytes);
993 PT_MQTT_WRITE_BYTE(conn, conn->out_packet.auth_reason_code);
996 write_out_props(pt, conn, conn->out_props);
999 send_out_buffer(conn);
1008handle_connack(
struct mqtt_connection *conn)
1010 struct mqtt_connack_event connack_event;
1012 DBG(
"MQTT - Got CONNACK\n");
1014#if MQTT_PROTOCOL_VERSION <= MQTT_PROTOCOL_VERSION_3_1_1
1015 if(conn->in_packet.remaining_length != 2) {
1016 PRINTF(
"MQTT - CONNACK VHDR remaining length %i incorrect\n",
1017 conn->in_packet.remaining_length);
1021 abort_connection(conn);
1025 if(conn->in_packet.payload[1] != 0) {
1026 PRINTF(
"MQTT - Connection refused with Return Code %i\n",
1027 conn->in_packet.payload[1]);
1029 MQTT_EVENT_CONNECTION_REFUSED_ERROR,
1030 &conn->in_packet.payload[1]);
1031 abort_connection(conn);
1036#if MQTT_PROTOCOL_VERSION >= MQTT_PROTOCOL_VERSION_5
1042 if(conn->in_packet.remaining_length < 3) {
1043 PRINTF(
"MQTT - CONNACK VHDR remaining length %i incorrect\n",
1044 conn->in_packet.remaining_length);
1048 abort_connection(conn);
1053 conn->out_packet.qos_state = MQTT_QOS_STATE_GOT_ACK;
1055#if MQTT_PROTOCOL_VERSION >= MQTT_PROTOCOL_VERSION_3_1_1
1056 connack_event.session_present = conn->in_packet.payload[0] & MQTT_VHDR_CONNACK_SESSION_PRESENT;
1059#if MQTT_PROTOCOL_VERSION >= MQTT_PROTOCOL_VERSION_5
1060 mqtt_prop_parse_connack_props(conn);
1064 keep_alive_callback, conn);
1067 conn->state = MQTT_CONN_STATE_CONNECTED_TO_BROKER;
1068 call_event(conn, MQTT_EVENT_CONNECTED, &connack_event);
1072handle_pingresp(
struct mqtt_connection *conn)
1074 DBG(
"MQTT - Got PINGRESP\n");
1078handle_suback(
struct mqtt_connection *conn)
1080 struct mqtt_suback_event suback_event;
1082 DBG(
"MQTT - Got SUBACK\n");
1086 if(conn->in_packet.remaining_length > MQTT_MID_SIZE +
1087 MQTT_MAX_TOPICS_PER_SUBSCRIBE * MQTT_QOS_SIZE +
1088 conn->in_packet.properties_len + conn->in_packet.properties_enc_len) {
1090 if(conn->in_packet.remaining_length > MQTT_MID_SIZE +
1091 MQTT_MAX_TOPICS_PER_SUBSCRIBE * MQTT_QOS_SIZE) {
1093 DBG(
"MQTT - Error, SUBACK with > 1 topic, not supported.\n");
1096 conn->out_packet.qos_state = MQTT_QOS_STATE_GOT_ACK;
1098 suback_event.mid = conn->in_packet.mid;
1101 suback_event.success = 0;
1103 switch(conn->in_packet.payload_start[0]) {
1104 case MQTT_SUBACK_RET_FAIL:
1105 PRINTF(
"MQTT - Error, SUBSCRIBE failed with SUBACK return code '%x'", conn->in_packet.payload_start[0]);
1108 case MQTT_SUBACK_RET_QOS_0:
1109 case MQTT_SUBACK_RET_QOS_1:
1110 case MQTT_SUBACK_RET_QOS_2:
1111 suback_event.qos_level = conn->in_packet.payload_start[0] & 0x03;
1112 suback_event.success = 1;
1116 PRINTF(
"MQTT - Error, Unrecognised SUBACK return code '%x'", conn->in_packet.payload_start[0]);
1120 suback_event.return_code = conn->in_packet.payload_start[0];
1122 suback_event.qos_level = conn->in_packet.payload_start[0];
1125 if(conn->in_packet.mid != conn->out_packet.mid) {
1126 DBG(
"MQTT - Warning, got SUBACK with none matching MID. Currently there is"
1127 "no support for several concurrent SUBSCRIBE messages.\n");
1131 call_event(conn, MQTT_EVENT_SUBACK, &suback_event);
1135handle_unsuback(
struct mqtt_connection *conn)
1137 DBG(
"MQTT - Got UNSUBACK\n");
1139 conn->out_packet.qos_state = MQTT_QOS_STATE_GOT_ACK;
1141 if(conn->in_packet.mid != conn->out_packet.mid) {
1142 DBG(
"MQTT - Warning, got UNSUBACK with none matching MID. Currently there is"
1143 "no support for several concurrent UNSUBSCRIBE messages.\n");
1146 call_event(conn, MQTT_EVENT_UNSUBACK, &conn->in_packet.mid);
1150handle_puback(
struct mqtt_connection *conn)
1152 DBG(
"MQTT - Got PUBACK\n");
1154 conn->out_packet.qos_state = MQTT_QOS_STATE_GOT_ACK;
1156 call_event(conn, MQTT_EVENT_PUBACK, &conn->in_packet.mid);
1159static mqtt_pub_status_t
1160handle_publish(
struct mqtt_connection *conn)
1162 DBG(
"MQTT - Got PUBLISH, called once per manageable chunk of message.\n");
1163 DBG(
"MQTT - Handling publish on topic '%s'\n", conn->in_publish_msg.topic);
1165#if MQTT_PROTOCOL_VERSION >= MQTT_PROTOCOL_VERSION_3_1_1
1166 if(strlen(conn->in_publish_msg.topic) < conn->in_packet.topic_len) {
1167 DBG(
"NULL detected in received PUBLISH topic\n");
1173 return MQTT_PUBLISH_ERR;
1177 DBG(
"MQTT - This chunk is %i bytes\n", conn->in_publish_msg.payload_chunk_length);
1179 if(((conn->in_packet.fhdr & 0x09) >> 1) != 0) {
1180 PRINTF(
"MQTT - Error, got incoming PUBLISH with QoS > 0, not supported atm!\n");
1183 call_event(conn, MQTT_EVENT_PUBLISH, &conn->in_publish_msg);
1185 if(conn->in_publish_msg.first_chunk == 1) {
1186 conn->in_publish_msg.first_chunk = 0;
1190 if(conn->in_publish_msg.payload_left == 0) {
1195 DBG(
"MQTT - (handle_publish) resetting packet.\n");
1196 reset_packet(&conn->in_packet);
1199 return MQTT_PUBLISH_OK;
1203parse_publish_vhdr(
struct mqtt_connection *conn,
1205 const uint8_t *input_data_ptr,
1208 uint16_t copy_bytes;
1211 if(conn->in_packet.topic_len_received == 0) {
1212 if(!conn->in_packet.topic_len_msb_received) {
1213 conn->in_packet.topic_pos = 0;
1214 conn->in_packet.topic_len = (input_data_ptr[(*pos)++] << 8);
1215 conn->in_packet.byte_counter++;
1216 conn->in_packet.topic_len_msb_received = 1;
1217 if(*pos >= input_data_len) {
1221 conn->in_packet.topic_len |= input_data_ptr[(*pos)++];
1222 conn->in_packet.byte_counter++;
1223 conn->in_packet.topic_len_received = 1;
1226 if((uint32_t)conn->in_packet.topic_len + 2 >
1227 conn->in_packet.remaining_length) {
1228 PRINTF(
"MQTT - PUBLISH topic_len %u exceeds remaining_length %u\n",
1229 conn->in_packet.topic_len, conn->in_packet.remaining_length);
1230 call_event(conn, MQTT_EVENT_ERROR, NULL);
1231 abort_connection(conn);
1235 if(conn->in_packet.topic_len > MQTT_MAX_TOPIC_LENGTH) {
1236 PRINTF(
"MQTT - PUBLISH topic too long %u/%u, aborting\n",
1237 conn->in_packet.topic_len, MQTT_MAX_TOPIC_LENGTH);
1238 call_event(conn, MQTT_EVENT_ERROR, NULL);
1239 abort_connection(conn);
1242 DBG(
"MQTT - Read PUBLISH topic len %i\n", conn->in_packet.topic_len);
1246 if(conn->in_packet.topic_len_received == 1 &&
1247 conn->in_packet.topic_received == 0) {
1248 copy_bytes = MIN(conn->in_packet.topic_len - conn->in_packet.topic_pos,
1249 input_data_len - *pos);
1250 DBG(
"MQTT - topic_pos: %i copy_bytes: %i\n", conn->in_packet.topic_pos,
1252 memcpy(&conn->in_publish_msg.topic[conn->in_packet.topic_pos],
1253 &input_data_ptr[*pos],
1255 (*pos) += copy_bytes;
1256 conn->in_packet.byte_counter += copy_bytes;
1257 conn->in_packet.topic_pos += copy_bytes;
1259 if(conn->in_packet.topic_len - conn->in_packet.topic_pos == 0) {
1260 DBG(
"MQTT - Got topic '%s'", conn->in_publish_msg.topic);
1261 conn->in_packet.topic_received = 1;
1262 conn->in_publish_msg.topic[conn->in_packet.topic_pos] =
'\0';
1263 conn->in_publish_msg.payload_length =
1264 conn->in_packet.remaining_length - conn->in_packet.topic_len - 2;
1265 conn->in_publish_msg.payload_left = conn->in_publish_msg.payload_length;
1269 conn->in_publish_msg.first_chunk = 1;
1278handle_disconnect(
struct mqtt_connection *conn)
1280 DBG(
"MQTT - (handle_disconnect) Got DISCONNECT.\n");
1281 call_event(conn, MQTT_EVENT_DISCONNECTED, NULL);
1282 abort_connection(conn);
1286handle_auth(
struct mqtt_connection *conn)
1288 struct mqtt_prop_auth_event event;
1290 DBG(
"MQTT - (handle_auth) Got AUTH.\n");
1292 if((conn->in_packet.fhdr & 0x0F) != 0x0) {
1296 abort_connection(conn);
1301 if(conn->state == MQTT_CONN_STATE_CONNECTING_TO_BROKER &&
1302 (!conn->in_packet.has_reason_code ||
1303 conn->in_packet.reason_code != MQTT_VHDR_RC_CONTINUE_AUTH)) {
1304 DBG(
"MQTT - (handle_auth) Not reauth - Reason Code 0x18 expected!\n");
1307 mqtt_prop_parse_auth_props(conn, &event);
1308 call_event(conn, MQTT_EVENT_AUTH, &event);
1313parse_vhdr(
struct mqtt_connection *conn)
1315 conn->in_packet.payload_start = conn->in_packet.payload;
1318 switch(conn->in_packet.fhdr & 0xF0) {
1319 case MQTT_FHDR_MSG_TYPE_PUBACK:
1320 case MQTT_FHDR_MSG_TYPE_SUBACK:
1321 case MQTT_FHDR_MSG_TYPE_UNSUBACK:
1322 conn->in_packet.mid = (conn->in_packet.payload[0] << 8) |
1323 (conn->in_packet.payload[1]);
1324 conn->in_packet.payload_start += 2;
1338 switch(conn->in_packet.fhdr & 0xF0) {
1339 case MQTT_FHDR_MSG_TYPE_CONNACK:
1340 case MQTT_FHDR_MSG_TYPE_PUBACK:
1341 case MQTT_FHDR_MSG_TYPE_PUBREC:
1342 case MQTT_FHDR_MSG_TYPE_PUBREL:
1343 case MQTT_FHDR_MSG_TYPE_PUBCOMP:
1344 case MQTT_FHDR_MSG_TYPE_DISCONNECT:
1345 case MQTT_FHDR_MSG_TYPE_AUTH:
1346 conn->in_packet.reason_code = conn->in_packet.payload_start[0];
1347 conn->in_packet.has_reason_code = 1;
1348 conn->in_packet.payload_start += 1;
1352 conn->in_packet.has_reason_code = 0;
1356 if(!conn->in_packet.has_props) {
1357 mqtt_prop_decode_input_props(conn);
1364trim_publish_props(
struct mqtt_connection *conn)
1366 uint32_t prop_total = (uint32_t)conn->in_packet.properties_len +
1367 conn->in_packet.properties_enc_len;
1369 if(prop_total > conn->in_publish_msg.payload_chunk_length) {
1370 PRINTF(
"MQTT - Error, PUBLISH properties length exceeds buffered chunk\n");
1371 call_event(conn, MQTT_EVENT_ERROR, NULL);
1372 abort_connection(conn);
1376 conn->in_publish_msg.payload_chunk_length -= prop_total;
1377 conn->in_publish_msg.payload_chunk += prop_total;
1383tcp_input(
struct tcp_socket *s,
1385 const uint8_t *input_data_ptr,
1388 struct mqtt_connection *conn = ptr;
1390 uint32_t copy_bytes = 0;
1391 mqtt_pub_status_t pub_status;
1392 uint8_t remaining_length_bytes;
1394 if(input_data_len == 0) {
1398 DBG(
"tcp_input with %i bytes of data:\n", input_data_len);
1401 if(conn->in_packet.packet_received) {
1402 reset_packet(&conn->in_packet);
1406 if(!conn->in_packet.fhdr) {
1407 conn->in_packet.fhdr = input_data_ptr[pos++];
1408 conn->in_packet.byte_counter++;
1410 DBG(
"MQTT - Read VHDR '%02X'\n", conn->in_packet.fhdr);
1412 if(pos >= input_data_len) {
1423 if(!conn->in_packet.has_remaining_length) {
1424 remaining_length_bytes =
1425 mqtt_decode_var_byte_int(input_data_ptr, input_data_len, &pos,
1427 &conn->in_packet.remaining_length);
1429 if(remaining_length_bytes == 0) {
1430 call_event(conn, MQTT_EVENT_ERROR, NULL);
1434 DBG(
"MQTT - Finished reading remaining length byte\n");
1435 conn->in_packet.has_remaining_length = 1;
1444 if((conn->in_packet.remaining_length > MQTT_INPUT_BUFF_SIZE) &&
1445 (conn->in_packet.fhdr & 0xF0) != MQTT_FHDR_MSG_TYPE_PUBLISH) {
1446 uint32_t pkt_total = MQTT_FHDR_SIZE + conn->in_packet.remaining_length;
1447 uint32_t drain = MIN((uint32_t)input_data_len - pos,
1448 pkt_total - conn->in_packet.byte_counter);
1450 PRINTF(
"MQTT - Error, unsupported payload size for non-PUBLISH message\n");
1453 conn->in_packet.byte_counter += drain;
1454 if(conn->in_packet.byte_counter >= pkt_total) {
1455 conn->in_packet.packet_received = 1;
1456 if(pos < input_data_len) {
1469 while(conn->in_packet.byte_counter <
1470 (MQTT_FHDR_SIZE + conn->in_packet.remaining_length)) {
1472 if((conn->in_packet.fhdr & 0xF0) == MQTT_FHDR_MSG_TYPE_PUBLISH &&
1473 conn->in_packet.topic_received == 0) {
1474 if(parse_publish_vhdr(conn, &pos, input_data_ptr, input_data_len) < 0) {
1486 copy_bytes = MIN(input_data_len - pos,
1487 MQTT_INPUT_BUFF_SIZE - conn->in_packet.payload_pos);
1488 copy_bytes = MIN(copy_bytes,
1489 (MQTT_FHDR_SIZE + conn->in_packet.remaining_length) -
1490 conn->in_packet.byte_counter);
1491 DBG(
"- Copied %i payload bytes\n", copy_bytes);
1492 memcpy(&conn->in_packet.payload[conn->in_packet.payload_pos],
1493 &input_data_ptr[pos],
1495 conn->in_packet.byte_counter += copy_bytes;
1496 conn->in_packet.payload_pos += copy_bytes;
1501 DBG(
"MQTT - Copied bytes: \n");
1502 for(i = 0; i < copy_bytes; i++) {
1503 DBG(
"%02X ", conn->in_packet.payload[i]);
1509 if(MQTT_INPUT_BUFF_SIZE - conn->in_packet.payload_pos == 0) {
1510 conn->in_publish_msg.payload_chunk = conn->in_packet.payload;
1511 conn->in_publish_msg.payload_chunk_length = MQTT_INPUT_BUFF_SIZE;
1512 conn->in_publish_msg.payload_left -= MQTT_INPUT_BUFF_SIZE;
1515 if(!conn->in_packet.has_props) {
1516 mqtt_prop_decode_input_props(conn);
1519 if(conn->in_publish_msg.first_chunk) {
1520 if(trim_publish_props(conn) < 0) {
1526 pub_status = handle_publish(conn);
1528 conn->in_publish_msg.payload_chunk = conn->in_packet.payload;
1529 conn->in_packet.payload_pos = 0;
1531 if(pub_status != MQTT_PUBLISH_OK) {
1536 if(pos >= input_data_len &&
1537 (conn->in_packet.byte_counter < (MQTT_FHDR_SIZE + conn->in_packet.remaining_length))) {
1547 DBG(
"MQTT - Finished reading packet!\n");
1549 DBG(
"MQTT - total data was %i bytes of data. \n",
1550 (MQTT_FHDR_SIZE + conn->in_packet.remaining_length));
1553 if(conn->in_packet.has_reason_code &&
1554 conn->in_packet.reason_code >= MQTT_VHDR_RC_UNSPEC_ERR) {
1555 PRINTF(
"MQTT - Reason Code indicated error %i\n",
1556 conn->in_packet.reason_code);
1560 abort_connection(conn);
1566 switch(conn->in_packet.fhdr & 0xF0) {
1567 case MQTT_FHDR_MSG_TYPE_CONNACK:
1568 handle_connack(conn);
1570 case MQTT_FHDR_MSG_TYPE_PUBLISH:
1572 conn->in_publish_msg.payload_chunk = conn->in_packet.payload;
1573 conn->in_publish_msg.payload_chunk_length = conn->in_packet.payload_pos;
1574 conn->in_publish_msg.payload_left = 0;
1576 DBG(
"MQTT - First chunk? %i\n", conn->in_publish_msg.first_chunk);
1578 if(conn->in_publish_msg.first_chunk) {
1579 if(trim_publish_props(conn) < 0) {
1584 (void)handle_publish(conn);
1586 case MQTT_FHDR_MSG_TYPE_PUBACK:
1587 handle_puback(conn);
1589 case MQTT_FHDR_MSG_TYPE_SUBACK:
1590 handle_suback(conn);
1592 case MQTT_FHDR_MSG_TYPE_UNSUBACK:
1593 handle_unsuback(conn);
1595 case MQTT_FHDR_MSG_TYPE_PINGRESP:
1596 handle_pingresp(conn);
1600 case MQTT_FHDR_MSG_TYPE_PUBREC:
1601 case MQTT_FHDR_MSG_TYPE_PUBREL:
1602 case MQTT_FHDR_MSG_TYPE_PUBCOMP:
1603 call_event(conn, MQTT_EVENT_NOT_IMPLEMENTED_ERROR, NULL);
1604 PRINTF(
"MQTT - Got unhandled MQTT Message Type '%i'",
1605 (conn->in_packet.fhdr & 0xF0));
1608#if MQTT_PROTOCOL_VERSION >= MQTT_PROTOCOL_VERSION_5
1609 case MQTT_FHDR_MSG_TYPE_DISCONNECT:
1610 handle_disconnect(conn);
1613 case MQTT_FHDR_MSG_TYPE_AUTH:
1620 PRINTF(
"MQTT - Got MQTT Message Type '%i'", (conn->in_packet.fhdr & 0xF0));
1624 conn->in_packet.packet_received = 1;
1632 if(conn->state == MQTT_CONN_STATE_NOT_CONNECTED) {
1636 if(pos < input_data_len) {
1647tcp_event(
struct tcp_socket *s,
void *ptr, tcp_socket_event_t event)
1649 struct mqtt_connection *conn = ptr;
1655 case TCP_SOCKET_CLOSED:
1656 case TCP_SOCKET_TIMEDOUT:
1657 case TCP_SOCKET_ABORTED: {
1659 DBG(
"MQTT - Disconnected by tcp event %d\n", event);
1660 process_post(&mqtt_process, mqtt_abort_now_event, conn);
1661 conn->state = MQTT_CONN_STATE_NOT_CONNECTED;
1663 call_event(conn, MQTT_EVENT_DISCONNECTED, &event);
1664 abort_connection(conn);
1667 if(conn->auto_reconnect == 1) {
1672 case TCP_SOCKET_CONNECTED: {
1673 conn->state = MQTT_CONN_STATE_TCP_CONNECTED;
1674 conn->out_buffer_sent = 1;
1676 process_post(&mqtt_process, mqtt_do_connect_mqtt_event, conn);
1679 case TCP_SOCKET_DATA_SENT: {
1680 DBG(
"MQTT - Got TCP_DATA_SENT\n");
1682 if(conn->socket.output_data_len == 0) {
1683 conn->out_buffer_sent = 1;
1684 conn->out_buffer_ptr = conn->out_buffer;
1692 DBG(
"MQTT - TCP Event %d is currently not managed by the tcp event callback\n",
1700 static struct mqtt_connection *conn;
1707 if(ev == mqtt_abort_now_event) {
1708 DBG(
"MQTT - Abort\n");
1710 conn->state = MQTT_CONN_STATE_ABORT_IMMEDIATE;
1712 abort_connection(conn);
1714 if(ev == mqtt_do_connect_tcp_event) {
1716 DBG(
"MQTT - Got mqtt_do_connect_tcp_event!\n");
1719 if(ev == mqtt_do_connect_mqtt_event) {
1721 conn->socket.output_data_max_seg = conn->max_segment_size;
1722 DBG(
"MQTT - Got mqtt_do_connect_mqtt_event!\n");
1724 if(conn->out_buffer_sent == 1) {
1725 PT_INIT(&conn->out_proto_thread);
1726 while(connect_pt(&conn->out_proto_thread, conn) < PT_EXITED &&
1727 conn->state != MQTT_CONN_STATE_ABORT_IMMEDIATE) {
1728 PT_MQTT_WAIT_SEND();
1732 if(ev == mqtt_do_disconnect_mqtt_event) {
1734 DBG(
"MQTT - Got mqtt_do_disconnect_mqtt_event!\n");
1737 if(conn->state == MQTT_CONN_STATE_SENDING_MQTT_DISCONNECT) {
1738 if(conn->out_buffer_sent == 1) {
1739 PT_INIT(&conn->out_proto_thread);
1740 while(conn->state != MQTT_CONN_STATE_ABORT_IMMEDIATE &&
1741 disconnect_pt(&conn->out_proto_thread, conn) < PT_EXITED) {
1742 PT_MQTT_WAIT_SEND();
1744 abort_connection(conn);
1745 call_event(conn, MQTT_EVENT_DISCONNECTED, &ev);
1747 process_post(&mqtt_process, mqtt_do_disconnect_mqtt_event, conn);
1751 if(ev == mqtt_do_pingreq_event) {
1753 DBG(
"MQTT - Got mqtt_do_pingreq_event!\n");
1755 if(conn->out_buffer_sent == 1 &&
1756 conn->state == MQTT_CONN_STATE_CONNECTED_TO_BROKER) {
1757 PT_INIT(&conn->out_proto_thread);
1758 while(conn->state == MQTT_CONN_STATE_CONNECTED_TO_BROKER &&
1759 pingreq_pt(&conn->out_proto_thread, conn) < PT_EXITED) {
1760 PT_MQTT_WAIT_SEND();
1764 if(ev == mqtt_do_subscribe_event) {
1766 DBG(
"MQTT - Got mqtt_do_subscribe_mqtt_event!\n");
1768 if(conn->out_buffer_sent == 1 &&
1769 conn->state == MQTT_CONN_STATE_CONNECTED_TO_BROKER) {
1770 PT_INIT(&conn->out_proto_thread);
1771 while(conn->state == MQTT_CONN_STATE_CONNECTED_TO_BROKER &&
1772 subscribe_pt(&conn->out_proto_thread, conn) < PT_EXITED) {
1773 PT_MQTT_WAIT_SEND();
1777 if(ev == mqtt_do_unsubscribe_event) {
1779 DBG(
"MQTT - Got mqtt_do_unsubscribe_mqtt_event!\n");
1781 if(conn->out_buffer_sent == 1 &&
1782 conn->state == MQTT_CONN_STATE_CONNECTED_TO_BROKER) {
1783 PT_INIT(&conn->out_proto_thread);
1784 while(conn->state == MQTT_CONN_STATE_CONNECTED_TO_BROKER &&
1785 unsubscribe_pt(&conn->out_proto_thread, conn) < PT_EXITED) {
1786 PT_MQTT_WAIT_SEND();
1790 if(ev == mqtt_do_publish_event) {
1792 DBG(
"MQTT - Got mqtt_do_publish_mqtt_event!\n");
1794 if(conn->out_buffer_sent == 1 &&
1795 conn->state == MQTT_CONN_STATE_CONNECTED_TO_BROKER) {
1796 PT_INIT(&conn->out_proto_thread);
1797 while(conn->state == MQTT_CONN_STATE_CONNECTED_TO_BROKER &&
1798 publish_pt(&conn->out_proto_thread, conn) < PT_EXITED) {
1799 PT_MQTT_WAIT_SEND();
1804 if(ev == mqtt_do_auth_event) {
1806 DBG(
"MQTT - Got mqtt_do_auth_event!\n");
1808 if(conn->out_buffer_sent == 1) {
1809 PT_INIT(&conn->out_proto_thread);
1810 while(auth_pt(&conn->out_proto_thread, conn) < PT_EXITED) {
1811 PT_MQTT_WAIT_SEND();
1816 conn->out_props = NULL;
1825 static uint8_t inited = 0;
1828 mqtt_event_min = mqtt_do_connect_tcp_event;
1838 mqtt_event_max = mqtt_abort_now_event;
1853 uint16_t max_segment_size)
1855#if MQTT_31 || !MQTT_SRV_SUPPORTS_EMPTY_CLIENT_ID
1856 if(strlen(client_id) < 1) {
1857 return MQTT_STATUS_INVALID_ARGS_ERROR;
1862 memset(conn, 0,
sizeof(
struct mqtt_connection));
1865 conn->srv_feature_en = -1;
1867 string_to_mqtt_string(&conn->client_id, client_id);
1868 conn->event_callback = event_callback;
1869 conn->app_process = app_process;
1870 conn->auto_reconnect = 1;
1871 conn->max_segment_size = max_segment_size;
1873 reset_defaults(conn);
1879 DBG(
"MQTT - Registered successfully\n");
1881 return MQTT_STATUS_OK;
1891 uint16_t keep_alive,
1893 uint8_t clean_session,
1894 struct mqtt_prop_list *prop_list)
1896 uint8_t clean_session)
1899 uip_ip6addr_t ip6addr;
1904 if(conn->state > MQTT_CONN_STATE_NOT_CONNECTED) {
1905 return MQTT_STATUS_OK;
1908 conn->server_host = host;
1909 conn->keep_alive = keep_alive;
1910 conn->server_port = port;
1911 conn->out_buffer_ptr = conn->out_buffer;
1912 conn->out_packet.qos_state = MQTT_QOS_STATE_NO_ACK;
1915 if(clean_session || (conn->client_id.length == 0)) {
1916 conn->connect_vhdr_flags |= MQTT_VHDR_CLEAN_SESSION_FLAG;
1920 if(uiplib_ip6addrconv(host, &ip6addr) == 0) {
1921 return MQTT_STATUS_ERROR;
1932 conn->out_props = prop_list;
1935 process_post(&mqtt_process, mqtt_do_connect_tcp_event, conn);
1937 return MQTT_STATUS_OK;
1943 struct mqtt_prop_list *prop_list)
1948 if(conn->state != MQTT_CONN_STATE_CONNECTED_TO_BROKER) {
1952 conn->state = MQTT_CONN_STATE_SENDING_MQTT_DISCONNECT;
1955 conn->out_props = prop_list;
1958 process_post(&mqtt_process, mqtt_do_disconnect_mqtt_event, conn);
1964 mqtt_qos_level_t qos_level,
1965 mqtt_nl_en_t nl, mqtt_rap_en_t rap,
1966 mqtt_retain_handling_t ret_handling,
1967 struct mqtt_prop_list *prop_list)
1969 mqtt_qos_level_t qos_level)
1972 if(conn->state != MQTT_CONN_STATE_CONNECTED_TO_BROKER) {
1973 return MQTT_STATUS_NOT_CONNECTED_ERROR;
1976 DBG(
"MQTT - Call to mqtt_subscribe...\n");
1979 if(conn->out_queue_full) {
1980 DBG(
"MQTT - Not accepted!\n");
1981 return MQTT_STATUS_OUT_QUEUE_FULL;
1983 conn->out_queue_full = 1;
1984 DBG(
"MQTT - Accepted!\n");
1986 conn->out_packet.mid = INCREMENT_MID(conn);
1987 conn->out_packet.topic = topic;
1988 conn->out_packet.topic_length = strlen(topic);
1989 conn->out_packet.qos_state = MQTT_QOS_STATE_NO_ACK;
1992 *mid = conn->out_packet.mid;
1996 conn->out_packet.sub_options = 0x00;
1997 conn->out_packet.sub_options |= qos_level & MQTT_SUB_OPTION_QOS;
1998 conn->out_packet.sub_options |= nl & MQTT_SUB_OPTION_NL;
1999 conn->out_packet.sub_options |= rap & MQTT_SUB_OPTION_RAP;
2000 conn->out_packet.sub_options |= ret_handling & MQTT_SUB_OPTION_RETAIN_HANDLING;
2002 conn->out_packet.qos = qos_level;
2006 conn->out_props = prop_list;
2009 process_post(&mqtt_process, mqtt_do_subscribe_event, conn);
2010 return MQTT_STATUS_OK;
2017 struct mqtt_prop_list *prop_list)
2022 if(conn->state != MQTT_CONN_STATE_CONNECTED_TO_BROKER) {
2023 return MQTT_STATUS_NOT_CONNECTED_ERROR;
2026 DBG(
"MQTT - Call to mqtt_unsubscribe...\n");
2028 if(conn->out_queue_full) {
2029 DBG(
"MQTT - Not accepted!\n");
2030 return MQTT_STATUS_OUT_QUEUE_FULL;
2032 conn->out_queue_full = 1;
2033 DBG(
"MQTT - Accepted!\n");
2035 conn->out_packet.mid = INCREMENT_MID(conn);
2036 conn->out_packet.topic = topic;
2037 conn->out_packet.topic_length = strlen(topic);
2038 conn->out_packet.qos_state = MQTT_QOS_STATE_NO_ACK;
2041 *mid = conn->out_packet.mid;
2045 conn->out_props = prop_list;
2048 process_post(&mqtt_process, mqtt_do_unsubscribe_event, conn);
2049 return MQTT_STATUS_OK;
2054 uint8_t *payload, uint32_t payload_size,
2055 mqtt_qos_level_t qos_level,
2057 mqtt_retain_t retain,
2058 uint8_t topic_alias, mqtt_topic_alias_en_t topic_alias_en,
2059 struct mqtt_prop_list *prop_list)
2061 mqtt_retain_t retain)
2064 if(conn->state != MQTT_CONN_STATE_CONNECTED_TO_BROKER) {
2065 return MQTT_STATUS_NOT_CONNECTED_ERROR;
2068 DBG(
"MQTT - Call to mqtt_publish...\n");
2071 if(conn->out_queue_full) {
2072 DBG(
"MQTT - Not accepted!\n");
2073 return MQTT_STATUS_OUT_QUEUE_FULL;
2075 conn->out_queue_full = 1;
2076 DBG(
"MQTT - Accepted!\n");
2078 conn->out_packet.mid = INCREMENT_MID(conn);
2079 conn->out_packet.retain = retain;
2081 if(topic_alias_en == MQTT_TOPIC_ALIAS_ON) {
2082 conn->out_packet.topic =
"";
2083 conn->out_packet.topic_length = 0;
2084 conn->out_packet.topic_alias = topic_alias;
2085 if(topic_alias == 0) {
2086 DBG(
"MQTT - Error, a topic alias of 0 is not permitted! It won't be sent.\n");
2089 conn->out_packet.topic = topic;
2090 conn->out_packet.topic_length = strlen(topic);
2091 conn->out_packet.topic_alias = 0;
2094 conn->out_packet.topic = topic;
2095 conn->out_packet.topic_length = strlen(topic);
2097 conn->out_packet.payload = payload;
2098 conn->out_packet.payload_size = payload_size;
2099 conn->out_packet.qos = qos_level;
2100 conn->out_packet.qos_state = MQTT_QOS_STATE_NO_ACK;
2103 *mid = conn->out_packet.mid;
2107 conn->out_props = prop_list;
2110 process_post(&mqtt_process, mqtt_do_publish_event, conn);
2111 return MQTT_STATUS_OK;
2119 string_to_mqtt_string(&conn->credentials.username, username);
2120 string_to_mqtt_string(&conn->credentials.password, password);
2123 if(username != NULL) {
2124 conn->connect_vhdr_flags |= MQTT_VHDR_USERNAME_FLAG;
2126 conn->connect_vhdr_flags &= ~MQTT_VHDR_USERNAME_FLAG;
2128 if(password != NULL) {
2129 conn->connect_vhdr_flags |= MQTT_VHDR_PASSWORD_FLAG;
2131 conn->connect_vhdr_flags &= ~MQTT_VHDR_PASSWORD_FLAG;
2138 mqtt_qos_level_t qos,
struct mqtt_prop_list *will_props)
2140 mqtt_qos_level_t qos)
2144 string_to_mqtt_string(&conn->will.topic, topic);
2145 string_to_mqtt_string(&conn->will.message, message);
2148 conn->will.qos = qos;
2151 conn->connect_vhdr_flags |= MQTT_VHDR_WILL_FLAG |
2152 MQTT_VHDR_WILL_RETAIN_FLAG;
2155 conn->will.properties = (
list_t)will_props;
2170 mqtt_auth_type_t auth_type,
2171 struct mqtt_prop_list *prop_list)
2173 DBG(
"MQTT - Call to mqtt_auth...\n");
2175 conn->out_packet.fhdr = MQTT_FHDR_MSG_TYPE_AUTH;
2176 conn->out_packet.remaining_length = 1;
2177 conn->out_packet.auth_reason_code = MQTT_VHDR_RC_CONTINUE_AUTH + auth_type;
2179 conn->out_props = prop_list;
2182 return MQTT_STATUS_OK;
Default definitions of C compiler quirk work-arounds.
Header file for the callback timer.
#define CLOCK_SECOND
A second, measured in system clock time.
void ctimer_stop(struct ctimer *c)
Stop a pending callback timer.
static void ctimer_set(struct ctimer *c, clock_time_t t, void(*f)(void *), void *ptr)
Set a callback timer.
void ctimer_restart(struct ctimer *c)
Restart a callback timer from the current point in time.
static void list_init(list_t list)
Initialize a list.
#define LIST(name)
Declare a linked list.
static void * list_item_next(const void *item)
Get the next item following this item.
void list_add(list_t list, void *item)
Add an item at the end of a list.
void ** list_t
The linked list type.
static void * list_head(const_list_t list)
Get a pointer to the first element of a list.
mqtt_status_t mqtt_auth(struct mqtt_connection *conn, mqtt_auth_type_t auth_type, struct mqtt_prop_list *prop_list)
Send authentication message (MQTTv5-only).
mqtt_status_t mqtt_connect(struct mqtt_connection *conn, char *host, uint16_t port, uint16_t keep_alive, uint8_t clean_session, struct mqtt_prop_list *prop_list)
Connects to a MQTT broker.
mqtt_status_t mqtt_register(struct mqtt_connection *conn, struct process *app_process, char *client_id, mqtt_event_callback_t event_callback, uint16_t max_segment_size)
Initializes the MQTT engine.
mqtt_status_t mqtt_unsubscribe(struct mqtt_connection *conn, uint16_t *mid, char *topic, struct mqtt_prop_list *prop_list)
Unsubscribes from a MQTT topic.
void(* mqtt_event_callback_t)(struct mqtt_connection *m, mqtt_event_t event, void *data)
MQTT event callback function.
mqtt_status_t mqtt_subscribe(struct mqtt_connection *conn, uint16_t *mid, char *topic, mqtt_qos_level_t qos_level, mqtt_nl_en_t nl, mqtt_rap_en_t rap, mqtt_retain_handling_t ret_handling, struct mqtt_prop_list *prop_list)
Subscribes to a MQTT topic.
void mqtt_disconnect(struct mqtt_connection *conn, struct mqtt_prop_list *prop_list)
Disconnects from a MQTT broker.
mqtt_status_t mqtt_publish(struct mqtt_connection *conn, uint16_t *mid, char *topic, uint8_t *payload, uint32_t payload_size, mqtt_qos_level_t qos_level, mqtt_retain_t retain, uint8_t topic_alias, mqtt_topic_alias_en_t topic_alias_en, struct mqtt_prop_list *prop_list)
Publish to a MQTT topic.
mqtt_event_t
MQTT engine events.
void mqtt_set_username_password(struct mqtt_connection *conn, char *username, char *password)
Set the user name and password for a MQTT client.
void mqtt_set_last_will(struct mqtt_connection *conn, char *topic, char *message, mqtt_qos_level_t qos, struct mqtt_prop_list *will_props)
Set the last will topic and message for a MQTT client.
#define PROCESS(name, strname)
Declare a process.
#define PROCESS_WAIT_EVENT()
Wait for an event to be posted to the process.
int process_post(struct process *p, process_event_t ev, process_data_t data)
Post an asynchronous event.
process_event_t process_alloc_event(void)
Allocate a global event number.
#define PROCESS_BEGIN()
Define the beginning of a process.
#define PROCESS_END()
Define the end of a process.
void process_start(struct process *p, process_data_t data)
Start a process.
#define PROCESS_THREAD(name, ev, data)
Define the body of a process.
#define PT_BEGIN(pt)
Declare the start of a protothread inside the C function implementing the protothread.
#define PT_THREAD(name_args)
Declaration of a protothread.
#define PT_END(pt)
Declare the end of a protothread.
#define PT_EXIT(pt)
Exit the protothread.
#define PT_WAIT_UNTIL(pt, condition)
Block and wait until condition is true.
#define PT_INIT(pt)
Initialize a protothread.
void timer_set(struct timer *t, clock_time_t interval)
Set a timer.
bool timer_expired(struct timer *t)
Check if a timer has expired.
#define uip_ipaddr_copy(dest, src)
Copy an IP address from one place to another.
Header file for the LED HAL.
Linked list manipulation routines.
Header file for the Contiki MQTT engine.
Protothreads implementation.
Header file for generating non-cryptographic random numbers.
Header file for IPv6-related data structures.
static uip_ipaddr_t ipaddr
Pointer to prefix information option in uip_buf.
Header file for the uIP TCP/IP stack.