mqtt_outbox.c 7.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301
  1. /* This is a modification of https://github.com/espressif/esp-mqtt/blob/master/lib/mqtt_outbox.c
  2. * to use the PSRAM instead of the internal heap.
  3. */
  4. #include "mqtt_outbox.h"
  5. #include <stdlib.h>
  6. #include <string.h>
  7. #include "sys/queue.h"
  8. #include "esp_log.h"
  9. #include "esp_heap_caps.h"
  10. #define USE_PSRAM
  11. #ifdef CONFIG_MQTT_CUSTOM_OUTBOX
  12. static const char *TAG = "outbox";
  13. typedef struct outbox_item {
  14. char *buffer;
  15. int len;
  16. int msg_id;
  17. int msg_type;
  18. int msg_qos;
  19. outbox_tick_t tick;
  20. pending_state_t pending;
  21. STAILQ_ENTRY(outbox_item) next;
  22. } outbox_item_t;
  23. STAILQ_HEAD(outbox_list_t, outbox_item);
  24. outbox_handle_t outbox_init(void)
  25. {
  26. #ifdef USE_PSRAM
  27. outbox_handle_t outbox = heap_caps_calloc(1, sizeof(struct outbox_list_t), MALLOC_CAP_8BIT | MALLOC_CAP_SPIRAM);
  28. #else
  29. outbox_handle_t outbox = calloc(1, sizeof(struct outbox_list_t));
  30. #endif
  31. //ESP_MEM_CHECK(TAG, outbox, return NULL);
  32. STAILQ_INIT(outbox);
  33. return outbox;
  34. }
  35. outbox_item_handle_t outbox_enqueue(outbox_handle_t outbox, outbox_message_handle_t message, outbox_tick_t tick)
  36. {
  37. #ifdef USE_PSRAM
  38. outbox_item_handle_t item = heap_caps_calloc(1, sizeof(outbox_item_t), MALLOC_CAP_8BIT | MALLOC_CAP_SPIRAM);
  39. #else
  40. outbox_item_handle_t item = calloc(1, sizeof(outbox_item_t));
  41. #endif
  42. //ESP_MEM_CHECK(TAG, item, return NULL);
  43. item->msg_id = message->msg_id;
  44. item->msg_type = message->msg_type;
  45. item->msg_qos = message->msg_qos;
  46. item->tick = tick;
  47. item->len = message->len + message->remaining_len;
  48. item->pending = QUEUED;
  49. #ifdef USE_PSRAM
  50. item->buffer = heap_caps_malloc(message->len + message->remaining_len, MALLOC_CAP_8BIT | MALLOC_CAP_SPIRAM);
  51. #else
  52. item->buffer = malloc(message->len + message->remaining_len);
  53. #endif
  54. /*ESP_MEM_CHECK(TAG, item->buffer, {
  55. free(item);
  56. return NULL;
  57. });*/
  58. memcpy(item->buffer, message->data, message->len);
  59. if (message->remaining_data) {
  60. memcpy(item->buffer + message->len, message->remaining_data, message->remaining_len);
  61. }
  62. STAILQ_INSERT_TAIL(outbox, item, next);
  63. ESP_LOGD(TAG, "ENQUEUE msgid=%d, msg_type=%d, len=%d, size=%d", message->msg_id, message->msg_type, message->len + message->remaining_len, outbox_get_size(outbox));
  64. return item;
  65. }
  66. outbox_item_handle_t outbox_get(outbox_handle_t outbox, int msg_id)
  67. {
  68. outbox_item_handle_t item;
  69. STAILQ_FOREACH(item, outbox, next) {
  70. if (item->msg_id == msg_id) {
  71. return item;
  72. }
  73. }
  74. return NULL;
  75. }
  76. outbox_item_handle_t outbox_dequeue(outbox_handle_t outbox, pending_state_t pending, outbox_tick_t *tick)
  77. {
  78. outbox_item_handle_t item;
  79. STAILQ_FOREACH(item, outbox, next) {
  80. if (item->pending == pending) {
  81. if (tick) {
  82. *tick = item->tick;
  83. }
  84. return item;
  85. }
  86. }
  87. return NULL;
  88. }
  89. esp_err_t outbox_delete_item(outbox_handle_t outbox, outbox_item_handle_t item_to_delete)
  90. {
  91. outbox_item_handle_t item;
  92. STAILQ_FOREACH(item, outbox, next) {
  93. if (item == item_to_delete) {
  94. STAILQ_REMOVE(outbox, item, outbox_item, next);
  95. #ifdef USE_PSRAM
  96. heap_caps_free(item->buffer);
  97. heap_caps_free(item);
  98. #else
  99. free(item->buffer);
  100. free(item);
  101. #endif
  102. return ESP_OK;
  103. }
  104. }
  105. return ESP_FAIL;
  106. }
  107. uint8_t *outbox_item_get_data(outbox_item_handle_t item, size_t *len, uint16_t *msg_id, int *msg_type, int *qos)
  108. {
  109. if (item) {
  110. *len = item->len;
  111. *msg_id = item->msg_id;
  112. *msg_type = item->msg_type;
  113. *qos = item->msg_qos;
  114. return (uint8_t *)item->buffer;
  115. }
  116. return NULL;
  117. }
  118. esp_err_t outbox_delete(outbox_handle_t outbox, int msg_id, int msg_type)
  119. {
  120. outbox_item_handle_t item, tmp;
  121. STAILQ_FOREACH_SAFE(item, outbox, next, tmp) {
  122. if (item->msg_id == msg_id && (0xFF & (item->msg_type)) == msg_type) {
  123. STAILQ_REMOVE(outbox, item, outbox_item, next);
  124. #ifdef USE_PSRAM
  125. heap_caps_free(item->buffer);
  126. heap_caps_free(item);
  127. #else
  128. free(item->buffer);
  129. free(item);
  130. #endif
  131. ESP_LOGD(TAG, "DELETED msgid=%d, msg_type=%d, remain size=%d", msg_id, msg_type, outbox_get_size(outbox));
  132. return ESP_OK;
  133. }
  134. }
  135. return ESP_FAIL;
  136. }
  137. esp_err_t outbox_delete_msgid(outbox_handle_t outbox, int msg_id)
  138. {
  139. outbox_item_handle_t item, tmp;
  140. STAILQ_FOREACH_SAFE(item, outbox, next, tmp) {
  141. if (item->msg_id == msg_id) {
  142. STAILQ_REMOVE(outbox, item, outbox_item, next);
  143. #ifdef USE_PSRAM
  144. heap_caps_free(item->buffer);
  145. heap_caps_free(item);
  146. #else
  147. free(item->buffer);
  148. free(item);
  149. #endif
  150. }
  151. }
  152. return ESP_OK;
  153. }
  154. esp_err_t outbox_set_pending(outbox_handle_t outbox, int msg_id, pending_state_t pending)
  155. {
  156. outbox_item_handle_t item = outbox_get(outbox, msg_id);
  157. if (item) {
  158. item->pending = pending;
  159. return ESP_OK;
  160. }
  161. return ESP_FAIL;
  162. }
  163. pending_state_t outbox_item_get_pending(outbox_item_handle_t item)
  164. {
  165. if (item) {
  166. return item->pending;
  167. }
  168. return QUEUED;
  169. }
  170. esp_err_t outbox_set_tick(outbox_handle_t outbox, int msg_id, outbox_tick_t tick)
  171. {
  172. outbox_item_handle_t item = outbox_get(outbox, msg_id);
  173. if (item) {
  174. item->tick = tick;
  175. return ESP_OK;
  176. }
  177. return ESP_FAIL;
  178. }
  179. esp_err_t outbox_delete_msgtype(outbox_handle_t outbox, int msg_type)
  180. {
  181. outbox_item_handle_t item, tmp;
  182. STAILQ_FOREACH_SAFE(item, outbox, next, tmp) {
  183. if (item->msg_type == msg_type) {
  184. STAILQ_REMOVE(outbox, item, outbox_item, next);
  185. #ifdef USE_PSRAM
  186. heap_caps_free(item->buffer);
  187. heap_caps_free(item);
  188. #else
  189. free(item->buffer);
  190. free(item);
  191. #endif
  192. }
  193. }
  194. return ESP_OK;
  195. }
  196. int outbox_delete_single_expired(outbox_handle_t outbox, outbox_tick_t current_tick, outbox_tick_t timeout)
  197. {
  198. int msg_id = -1;
  199. outbox_item_handle_t item;
  200. STAILQ_FOREACH(item, outbox, next) {
  201. if (current_tick - item->tick > timeout) {
  202. STAILQ_REMOVE(outbox, item, outbox_item, next);
  203. #ifdef USE_PSRAM
  204. heap_caps_free(item->buffer);
  205. #else
  206. free(item->buffer);
  207. #endif
  208. msg_id = item->msg_id;
  209. #ifdef USE_PSRAM
  210. heap_caps_free(item);
  211. #else
  212. free(item);
  213. #endif
  214. return msg_id;
  215. }
  216. }
  217. return msg_id;
  218. }
  219. int outbox_delete_expired(outbox_handle_t outbox, outbox_tick_t current_tick, outbox_tick_t timeout)
  220. {
  221. int deleted_items = 0;
  222. outbox_item_handle_t item, tmp;
  223. STAILQ_FOREACH_SAFE(item, outbox, next, tmp) {
  224. if (current_tick - item->tick > timeout) {
  225. STAILQ_REMOVE(outbox, item, outbox_item, next);
  226. #ifdef USE_PSRAM
  227. heap_caps_free(item->buffer);
  228. heap_caps_free(item);
  229. #else
  230. free(item->buffer);
  231. free(item);
  232. #endif
  233. deleted_items ++;
  234. }
  235. }
  236. return deleted_items;
  237. }
  238. int outbox_get_size(outbox_handle_t outbox)
  239. {
  240. int siz = 0;
  241. outbox_item_handle_t item;
  242. STAILQ_FOREACH(item, outbox, next) {
  243. // Suppressing "use after free" warning as this could happen only if queue is in inconsistent state
  244. // which never happens if STAILQ interface used
  245. siz += item->len; // NOLINT(clang-analyzer-unix.Malloc)
  246. }
  247. return siz;
  248. }
  249. void outbox_delete_all_items(outbox_handle_t outbox)
  250. {
  251. outbox_item_handle_t item, tmp;
  252. STAILQ_FOREACH_SAFE(item, outbox, next, tmp) {
  253. STAILQ_REMOVE(outbox, item, outbox_item, next);
  254. #ifdef USE_PSRAM
  255. heap_caps_free(item->buffer);
  256. heap_caps_free(item);
  257. #else
  258. free(item->buffer);
  259. free(item);
  260. #endif
  261. }
  262. }
  263. void outbox_destroy(outbox_handle_t outbox)
  264. {
  265. outbox_delete_all_items(outbox);
  266. #ifdef USE_PSRAM
  267. heap_caps_free(outbox);
  268. #else
  269. free(outbox);
  270. #endif
  271. }
  272. #endif /* CONFIG_MQTT_CUSTOM_OUTBOX */