@cryptotaxi247 / netdata-1 / commits / d066c655e

Detect missing ACLK MQTT packet acknowledgents (#20711)

Enhance MQTT buffer size configuration and improve packet timeout handling - Introduced configurable buffer size based on client type. - Implemented packet timeout monitoring using JudyL. - Adjusted transaction buffer handling for MQTT messages. - Minor formatting and readability improvements throughout the code. - Add function call to check for packet timeouts in message handling loop.

Stelios Fragkakis committed Jul 23, 2025 at 08:47 UTC d066c655e417a6b140db2da8f2f9854c1d4bce3d
2 files changed +260 -191
src/aclk/aclk.c
+5 -3
@@ -21,6 +21,9 @@
21 #include <fcntl.h>
22 #endif
23
24 +#define MQTT_DEFAULT_MAX_BUF_SIZE (25 * 1024 * 1024)
25 +#define MQTT_PARENT_MAX_BUF_SIZE (128 * 1024 * 1024)
26 +
27 int aclk_pubacks_per_conn = 0; // How many PubAcks we got since MQTT conn est.
28 int aclk_rcvd_cloud_msgs = 0;
29 int aclk_connection_counter = 0;
@@ -864,9 +867,8 @@ void *aclk_main(void *ptr)
867 #endif
868
869 // Enable MQTT buffer growth if necessary
867 - // e.g. old cloud architecture clients with huge nodes
868 - // that send JSON payloads of 10 MB as single messages
869 - mqtt_wss_set_max_buf_size(mqttwss_client, 25*1024*1024);
870 + size_t max_buf_size = netdata_conf_is_parent() ? MQTT_PARENT_MAX_BUF_SIZE : MQTT_DEFAULT_MAX_BUF_SIZE;
871 + mqtt_wss_set_max_buf_size(mqttwss_client, max_buf_size);
872
873 // Keep reconnecting and talking until our time has come
874 // and the Grim Reaper (exit_initiated) calls
src/aclk/mqtt_websockets/mqtt_ng.c
+255 -188
@@ -12,6 +12,8 @@ void pulse_aclk_sent_message_acked(usec_t usec, size_t len);
12 #include "mqtt_ng.h"
13 #include "aclk_mqtt_workers.h"
14
15 +
16 +#define PACKET_ACK_TIMEOUT_SECS (60)
17 #define SMALL_STRING_DONT_FRAGMENT_LIMIT 128
18
19 #define LOCK_HDR_BUFFER(buffer) spinlock_lock(&((buffer)->spinlock))
@@ -238,6 +240,11 @@ struct mqtt_ng_client {
240 struct mqtt_ng_stats stats;
241 SPINLOCK stats_spinlock;
242
243 + struct {
244 + SPINLOCK spinlock;
245 + Pvoid_t JudyL;
246 + } pending_packets;
247 +
248 struct topic_aliases_data tx_topic_aliases;
249 c_rhash rx_aliases;
250
@@ -418,7 +425,7 @@ static void buffer_frag_free_data(struct buffer_fragment *frag)
425 }
426 }
427
421 -#define HEADER_BUFFER_SIZE 1024*1024
428 +#define HEADER_BUFFER_SIZE (1024*1024)
429 #define GROWTH_FACTOR 1.25
430
431 #define BUFFER_BYTES_USED(buf) ((size_t)((buf)->tail - (buf)->data))
@@ -527,7 +534,7 @@ static int transaction_buffer_grow(struct transaction_buffer *buf, float rate, s
534 if (buf->sending_frag != &ping_frag)
535 buf->sending_frag = NULL;
536
530 - buf->hdr_buffer.size *= rate;
537 + buf->hdr_buffer.size = (size_t)((float)buf->hdr_buffer.size * rate);
538 if (buf->hdr_buffer.size > max)
539 buf->hdr_buffer.size = max;
540
@@ -563,9 +570,11 @@ static void transaction_buffer_destroy(struct transaction_buffer *to_init)
570 // Creates transaction
571 // saves state of buffer before any operation was done
572 // allowing for rollback if things go wrong
566 -#define transaction_buffer_transaction_start(buf) \
567 - { LOCK_HDR_BUFFER(buf); \
568 - memcpy(&(buf)->state_backup, &(buf)->hdr_buffer, sizeof((buf)->hdr_buffer)); }
573 +#define transaction_buffer_transaction_start(buf) \
574 + { \
575 + LOCK_HDR_BUFFER(buf); \
576 + memcpy(&(buf)->state_backup, &(buf)->hdr_buffer, sizeof((buf)->hdr_buffer)); \
577 + }
578
579 #define transaction_buffer_transaction_commit(buf) UNLOCK_HDR_BUFFER(buf);
580
@@ -611,6 +620,8 @@ struct mqtt_ng_client *mqtt_ng_init(struct mqtt_ng_init *settings)
620 client->puback_callback = settings->puback_callback;
621 client->connack_callback = settings->connack_callback;
622 client->msg_callback = settings->msg_callback;
623 + spinlock_init(&client->pending_packets.spinlock);
624 + client->pending_packets.JudyL = NULL;
625
626 return client;
627 }
@@ -650,7 +661,9 @@ void mqtt_ng_destroy(struct mqtt_ng_client *client)
661
662 mqtt_ng_destroy_tx_alias_hash(client->tx_topic_aliases.stoi_dict);
663 mqtt_ng_destroy_rx_alias_hash(client->rx_aliases);
653 -
664 + spinlock_lock(&client->pending_packets.spinlock);
665 + (void) JudyLFreeArray(&client->pending_packets.JudyL, PJE0);
666 + spinlock_unlock(&client->pending_packets.spinlock);
667 freez(client);
668 }
669
@@ -736,9 +749,12 @@ static size_t mqtt_ng_connect_size(struct mqtt_auth_properties *auth,
749 frag = buffer_new_frag(buf, (flags)); } \
750 if(frag==NULL) { on_fail; }}
751
739 -#define CHECK_BYTES_AVAILABLE(buf, needed, fail) \
740 - { if (BUFFER_BYTES_AVAILABLE(buf) < (size_t)needed) { \
741 - fail; } }
752 +#define CHECK_BYTES_AVAILABLE(buf, needed, fail) \
753 + { \
754 + if (BUFFER_BYTES_AVAILABLE(buf) < (size_t)needed) { \
755 + fail; \
756 + } \
757 + }
758
759 #define DATA_ADVANCE(buf, bytes, frag) { size_t b = (bytes); (buf)->tail += b; (frag)->len += b; }
760
@@ -746,15 +762,14 @@ static size_t mqtt_ng_connect_size(struct mqtt_auth_properties *auth,
762 #define WRITE_POS(frag) (&(frag->data[frag->len]))
763
764 // [MQTT-1.5.2] Two Byte Integer
749 -#define PACK_2B_INT(buffer, integer, frag) { \
750 - uint16_t temp = htobe16((integer)); \
751 - memcpy(WRITE_POS(frag), &temp, sizeof(uint16_t)); \
752 - DATA_ADVANCE(buffer, sizeof(uint16_t), frag); \
753 -}
754 -// #define PACK_2B_INT(buffer, integer, frag) { *(uint16_t *)WRITE_POS(frag) = htobe16((integer));
755 -// DATA_ADVANCE(buffer, sizeof(uint16_t), frag); }
765 +#define PACK_2B_INT(buffer, integer, frag) \
766 + { \
767 + uint16_t temp = htobe16((integer)); \
768 + memcpy(WRITE_POS(frag), &temp, sizeof(uint16_t)); \
769 + DATA_ADVANCE(buffer, sizeof(uint16_t), frag); \
770 + }
771
757 -static int _optimized_add(struct header_buffer *buf, void *data, size_t data_len, free_fnc_t data_free_fnc, struct buffer_fragment **frag)
772 +static int optimized_add(struct header_buffer *buf, void *data, size_t data_len, free_fnc_t data_free_fnc, struct buffer_fragment **frag)
773 {
774 if (data_len > SMALL_STRING_DONT_FRAGMENT_LIMIT) {
775 buffer_frag_flag_t flags = BUFFER_FRAG_DATA_EXTERNAL;
@@ -773,40 +788,64 @@ static int _optimized_add(struct header_buffer *buf, void *data, size_t data_len
788 } else if (data_len) {
789 // if the data are small dont bother creating new fragments
790 // store in buffer directly
776 - CHECK_BYTES_AVAILABLE(buf, data_len, return 1);
791 + CHECK_BYTES_AVAILABLE(buf, data_len, return 1)
792 memcpy(buf->tail, data, data_len);
778 - DATA_ADVANCE(buf, data_len, *frag);
793 + DATA_ADVANCE(buf, data_len, *frag)
794 }
795 return 0;
796 }
797
798 +static void remove_packet_from_timeout_monitor_list(struct mqtt_ng_client *client, uint16_t packet_id)
799 +{
800 + spinlock_lock(&client->pending_packets.spinlock);
801 + (void) JudyLDel(&client->pending_packets.JudyL, (Word_t) packet_id, PJE0);
802 + spinlock_unlock(&client->pending_packets.spinlock);
803 +}
804 +
805 +static void add_packet_to_timeout_monitor_list(struct mqtt_ng_client *client, uint16_t packet_id)
806 +{
807 + spinlock_lock(&client->pending_packets.spinlock);
808 + time_t now = now_realtime_sec();
809 + // Add it to the JudyL array
810 + time_t *Pvalue = (time_t *) JudyLIns(&client->pending_packets.JudyL, (Word_t) packet_id, PJE0);
811 + if (!Pvalue || Pvalue == PJERR) {
812 + nd_log(NDLS_DAEMON, NDLP_ERR, "Error inserting packet_id (%" PRIu16 ") into JudyL array.", packet_id);
813 + spinlock_unlock(&client->pending_packets.spinlock);
814 + return;
815 + }
816 + *Pvalue = now + PACKET_ACK_TIMEOUT_SECS;
817 + spinlock_unlock(&client->pending_packets.spinlock);
818 +}
819 +
820 #define TRY_GENERATE_MESSAGE(generator_function, ...) \
784 - int rc = generator_function(&client->main_buffer, ##__VA_ARGS__); \
785 - if (rc == MQTT_NG_MSGGEN_BUFFER_OOM) { \
786 - LOCK_HDR_BUFFER(&client->main_buffer); \
787 - transaction_buffer_garbage_collect((&client->main_buffer)); \
788 - UNLOCK_HDR_BUFFER(&client->main_buffer); \
789 - rc = generator_function(&client->main_buffer, ##__VA_ARGS__); \
790 - if (rc == MQTT_NG_MSGGEN_BUFFER_OOM && client->max_mem_bytes) { \
821 + ({ \
822 + int _rc = generator_function(&client->main_buffer, ##__VA_ARGS__); \
823 + if (_rc == MQTT_NG_MSGGEN_BUFFER_OOM) { \
824 LOCK_HDR_BUFFER(&client->main_buffer); \
792 - transaction_buffer_grow((&client->main_buffer), GROWTH_FACTOR, client->max_mem_bytes); \
825 + transaction_buffer_garbage_collect((&client->main_buffer)); \
826 UNLOCK_HDR_BUFFER(&client->main_buffer); \
794 - rc = generator_function(&client->main_buffer, ##__VA_ARGS__); \
827 + _rc = generator_function(&client->main_buffer, ##__VA_ARGS__); \
828 + if (_rc == MQTT_NG_MSGGEN_BUFFER_OOM && client->max_mem_bytes) { \
829 + LOCK_HDR_BUFFER(&client->main_buffer); \
830 + transaction_buffer_grow((&client->main_buffer), GROWTH_FACTOR, client->max_mem_bytes); \
831 + UNLOCK_HDR_BUFFER(&client->main_buffer); \
832 + _rc = generator_function(&client->main_buffer, ##__VA_ARGS__); \
833 + } \
834 + if (_rc == MQTT_NG_MSGGEN_BUFFER_OOM) \
835 + nd_log( \
836 + NDLS_DAEMON, \
837 + NDLP_ERR, \
838 + "%s failed to generate message due to insufficient buffer space (line %d)", \
839 + __FUNCTION__, \
840 + __LINE__); \
841 } \
796 - if (rc == MQTT_NG_MSGGEN_BUFFER_OOM) \
797 - nd_log( \
798 - NDLS_DAEMON, \
799 - NDLP_ERR, \
800 - "%s failed to generate message due to insufficient buffer space (line %d)", \
801 - __FUNCTION__, \
802 - __LINE__); \
803 - } \
804 - if (rc == MQTT_NG_MSGGEN_OK) { \
805 - spinlock_lock(&client->stats_spinlock); \
806 - client->stats.tx_messages_queued++; \
807 - spinlock_unlock(&client->stats_spinlock); \
808 - } \
809 - return rc;
842 + if (_rc == MQTT_NG_MSGGEN_OK) { \
843 + spinlock_lock(&client->stats_spinlock); \
844 + client->stats.tx_messages_queued++; \
845 + spinlock_unlock(&client->stats_spinlock); \
846 + } \
847 + _rc; \
848 + })
849
850 mqtt_msg_data mqtt_ng_generate_connect(struct transaction_buffer *trx_buf,
851 struct mqtt_auth_properties *auth,
@@ -849,7 +888,7 @@ mqtt_msg_data mqtt_ng_generate_connect(struct transaction_buffer *trx_buf,
888 }
889
890 // >> START THE RODEO <<
852 - transaction_buffer_transaction_start(trx_buf);
891 + transaction_buffer_transaction_start(trx_buf)
892
893 // Calculate the resulting message size sans fixed MQTT header
894 size_t size = mqtt_ng_connect_size(auth, lwt);
@@ -858,19 +897,19 @@ mqtt_msg_data mqtt_ng_generate_connect(struct transaction_buffer *trx_buf,
897 struct buffer_fragment *frag = NULL;
898 mqtt_msg_data ret = NULL;
899
861 - BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, BUFFER_FRAG_MQTT_PACKET_HEAD, frag, goto fail_rollback );
900 + BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, BUFFER_FRAG_MQTT_PACKET_HEAD, frag, goto fail_rollback)
901 ret = frag;
902
903 // MQTT Fixed Header
904 size_t needed_bytes = 1 /* Packet type */ + MQTT_VARSIZE_INT_BYTES(size) + sizeof(mqtt_protocol_name_frag) + 1 /* CONNECT FLAGS */ + 2 /* keepalive */ + 1 /* Properties TODO now fixed 0*/;
866 - CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, needed_bytes, goto fail_rollback);
905 + CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, needed_bytes, goto fail_rollback)
906
907 *WRITE_POS(frag) = MQTT_CPT_CONNECT << 4;
869 - DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
870 - DATA_ADVANCE(&trx_buf->hdr_buffer, uint32_to_mqtt_vbi(size, WRITE_POS(frag)), frag);
908 + DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag)
909 + DATA_ADVANCE(&trx_buf->hdr_buffer, uint32_to_mqtt_vbi(size, WRITE_POS(frag)), frag)
910
911 memcpy(WRITE_POS(frag), mqtt_protocol_name_frag, sizeof(mqtt_protocol_name_frag));
873 - DATA_ADVANCE(&trx_buf->hdr_buffer, sizeof(mqtt_protocol_name_frag), frag);
912 + DATA_ADVANCE(&trx_buf->hdr_buffer, sizeof(mqtt_protocol_name_frag), frag)
913
914 // [MQTT-3.1.2.3] Connect flags
915 unsigned char *connect_flags = WRITE_POS(frag);
@@ -890,66 +929,67 @@ mqtt_msg_data mqtt_ng_generate_connect(struct transaction_buffer *trx_buf,
929 if (clean_start)
930 *connect_flags |= MQTT_CONNECT_FLAG_CLEAN_START;
931
893 - DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
932 + DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag)
933
895 - PACK_2B_INT(&trx_buf->hdr_buffer, keep_alive, frag);
934 + PACK_2B_INT(&trx_buf->hdr_buffer, keep_alive, frag)
935
936 // TODO Property Length [MQTT-3.1.3.2.1] temporary fixed to 3 (one property topic alias max)
898 - DATA_ADVANCE(&trx_buf->hdr_buffer, uint32_to_mqtt_vbi(3, WRITE_POS(frag)), frag);
937 + DATA_ADVANCE(&trx_buf->hdr_buffer, uint32_to_mqtt_vbi(3, WRITE_POS(frag)), frag)
938 *WRITE_POS(frag) = MQTT_PROP_TOPIC_ALIAS_MAX;
900 - DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
939 + DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag)
940
902 - PACK_2B_INT(&trx_buf->hdr_buffer, 65535, frag);
941 + PACK_2B_INT(&trx_buf->hdr_buffer, 65535, frag)
942
943 // [MQTT-3.1.3.1] Client identifier
905 - CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, 2, goto fail_rollback);
906 - PACK_2B_INT(&trx_buf->hdr_buffer, strlen(auth->client_id), frag);
907 - if (_optimized_add(&trx_buf->hdr_buffer, auth->client_id, strlen(auth->client_id), auth->client_id_free, &frag))
944 + CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, 2, goto fail_rollback)
945 + PACK_2B_INT(&trx_buf->hdr_buffer, strlen(auth->client_id), frag)
946 + if (optimized_add(&trx_buf->hdr_buffer, auth->client_id, strlen(auth->client_id), auth->client_id_free, &frag))
947 goto fail_rollback;
948
949 if (lwt != NULL) {
950 // Will Properties [MQTT-3.1.3.2]
951 // TODO for now fixed 0
913 - BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, 0, frag, goto fail_rollback);
914 - CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, 1, goto fail_rollback);
952 + BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, 0, frag, goto fail_rollback)
953 + CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, 1, goto fail_rollback)
954 *WRITE_POS(frag) = 0;
916 - DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
955 + DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag)
956
957 // Will Topic [MQTT-3.1.3.3]
919 - CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, 2, goto fail_rollback);
920 - PACK_2B_INT(&trx_buf->hdr_buffer, strlen(lwt->will_topic), frag);
921 - if (_optimized_add(&trx_buf->hdr_buffer, lwt->will_topic, strlen(lwt->will_topic), lwt->will_topic_free, &frag))
958 + CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, 2, goto fail_rollback)
959 + PACK_2B_INT(&trx_buf->hdr_buffer, strlen(lwt->will_topic), frag)
960 + if (optimized_add(&trx_buf->hdr_buffer, lwt->will_topic, strlen(lwt->will_topic), lwt->will_topic_free, &frag))
961 goto fail_rollback;
962
963 // Will Payload [MQTT-3.1.3.4]
964 if (lwt->will_message_size) {
926 - BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, 0, frag, goto fail_rollback);
927 - CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, 2, goto fail_rollback);
928 - PACK_2B_INT(&trx_buf->hdr_buffer, lwt->will_message_size, frag);
929 - if (_optimized_add(&trx_buf->hdr_buffer, lwt->will_message, lwt->will_message_size, lwt->will_topic_free, &frag))
965 + BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, 0, frag, goto fail_rollback)
966 + CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, 2, goto fail_rollback)
967 + PACK_2B_INT(&trx_buf->hdr_buffer, lwt->will_message_size, frag)
968 + if (optimized_add(
969 + &trx_buf->hdr_buffer, lwt->will_message, lwt->will_message_size, lwt->will_topic_free, &frag))
970 goto fail_rollback;
971 }
972 }
973
974 // [MQTT-3.1.3.5]
975 if (auth->username) {
936 - BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, 0, frag, goto fail_rollback);
937 - CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, 2, goto fail_rollback);
938 - PACK_2B_INT(&trx_buf->hdr_buffer, strlen(auth->username), frag);
939 - if (_optimized_add(&trx_buf->hdr_buffer, auth->username, strlen(auth->username), auth->username_free, &frag))
976 + BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, 0, frag, goto fail_rollback)
977 + CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, 2, goto fail_rollback)
978 + PACK_2B_INT(&trx_buf->hdr_buffer, strlen(auth->username), frag)
979 + if (optimized_add(&trx_buf->hdr_buffer, auth->username, strlen(auth->username), auth->username_free, &frag))
980 goto fail_rollback;
981 }
982
983 // [MQTT-3.1.3.6]
984 if (auth->password) {
945 - BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, 0, frag, goto fail_rollback);
946 - CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, 2, goto fail_rollback);
947 - PACK_2B_INT(&trx_buf->hdr_buffer, strlen(auth->password), frag);
948 - if (_optimized_add(&trx_buf->hdr_buffer, auth->password, strlen(auth->password), auth->password_free, &frag))
985 + BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, 0, frag, goto fail_rollback)
986 + CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, 2, goto fail_rollback)
987 + PACK_2B_INT(&trx_buf->hdr_buffer, strlen(auth->password), frag)
988 + if (optimized_add(&trx_buf->hdr_buffer, auth->password, strlen(auth->password), auth->password_free, &frag))
989 goto fail_rollback;
990 }
991 trx_buf->hdr_buffer.tail_frag->flags |= BUFFER_FRAG_MQTT_PACKET_TAIL;
952 - transaction_buffer_transaction_commit(trx_buf);
992 + transaction_buffer_transaction_commit(trx_buf)
993 return ret;
994 fail_rollback:
995 transaction_buffer_transaction_rollback(trx_buf, ret);
@@ -1035,7 +1075,7 @@ int mqtt_ng_generate_publish(struct transaction_buffer *trx_buf,
1075 uint16_t topic_alias)
1076 {
1077 // >> START THE RODEO <<
1038 - transaction_buffer_transaction_start(trx_buf);
1078 + transaction_buffer_transaction_start(trx_buf)
1079
1080 // Calculate the resulting message size sans fixed MQTT header
1081 size_t size = mqtt_ng_publish_size(topic, msg_len, topic_alias);
@@ -1044,7 +1084,7 @@ int mqtt_ng_generate_publish(struct transaction_buffer *trx_buf,
1084 struct buffer_fragment *frag = NULL;
1085 mqtt_msg_data mqtt_msg = NULL;
1086
1047 - BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, BUFFER_FRAG_MQTT_PACKET_HEAD, frag, goto fail_rollback );
1087 + BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, BUFFER_FRAG_MQTT_PACKET_HEAD, frag, goto fail_rollback )
1088 // in case of QOS 0 we can garbage collect immediatelly after sending
1089 uint8_t qos = (publish_flags >> 1) & 0x03;
1090 if (!qos)
@@ -1053,35 +1093,35 @@ int mqtt_ng_generate_publish(struct transaction_buffer *trx_buf,
1093
1094 // MQTT Fixed Header
1095 size_t needed_bytes = 1 /* Packet type */ + MQTT_VARSIZE_INT_BYTES(size) + size - msg_len;
1056 - CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, needed_bytes, goto fail_rollback);
1096 + CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, needed_bytes, goto fail_rollback)
1097
1098 *WRITE_POS(frag) = (MQTT_CPT_PUBLISH << 4) | (publish_flags & 0xF);
1059 - DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
1060 - DATA_ADVANCE(&trx_buf->hdr_buffer, uint32_to_mqtt_vbi(size, WRITE_POS(frag)), frag);
1099 + DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag)
1100 + DATA_ADVANCE(&trx_buf->hdr_buffer, uint32_to_mqtt_vbi(size, WRITE_POS(frag)), frag)
1101
1102 // MQTT Variable Header
1103 // [MQTT-3.3.2.1]
1064 - PACK_2B_INT(&trx_buf->hdr_buffer, topic == NULL ? 0 : strlen(topic), frag);
1104 + PACK_2B_INT(&trx_buf->hdr_buffer, topic == NULL ? 0 : strlen(topic), frag)
1105 if (topic != NULL) {
1066 - if (_optimized_add(&trx_buf->hdr_buffer, topic, strlen(topic), topic_free, &frag))
1106 + if (optimized_add(&trx_buf->hdr_buffer, topic, strlen(topic), topic_free, &frag))
1107 goto fail_rollback;
1068 - BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, 0, frag, goto fail_rollback);
1108 + BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, 0, frag, goto fail_rollback)
1109 }
1110
1111 // [MQTT-3.3.2.2]
1112 mqtt_msg->packet_id = get_unused_packet_id();
1113 *packet_id = mqtt_msg->packet_id;
1074 - PACK_2B_INT(&trx_buf->hdr_buffer, mqtt_msg->packet_id, frag);
1114 + PACK_2B_INT(&trx_buf->hdr_buffer, mqtt_msg->packet_id, frag)
1115
1116 // [MQTT-3.3.2.3.1] TODO Property Length for now fixed 0
1117 *WRITE_POS(frag) = topic_alias ? 3 : 0;
1078 - DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
1118 + DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag)
1119
1120 if(topic_alias) {
1121 *WRITE_POS(frag) = MQTT_PROP_TOPIC_ALIAS;
1082 - DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
1122 + DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag)
1123
1084 - PACK_2B_INT(&trx_buf->hdr_buffer, topic_alias, frag);
1124 + PACK_2B_INT(&trx_buf->hdr_buffer, topic_alias, frag)
1125 }
1126
1127 if( (frag = buffer_new_frag(&trx_buf->hdr_buffer, BUFFER_FRAG_DATA_EXTERNAL)) == NULL )
@@ -1093,13 +1133,76 @@ int mqtt_ng_generate_publish(struct transaction_buffer *trx_buf,
1133 trx_buf->hdr_buffer.tail_frag->flags |= BUFFER_FRAG_MQTT_PACKET_TAIL;
1134 if (!qos)
1135 trx_buf->hdr_buffer.tail_frag->flags |= BUFFER_FRAG_GARBAGE_COLLECT_ON_SEND;
1096 - transaction_buffer_transaction_commit(trx_buf);
1136 + transaction_buffer_transaction_commit(trx_buf)
1137 return MQTT_NG_MSGGEN_OK;
1138 fail_rollback:
1139 transaction_buffer_transaction_rollback(trx_buf, mqtt_msg);
1140 return MQTT_NG_MSGGEN_BUFFER_OOM;
1141 }
1142
1143 +static void mark_message_for_gc(struct buffer_fragment *frag)
1144 +{
1145 + while (frag) {
1146 + frag->flags |= BUFFER_FRAG_GARBAGE_COLLECT;
1147 + buffer_frag_free_data(frag);
1148 + if (frag->flags & BUFFER_FRAG_MQTT_PACKET_TAIL)
1149 + return;
1150 + frag = frag->next;
1151 + }
1152 +}
1153 +
1154 +static int mark_packet_acked(struct mqtt_ng_client *client, uint16_t packet_id)
1155 +{
1156 + size_t reclaimable = 0;
1157 + LOCK_HDR_BUFFER(&client->main_buffer);
1158 + struct buffer_fragment *frag = BUFFER_FIRST_FRAG(&client->main_buffer.hdr_buffer);
1159 + while (frag) {
1160 + if ( (frag->flags & BUFFER_FRAG_MQTT_PACKET_HEAD) && frag->packet_id == packet_id) {
1161 + if (!frag->sent) {
1162 + nd_log(NDLS_DAEMON, NDLP_ERR, "Received packet_id (%" PRIu16 ") belongs to MQTT packet which was not yet sent!", packet_id);
1163 + UNLOCK_HDR_BUFFER(&client->main_buffer);
1164 + return 1;
1165 + }
1166 + pulse_aclk_sent_message_acked(frag->sent_monotonic_ut, frag->len);
1167 + mark_message_for_gc(frag);
1168 +
1169 + size_t used = BUFFER_BYTES_USED(&client->main_buffer.hdr_buffer);
1170 + if (reclaimable >= (used / 4))
1171 + transaction_buffer_garbage_collect(&client->main_buffer);
1172 +
1173 + UNLOCK_HDR_BUFFER(&client->main_buffer);
1174 + remove_packet_from_timeout_monitor_list(client, packet_id);
1175 + return 0;
1176 + }
1177 +
1178 + if(frag_is_marked_for_gc(frag))
1179 + reclaimable += FRAG_SIZE_IN_BUFFER(frag);
1180 +
1181 + frag = frag->next;
1182 + }
1183 + nd_log(NDLS_DAEMON, NDLP_ERR, "Received packet_id (%" PRIu16 ") is unknown!", packet_id);
1184 + UNLOCK_HDR_BUFFER(&client->main_buffer);
1185 + return 1;
1186 +}
1187 +
1188 +static void check_packet_monitor_list_for_timeouts(struct mqtt_ng_client *client)
1189 +{
1190 + spinlock_lock(&client->pending_packets.spinlock);
1191 + bool first_then_next = true;
1192 + time_t *Pvalue;
1193 + Word_t packet_id = 0;
1194 + time_t now = now_realtime_sec();
1195 + while ((Pvalue = (time_t *) JudyLFirstThenNext(client->pending_packets.JudyL, &packet_id, &first_then_next))) {
1196 + time_t expire_time = *Pvalue;
1197 + if (now >= expire_time) {
1198 + spinlock_unlock(&client->pending_packets.spinlock);
1199 + (void) mark_packet_acked(client, (uint16_t) packet_id);
1200 + spinlock_lock(&client->pending_packets.spinlock);
1201 + }
1202 + }
1203 + spinlock_unlock(&client->pending_packets.spinlock);
1204 +}
1205 +
1206 #define PUBLISH_SP_SIZE 64
1207 int mqtt_ng_publish(struct mqtt_ng_client *client,
1208 char *topic,
@@ -1133,7 +1236,15 @@ int mqtt_ng_publish(struct mqtt_ng_client *client,
1236 return MQTT_NG_MSGGEN_MSG_TOO_BIG;
1237 }
1238
1136 - TRY_GENERATE_MESSAGE(mqtt_ng_generate_publish, topic, topic_free, msg, msg_free, msg_len, publish_flags, packet_id, topic_id);
1239 + int rc = TRY_GENERATE_MESSAGE(mqtt_ng_generate_publish, topic, topic_free, msg, msg_free, msg_len, publish_flags, packet_id, topic_id);
1240 + if (rc == MQTT_NG_MSGGEN_BUFFER_OOM) {
1241 + check_packet_monitor_list_for_timeouts(client);
1242 + rc = TRY_GENERATE_MESSAGE(mqtt_ng_generate_publish, topic, topic_free, msg, msg_free, msg_len, publish_flags, packet_id, topic_id);
1243 + }
1244 +
1245 + if (rc == MQTT_NG_MSGGEN_OK)
1246 + add_packet_to_timeout_monitor_list(client, *packet_id);
1247 + return rc;
1248 }
1249
1250 static size_t mqtt_ng_subscribe_size(struct mqtt_sub *subs, size_t sub_count)
@@ -1150,7 +1261,7 @@ static size_t mqtt_ng_subscribe_size(struct mqtt_sub *subs, size_t sub_count)
1261 int mqtt_ng_generate_subscribe(struct transaction_buffer *trx_buf, struct mqtt_sub *subs, size_t sub_count)
1262 {
1263 // >> START THE RODEO <<
1153 - transaction_buffer_transaction_start(trx_buf);
1264 + transaction_buffer_transaction_start(trx_buf)
1265
1266 // Calculate the resulting message size sans fixed MQTT header
1267 size_t size = mqtt_ng_subscribe_size(subs, sub_count);
@@ -1159,38 +1270,38 @@ int mqtt_ng_generate_subscribe(struct transaction_buffer *trx_buf, struct mqtt_s
1270 struct buffer_fragment *frag = NULL;
1271 mqtt_msg_data ret = NULL;
1272
1162 - BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, BUFFER_FRAG_MQTT_PACKET_HEAD, frag, goto fail_rollback);
1273 + BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, BUFFER_FRAG_MQTT_PACKET_HEAD, frag, goto fail_rollback)
1274 ret = frag;
1275
1276 // MQTT Fixed Header
1277 size_t needed_bytes = 1 /* Packet type */ + MQTT_VARSIZE_INT_BYTES(size) + 3 /*Packet ID + Property Length*/;
1167 - CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, needed_bytes, goto fail_rollback);
1278 + CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, needed_bytes, goto fail_rollback)
1279
1280 *WRITE_POS(frag) = (MQTT_CPT_SUBSCRIBE << 4) | 0x2 /* [MQTT-3.8.1-1] */;
1170 - DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
1171 - DATA_ADVANCE(&trx_buf->hdr_buffer, uint32_to_mqtt_vbi(size, WRITE_POS(frag)), frag);
1281 + DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag)
1282 + DATA_ADVANCE(&trx_buf->hdr_buffer, uint32_to_mqtt_vbi(size, WRITE_POS(frag)), frag)
1283
1284 // MQTT Variable Header
1285 // [MQTT-3.8.2] PacketID
1286 ret->packet_id = get_unused_packet_id();
1176 - PACK_2B_INT(&trx_buf->hdr_buffer, ret->packet_id, frag);
1287 + PACK_2B_INT(&trx_buf->hdr_buffer, ret->packet_id, frag)
1288
1289 // [MQTT-3.8.2.1.1] Property Length // TODO for now fixed 0
1290 *WRITE_POS(frag) = 0;
1180 - DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
1291 + DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag)
1292
1293 for (size_t i = 0; i < sub_count; i++) {
1183 - BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, 0, frag, goto fail_rollback);
1184 - PACK_2B_INT(&trx_buf->hdr_buffer, strlen(subs[i].topic), frag);
1185 - if (_optimized_add(&trx_buf->hdr_buffer, subs[i].topic, strlen(subs[i].topic), subs[i].topic_free, &frag))
1294 + BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, 0, frag, goto fail_rollback)
1295 + PACK_2B_INT(&trx_buf->hdr_buffer, strlen(subs[i].topic), frag)
1296 + if (optimized_add(&trx_buf->hdr_buffer, subs[i].topic, strlen(subs[i].topic), subs[i].topic_free, &frag))
1297 goto fail_rollback;
1187 - BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, 0, frag, goto fail_rollback);
1298 + BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, 0, frag, goto fail_rollback)
1299 *WRITE_POS(frag) = subs[i].options;
1189 - DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
1300 + DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag)
1301 }
1302
1303 trx_buf->hdr_buffer.tail_frag->flags |= BUFFER_FRAG_MQTT_PACKET_TAIL;
1193 - transaction_buffer_transaction_commit(trx_buf);
1304 + transaction_buffer_transaction_commit(trx_buf)
1305 return MQTT_NG_MSGGEN_OK;
1306 fail_rollback:
1307 transaction_buffer_transaction_rollback(trx_buf, ret);
@@ -1199,13 +1310,13 @@ fail_rollback:
1310
1311 int mqtt_ng_subscribe(struct mqtt_ng_client *client, struct mqtt_sub *subs, size_t sub_count)
1312 {
1202 - TRY_GENERATE_MESSAGE(mqtt_ng_generate_subscribe, subs, sub_count);
1313 + return TRY_GENERATE_MESSAGE(mqtt_ng_generate_subscribe, subs, sub_count);
1314 }
1315
1316 int mqtt_ng_generate_disconnect(struct transaction_buffer *trx_buf, uint8_t reason_code)
1317 {
1318 // >> START THE RODEO <<
1208 - transaction_buffer_transaction_start(trx_buf);
1319 + transaction_buffer_transaction_start(trx_buf)
1320
1321 // Calculate the resulting message size sans fixed MQTT header
1322 size_t size = reason_code ? 1 : 0;
@@ -1214,26 +1325,26 @@ int mqtt_ng_generate_disconnect(struct transaction_buffer *trx_buf, uint8_t reas
1325 struct buffer_fragment *frag = NULL;
1326 mqtt_msg_data ret = NULL;
1327
1217 - BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, BUFFER_FRAG_MQTT_PACKET_HEAD, frag, goto fail_rollback);
1328 + BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, BUFFER_FRAG_MQTT_PACKET_HEAD, frag, goto fail_rollback)
1329 ret = frag;
1330
1331 // MQTT Fixed Header
1332 size_t needed_bytes = 1 /* Packet type */ + MQTT_VARSIZE_INT_BYTES(size) + (reason_code ? 1 : 0);
1222 - CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, needed_bytes, goto fail_rollback);
1333 + CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, needed_bytes, goto fail_rollback)
1334
1335 *WRITE_POS(frag) = MQTT_CPT_DISCONNECT << 4;
1225 - DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
1226 - DATA_ADVANCE(&trx_buf->hdr_buffer, uint32_to_mqtt_vbi(size, WRITE_POS(frag)), frag);
1336 + DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag)
1337 + DATA_ADVANCE(&trx_buf->hdr_buffer, uint32_to_mqtt_vbi(size, WRITE_POS(frag)), frag)
1338
1339 if (reason_code) {
1340 // MQTT Variable Header
1341 // [MQTT-3.14.2.1] PacketID
1342 *WRITE_POS(frag) = reason_code;
1232 - DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
1343 + DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag)
1344 }
1345
1346 trx_buf->hdr_buffer.tail_frag->flags |= BUFFER_FRAG_MQTT_PACKET_TAIL;
1236 - transaction_buffer_transaction_commit(trx_buf);
1347 + transaction_buffer_transaction_commit(trx_buf)
1348 return MQTT_NG_MSGGEN_OK;
1349 fail_rollback:
1350 transaction_buffer_transaction_rollback(trx_buf, ret);
@@ -1242,13 +1353,13 @@ fail_rollback:
1353
1354 int mqtt_ng_disconnect(struct mqtt_ng_client *client, uint8_t reason_code)
1355 {
1245 - TRY_GENERATE_MESSAGE(mqtt_ng_generate_disconnect, reason_code);
1356 + return TRY_GENERATE_MESSAGE(mqtt_ng_generate_disconnect, reason_code);
1357 }
1358
1359 static int mqtt_generate_puback(struct transaction_buffer *trx_buf, uint16_t packet_id, uint8_t reason_code)
1360 {
1361 // >> START THE RODEO <<
1251 - transaction_buffer_transaction_start(trx_buf);
1362 + transaction_buffer_transaction_start(trx_buf)
1363
1364 // Calculate the resulting message size sans fixed MQTT header
1365 size_t size = 2 /* Packet ID */ + (reason_code ? 1 : 0) /* reason code */;
@@ -1256,28 +1367,28 @@ static int mqtt_generate_puback(struct transaction_buffer *trx_buf, uint16_t pac
1367 // Start generating the message
1368 struct buffer_fragment *frag = NULL;
1369
1259 - BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, BUFFER_FRAG_MQTT_PACKET_HEAD | BUFFER_FRAG_GARBAGE_COLLECT_ON_SEND, frag, goto fail_rollback);
1370 + BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, BUFFER_FRAG_MQTT_PACKET_HEAD | BUFFER_FRAG_GARBAGE_COLLECT_ON_SEND, frag, goto fail_rollback)
1371
1372 // MQTT Fixed Header
1373 size_t needed_bytes = 1 /* Packet type */ + MQTT_VARSIZE_INT_BYTES(size) + size;
1263 - CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, needed_bytes, goto fail_rollback);
1374 + CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, needed_bytes, goto fail_rollback)
1375
1376 *WRITE_POS(frag) = MQTT_CPT_PUBACK << 4;
1266 - DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
1267 - DATA_ADVANCE(&trx_buf->hdr_buffer, uint32_to_mqtt_vbi(size, WRITE_POS(frag)), frag);
1377 + DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag)
1378 + DATA_ADVANCE(&trx_buf->hdr_buffer, uint32_to_mqtt_vbi(size, WRITE_POS(frag)), frag)
1379
1380 // MQTT Variable Header
1270 - PACK_2B_INT(&trx_buf->hdr_buffer, packet_id, frag);
1381 + PACK_2B_INT(&trx_buf->hdr_buffer, packet_id, frag)
1382
1383 if (reason_code) {
1384 // MQTT Variable Header
1385 // [MQTT-3.14.2.1] PacketID
1386 *WRITE_POS(frag) = reason_code;
1276 - DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
1387 + DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag)
1388 }
1389
1390 trx_buf->hdr_buffer.tail_frag->flags |= BUFFER_FRAG_MQTT_PACKET_TAIL;
1280 - transaction_buffer_transaction_commit(trx_buf);
1391 + transaction_buffer_transaction_commit(trx_buf)
1392 return MQTT_NG_MSGGEN_OK;
1393 fail_rollback:
1394 transaction_buffer_transaction_rollback(trx_buf, frag);
@@ -1286,7 +1397,7 @@ fail_rollback:
1397
1398 static int mqtt_ng_puback(struct mqtt_ng_client *client, uint16_t packet_id, uint8_t reason_code)
1399 {
1289 - TRY_GENERATE_MESSAGE(mqtt_generate_puback, packet_id, reason_code);
1400 + return TRY_GENERATE_MESSAGE(mqtt_generate_puback, packet_id, reason_code);
1401 }
1402
1403 int mqtt_ng_ping(struct mqtt_ng_client *client)
@@ -1300,10 +1411,10 @@ int mqtt_ng_ping(struct mqtt_ng_client *client)
1411 #define MQTT_NG_CLIENT_PARSE_DONE 0x12
1412 #define MQTT_NG_CLIENT_WANT_WRITE 0x13
1413 #define MQTT_NG_CLIENT_OK_CALL_AGAIN 0
1303 -#define MQTT_NG_CLIENT_PROTOCOL_ERROR -1
1304 -#define MQTT_NG_CLIENT_SERVER_RETURNED_ERROR -2
1305 -#define MQTT_NG_CLIENT_NOT_IMPL_YET -3
1306 -#define MQTT_NG_CLIENT_INTERNAL_ERROR -5
1414 +#define MQTT_NG_CLIENT_PROTOCOL_ERROR (-1)
1415 +#define MQTT_NG_CLIENT_SERVER_RETURNED_ERROR (-2)
1416 +#define MQTT_NG_CLIENT_NOT_IMPL_YET (-3)
1417 +#define MQTT_NG_CLIENT_INTERNAL_ERROR (-5)
1418
1419 #define BUF_READ_CHECK_AT_LEAST(buf, x) \
1420 if (rbuf_bytes_available(buf) < (x)) \
@@ -1318,7 +1429,7 @@ static int vbi_parser_parse(struct mqtt_vbi_parser_ctx *ctx, rbuf_t data)
1429 return MQTT_NG_CLIENT_PROTOCOL_ERROR;
1430 }
1431 if (!ctx->bytes || ctx->data[ctx->bytes-1] & MQTT_VBI_CONTINUATION_FLAG) {
1321 - BUF_READ_CHECK_AT_LEAST(data, 1);
1432 + BUF_READ_CHECK_AT_LEAST(data, 1)
1433 ctx->bytes++;
1434 rbuf_pop(data, &ctx->data[ctx->bytes-1], 1);
1435 if ( ctx->data[ctx->bytes-1] & MQTT_VBI_CONTINUATION_FLAG )
@@ -1429,7 +1540,7 @@ static int parse_properties_array(struct mqtt_properties_parser_ctx *ctx, rbuf_t
1540 }
1541 return rc;
1542 case PROPERTY_CREATE:
1432 - BUF_READ_CHECK_AT_LEAST(data, 1);
1543 + BUF_READ_CHECK_AT_LEAST(data, 1)
1544 struct mqtt_property *prop = callocz(1, sizeof(struct mqtt_property));
1545 if (ctx->head == NULL) {
1546 ctx->head = prop;
@@ -1471,7 +1582,7 @@ static int parse_properties_array(struct mqtt_properties_parser_ctx *ctx, rbuf_t
1582 }
1583 break;
1584 case PROPERTY_TYPE_STR_BIN_LEN:
1474 - BUF_READ_CHECK_AT_LEAST(data, sizeof(uint16_t));
1585 + BUF_READ_CHECK_AT_LEAST(data, sizeof(uint16_t))
1586 rbuf_pop(data, (char*)&ctx->tail->bindata_len, sizeof(uint16_t));
1587 ctx->tail->bindata_len = be16toh(ctx->tail->bindata_len);
1588 ctx->bytes_consumed += 2;
@@ -1489,7 +1600,7 @@ static int parse_properties_array(struct mqtt_properties_parser_ctx *ctx, rbuf_t
1600 }
1601 break;
1602 case PROPERTY_TYPE_STR:
1492 - BUF_READ_CHECK_AT_LEAST(data, ctx->tail->bindata_len);
1603 + BUF_READ_CHECK_AT_LEAST(data, ctx->tail->bindata_len)
1604 ctx->tail->data.strings[ctx->str_idx] = mallocz(ctx->tail->bindata_len + 1);
1605 rbuf_pop(data, ctx->tail->data.strings[ctx->str_idx], ctx->tail->bindata_len);
1606 ctx->tail->data.strings[ctx->str_idx][ctx->tail->bindata_len] = 0;
@@ -1502,7 +1613,7 @@ static int parse_properties_array(struct mqtt_properties_parser_ctx *ctx, rbuf_t
1613 ctx->state = PROPERTY_NEXT;
1614 break;
1615 case PROPERTY_TYPE_BIN:
1505 - BUF_READ_CHECK_AT_LEAST(data, ctx->tail->bindata_len);
1616 + BUF_READ_CHECK_AT_LEAST(data, ctx->tail->bindata_len)
1617 ctx->tail->data.bindata = mallocz(ctx->tail->bindata_len);
1618 rbuf_pop(data, ctx->tail->data.bindata, ctx->tail->bindata_len);
1619 ctx->bytes_consumed += ctx->tail->bindata_len;
@@ -1518,20 +1629,20 @@ static int parse_properties_array(struct mqtt_properties_parser_ctx *ctx, rbuf_t
1629 }
1630 return rc;
1631 case PROPERTY_TYPE_UINT8:
1521 - BUF_READ_CHECK_AT_LEAST(data, sizeof(uint8_t));
1632 + BUF_READ_CHECK_AT_LEAST(data, sizeof(uint8_t))
1633 rbuf_pop(data, (char*)&ctx->tail->data.uint8, sizeof(uint8_t));
1634 ctx->bytes_consumed += sizeof(uint8_t);
1635 ctx->state = PROPERTY_NEXT;
1636 break;
1637 case PROPERTY_TYPE_UINT32:
1527 - BUF_READ_CHECK_AT_LEAST(data, sizeof(uint32_t));
1638 + BUF_READ_CHECK_AT_LEAST(data, sizeof(uint32_t))
1639 rbuf_pop(data, (char*)&ctx->tail->data.uint32, sizeof(uint32_t));
1640 ctx->tail->data.uint32 = be32toh(ctx->tail->data.uint32);
1641 ctx->bytes_consumed += sizeof(uint32_t);
1642 ctx->state = PROPERTY_NEXT;
1643 break;
1644 case PROPERTY_TYPE_UINT16:
1534 - BUF_READ_CHECK_AT_LEAST(data, sizeof(uint16_t));
1645 + BUF_READ_CHECK_AT_LEAST(data, sizeof(uint16_t))
1646 rbuf_pop(data, (char*)&ctx->tail->data.uint16, sizeof(uint16_t));
1647 ctx->tail->data.uint16 = be16toh(ctx->tail->data.uint16);
1648 ctx->bytes_consumed += sizeof(uint16_t);
@@ -1552,7 +1663,7 @@ static int parse_connack_varhdr(struct mqtt_ng_client *client)
1663 struct mqtt_ng_parser *parser = &client->parser;
1664 switch (parser->varhdr_state) {
1665 case MQTT_PARSE_VARHDR_INITIAL:
1555 - BUF_READ_CHECK_AT_LEAST(parser->received_data, 2);
1666 + BUF_READ_CHECK_AT_LEAST(parser->received_data, 2)
1667 rbuf_pop(parser->received_data, (char*)&parser->mqtt_packet.connack.flags, 1);
1668 rbuf_pop(parser->received_data, (char*)&parser->mqtt_packet.connack.reason_code, 1);
1669 parser->varhdr_state = MQTT_PARSE_VARHDR_PROPS;
@@ -1577,7 +1688,7 @@ static int parse_disconnect_varhdr(struct mqtt_ng_client *client)
1688 parser->mqtt_packet.disconnect.reason_code = 0;
1689 return MQTT_NG_CLIENT_PARSE_DONE;
1690 }
1580 - BUF_READ_CHECK_AT_LEAST(parser->received_data, 1);
1691 + BUF_READ_CHECK_AT_LEAST(parser->received_data, 1)
1692 rbuf_pop(parser->received_data, (char*)&parser->mqtt_packet.disconnect.reason_code, 1);
1693 if (parser->mqtt_fixed_hdr_remaining_length == 1)
1694 return MQTT_NG_CLIENT_PARSE_DONE;
@@ -1598,7 +1709,7 @@ static int parse_puback_varhdr(struct mqtt_ng_client *client)
1709 struct mqtt_ng_parser *parser = &client->parser;
1710 switch (parser->varhdr_state) {
1711 case MQTT_PARSE_VARHDR_INITIAL:
1601 - BUF_READ_CHECK_AT_LEAST(parser->received_data, 2);
1712 + BUF_READ_CHECK_AT_LEAST(parser->received_data, 2)
1713 rbuf_pop(parser->received_data, (char*)&parser->mqtt_packet.puback.packet_id, 2);
1714 parser->mqtt_packet.puback.packet_id = be16toh(parser->mqtt_packet.puback.packet_id);
1715 if (parser->mqtt_fixed_hdr_remaining_length < 3) {
@@ -1611,7 +1722,7 @@ static int parse_puback_varhdr(struct mqtt_ng_client *client)
1722 parser->varhdr_state = MQTT_PARSE_VARHDR_OPTIONAL_REASON_CODE;
1723 /* FALLTHROUGH */
1724 case MQTT_PARSE_VARHDR_OPTIONAL_REASON_CODE:
1614 - BUF_READ_CHECK_AT_LEAST(parser->received_data, 1);
1725 + BUF_READ_CHECK_AT_LEAST(parser->received_data, 1)
1726 rbuf_pop(parser->received_data, (char*)&parser->mqtt_packet.puback.reason_code, 1);
1727 // LOL so in CONNACK you have to have 0 byte to
1728 // signify empty properties list
@@ -1640,7 +1751,7 @@ static int parse_suback_varhdr(struct mqtt_ng_client *client)
1751 switch (parser->varhdr_state) {
1752 case MQTT_PARSE_VARHDR_INITIAL:
1753 suback->reason_codes = NULL;
1643 - BUF_READ_CHECK_AT_LEAST(parser->received_data, 2);
1754 + BUF_READ_CHECK_AT_LEAST(parser->received_data, 2)
1755 rbuf_pop(parser->received_data, (char*)&suback->packet_id, 2);
1756 suback->packet_id = be16toh(suback->packet_id);
1757 parser->varhdr_state = MQTT_PARSE_VARHDR_PROPS;
@@ -1682,7 +1793,7 @@ static int parse_publish_varhdr(struct mqtt_ng_client *client)
1793 struct mqtt_publish *publish = &client->parser.mqtt_packet.publish;
1794 switch (parser->varhdr_state) {
1795 case MQTT_PARSE_VARHDR_INITIAL:
1685 - BUF_READ_CHECK_AT_LEAST(parser->received_data, 2);
1796 + BUF_READ_CHECK_AT_LEAST(parser->received_data, 2)
1797 publish->topic = NULL;
1798 publish->qos = ((parser->mqtt_control_packet_type >> 1) & 0x03);
1799 rbuf_pop(parser->received_data, (char*)&publish->topic_len, 2);
@@ -1697,7 +1808,7 @@ static int parse_publish_varhdr(struct mqtt_ng_client *client)
1808 /* FALLTHROUGH */
1809 case MQTT_PARSE_VARHDR_TOPICNAME:
1810 // TODO check empty topic can be valid? In which case we have to skip this step
1700 - BUF_READ_CHECK_AT_LEAST(parser->received_data, publish->topic_len);
1811 + BUF_READ_CHECK_AT_LEAST(parser->received_data, publish->topic_len)
1812 rbuf_pop(parser->received_data, publish->topic, publish->topic_len);
1813 parser->mqtt_parsed_len += publish->topic_len;
1814 parser->varhdr_state = MQTT_PARSE_VARHDR_POST_TOPICNAME;
@@ -1711,7 +1822,7 @@ static int parse_publish_varhdr(struct mqtt_ng_client *client)
1822 parser->varhdr_state = MQTT_PARSE_VARHDR_PACKET_ID;
1823 /* FALLTHROUGH */
1824 case MQTT_PARSE_VARHDR_PACKET_ID:
1714 - BUF_READ_CHECK_AT_LEAST(parser->received_data, 2);
1825 + BUF_READ_CHECK_AT_LEAST(parser->received_data, 2)
1826 rbuf_pop(parser->received_data, (char*)&publish->packet_id, 2);
1827 publish->packet_id = be16toh(publish->packet_id);
1828 parser->varhdr_state = MQTT_PARSE_VARHDR_PROPS;
@@ -1736,7 +1847,7 @@ static int parse_publish_varhdr(struct mqtt_ng_client *client)
1847 publish->data = NULL;
1848 return MQTT_NG_CLIENT_PARSE_DONE; // 0 length payload is OK [MQTT-3.3.3]
1849 }
1739 - BUF_READ_CHECK_AT_LEAST(parser->received_data, publish->data_len);
1850 + BUF_READ_CHECK_AT_LEAST(parser->received_data, publish->data_len)
1851
1852 publish->data = mallocz(publish->data_len);
1853 rbuf_pop(parser->received_data, publish->data, publish->data_len);
@@ -1758,7 +1869,7 @@ static int parse_data(struct mqtt_ng_client *client)
1869 struct mqtt_ng_parser *parser = &client->parser;
1870 switch(parser->state) {
1871 case MQTT_PARSE_FIXED_HEADER_PACKET_TYPE:
1761 - BUF_READ_CHECK_AT_LEAST(parser->received_data, 1);
1872 + BUF_READ_CHECK_AT_LEAST(parser->received_data, 1)
1873 rbuf_pop(parser->received_data, (char*)&parser->mqtt_control_packet_type, 1);
1874 vbi_parser_reset_ctx(&parser->vbi_parser);
1875 parser->state = MQTT_PARSE_FIXED_HEADER_LEN;
@@ -1812,6 +1923,7 @@ static int parse_data(struct mqtt_ng_client *client)
1923 }
1924 parser->state = MQTT_PARSE_MQTT_PACKET_DONE;
1925 ping_timeout = 0;
1926 + check_packet_monitor_list_for_timeouts(client);
1927 break;
1928 case MQTT_CPT_DISCONNECT:
1929 rc = parse_disconnect_varhdr(client);
@@ -1924,50 +2036,6 @@ static void try_send_all(struct mqtt_ng_client *client) {
2036 } while(send_all_message_fragments(client) >= 0);
2037 }
2038
1927 -static void mark_message_for_gc(struct buffer_fragment *frag)
1928 -{
1929 - while (frag) {
1930 - frag->flags |= BUFFER_FRAG_GARBAGE_COLLECT;
1931 - buffer_frag_free_data(frag);
1932 - if (frag->flags & BUFFER_FRAG_MQTT_PACKET_TAIL)
1933 - return;
1934 - frag = frag->next;
1935 - }
1936 -}
1937 -
1938 -static int mark_packet_acked(struct mqtt_ng_client *client, uint16_t packet_id)
1939 -{
1940 - size_t reclaimable = 0;
1941 - LOCK_HDR_BUFFER(&client->main_buffer);
1942 - struct buffer_fragment *frag = BUFFER_FIRST_FRAG(&client->main_buffer.hdr_buffer);
1943 - while (frag) {
1944 - if ( (frag->flags & BUFFER_FRAG_MQTT_PACKET_HEAD) && frag->packet_id == packet_id) {
1945 - if (!frag->sent) {
1946 - nd_log(NDLS_DAEMON, NDLP_ERR, "Received packet_id (%" PRIu16 ") belongs to MQTT packet which was not yet sent!", packet_id);
1947 - UNLOCK_HDR_BUFFER(&client->main_buffer);
1948 - return 1;
1949 - }
1950 - pulse_aclk_sent_message_acked(frag->sent_monotonic_ut, frag->len);
1951 - mark_message_for_gc(frag);
1952 -
1953 - size_t used = BUFFER_BYTES_USED(&client->main_buffer.hdr_buffer);
1954 - if (reclaimable >= (used / 4))
1955 - transaction_buffer_garbage_collect(&client->main_buffer);
1956 -
1957 - UNLOCK_HDR_BUFFER(&client->main_buffer);
1958 - return 0;
1959 - }
1960 -
1961 - if(frag_is_marked_for_gc(frag))
1962 - reclaimable += FRAG_SIZE_IN_BUFFER(frag);
1963 -
1964 - frag = frag->next;
1965 - }
1966 - nd_log(NDLS_DAEMON, NDLP_ERR, "Received packet_id (%" PRIu16 ") is unknown!", packet_id);
1967 - UNLOCK_HDR_BUFFER(&client->main_buffer);
1968 - return 1;
1969 -}
1970 -
2039 int handle_incoming_traffic(struct mqtt_ng_client *client)
2040 {
2041 int rc;
@@ -2012,8 +2080,7 @@ int handle_incoming_traffic(struct mqtt_ng_client *client)
2080 client->client_state = MQTT_STATE_CONNECTED;
2081 break;
2082 }
2015 - client->client_state = MQTT_STATE_ERROR;
2016 - return MQTT_NG_CLIENT_SERVER_RETURNED_ERROR;
2083 + client->client_state = MQTT_STATE_ERROR; return MQTT_NG_CLIENT_SERVER_RETURNED_ERROR;
2084
2085 case MQTT_CPT_PUBACK:
2086 worker_is_busy(WORKER_ACLK_CPT_PUBACK);