mqtt_ng.c 84 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832183318341835183618371838183918401841184218431844184518461847184818491850185118521853185418551856185718581859186018611862186318641865186618671868186918701871187218731874187518761877187818791880188118821883188418851886188718881889189018911892189318941895189618971898189919001901190219031904190519061907190819091910191119121913191419151916191719181919192019211922192319241925192619271928192919301931193219331934193519361937193819391940194119421943194419451946194719481949195019511952195319541955195619571958195919601961196219631964196519661967196819691970197119721973197419751976197719781979198019811982198319841985198619871988198919901991199219931994199519961997199819992000200120022003200420052006200720082009201020112012201320142015201620172018201920202021202220232024202520262027202820292030203120322033203420352036203720382039204020412042204320442045204620472048204920502051205220532054205520562057205820592060206120622063206420652066206720682069207020712072207320742075207620772078207920802081208220832084208520862087208820892090209120922093209420952096209720982099210021012102210321042105210621072108210921102111211221132114211521162117211821192120212121222123212421252126212721282129213021312132213321342135213621372138213921402141214221432144214521462147214821492150215121522153215421552156215721582159216021612162216321642165216621672168216921702171217221732174217521762177217821792180218121822183218421852186218721882189219021912192219321942195219621972198219922002201220222032204220522062207220822092210221122122213221422152216221722182219222022212222222322242225222622272228222922302231223222332234223522362237
  1. // Copyright: SPDX-License-Identifier: GPL-3.0-only
  2. #ifndef _GNU_SOURCE
  3. #define _GNU_SOURCE
  4. #endif
  5. #include <stdint.h>
  6. #include <stdlib.h>
  7. #include <string.h>
  8. #include <pthread.h>
  9. #include <inttypes.h>
  10. #include "c_rhash/c_rhash.h"
  11. #include "common_internal.h"
  12. #include "mqtt_constants.h"
  13. #include "mqtt_wss_log.h"
  14. #include "mqtt_ng.h"
  15. #define UNIT_LOG_PREFIX "mqtt_client: "
  16. #define FATAL(fmt, ...) mws_fatal(client->log, UNIT_LOG_PREFIX fmt, ##__VA_ARGS__)
  17. #define ERROR(fmt, ...) mws_error(client->log, UNIT_LOG_PREFIX fmt, ##__VA_ARGS__)
  18. #define WARN(fmt, ...) mws_warn (client->log, UNIT_LOG_PREFIX fmt, ##__VA_ARGS__)
  19. #define INFO(fmt, ...) mws_info (client->log, UNIT_LOG_PREFIX fmt, ##__VA_ARGS__)
  20. #define DEBUG(fmt, ...) mws_debug(client->log, UNIT_LOG_PREFIX fmt, ##__VA_ARGS__)
  21. #define SMALL_STRING_DONT_FRAGMENT_LIMIT 128
  22. #define MIN(a,b) (((a)<(b))?(a):(b))
  23. #define LOCK_HDR_BUFFER(buffer) pthread_mutex_lock(&((buffer)->mutex))
  24. #define UNLOCK_HDR_BUFFER(buffer) pthread_mutex_unlock(&((buffer)->mutex))
  25. #define BUFFER_FRAG_GARBAGE_COLLECT 0x01
  26. // some packets can be marked for garbage collection
  27. // immediately when they are sent (e.g. sent PUBACK on QoS1)
  28. #define BUFFER_FRAG_GARBAGE_COLLECT_ON_SEND 0x02
  29. // as buffer fragment can point to both
  30. // external data and data in the same buffer
  31. // we mark the former case with BUFFER_FRAG_DATA_EXTERNAL
  32. #define BUFFER_FRAG_DATA_EXTERNAL 0x04
  33. // as single MQTT Packet can be stored into multiple
  34. // buffer fragments (depending on copy requirements)
  35. // this marks this fragment to be the first/last
  36. #define BUFFER_FRAG_MQTT_PACKET_HEAD 0x10
  37. #define BUFFER_FRAG_MQTT_PACKET_TAIL 0x20
  38. typedef uint16_t buffer_frag_flag_t;
  39. struct buffer_fragment {
  40. size_t len;
  41. size_t sent;
  42. buffer_frag_flag_t flags;
  43. void (*free_fnc)(void *ptr);
  44. unsigned char *data;
  45. uint16_t packet_id;
  46. struct buffer_fragment *next;
  47. };
  48. typedef struct buffer_fragment *mqtt_msg_data;
  49. // buffer used for MQTT headers only
  50. // not for actual data sent
  51. struct header_buffer {
  52. size_t size;
  53. unsigned char *data;
  54. unsigned char *tail;
  55. struct buffer_fragment *tail_frag;
  56. };
  57. struct transaction_buffer {
  58. struct header_buffer hdr_buffer;
  59. // used while building new message
  60. // to be able to revert state easily
  61. // in case of error mid processing
  62. struct header_buffer state_backup;
  63. pthread_mutex_t mutex;
  64. struct buffer_fragment *sending_frag;
  65. };
  66. enum mqtt_client_state {
  67. RAW = 0,
  68. CONNECT_PENDING,
  69. CONNECTING,
  70. CONNECTED,
  71. ERROR,
  72. DISCONNECTED
  73. };
  74. enum parser_state {
  75. MQTT_PARSE_FIXED_HEADER_PACKET_TYPE = 0,
  76. MQTT_PARSE_FIXED_HEADER_LEN,
  77. MQTT_PARSE_VARIABLE_HEADER,
  78. MQTT_PARSE_MQTT_PACKET_DONE
  79. };
  80. enum varhdr_parser_state {
  81. MQTT_PARSE_VARHDR_INITIAL = 0,
  82. MQTT_PARSE_VARHDR_OPTIONAL_REASON_CODE,
  83. MQTT_PARSE_VARHDR_PROPS,
  84. MQTT_PARSE_VARHDR_TOPICNAME,
  85. MQTT_PARSE_VARHDR_POST_TOPICNAME,
  86. MQTT_PARSE_VARHDR_PACKET_ID,
  87. MQTT_PARSE_REASONCODES,
  88. MQTT_PARSE_PAYLOAD
  89. };
  90. struct mqtt_vbi_parser_ctx {
  91. char data[MQTT_VBI_MAXBYTES];
  92. uint8_t bytes;
  93. uint32_t result;
  94. };
  95. enum mqtt_datatype {
  96. MQTT_TYPE_UNKNOWN = 0,
  97. MQTT_TYPE_UINT_8,
  98. MQTT_TYPE_UINT_16,
  99. MQTT_TYPE_UINT_32,
  100. MQTT_TYPE_VBI,
  101. MQTT_TYPE_STR,
  102. MQTT_TYPE_STR_PAIR,
  103. MQTT_TYPE_BIN
  104. };
  105. struct mqtt_property {
  106. uint8_t id;
  107. enum mqtt_datatype type;
  108. union {
  109. char *strings[2];
  110. void *bindata;
  111. uint8_t uint8;
  112. uint16_t uint16;
  113. uint32_t uint32;
  114. } data;
  115. size_t bindata_len;
  116. struct mqtt_property *next;
  117. };
  118. enum mqtt_properties_parser_state {
  119. PROPERTIES_LENGTH = 0,
  120. PROPERTY_CREATE,
  121. PROPERTY_ID,
  122. PROPERTY_TYPE_UINT8,
  123. PROPERTY_TYPE_UINT16,
  124. PROPERTY_TYPE_UINT32,
  125. PROPERTY_TYPE_STR_BIN_LEN,
  126. PROPERTY_TYPE_STR,
  127. PROPERTY_TYPE_BIN,
  128. PROPERTY_TYPE_VBI,
  129. PROPERTY_NEXT
  130. };
  131. struct mqtt_properties_parser_ctx {
  132. enum mqtt_properties_parser_state state;
  133. struct mqtt_property *head;
  134. struct mqtt_property *tail;
  135. uint32_t properties_length;
  136. uint32_t vbi_length;
  137. struct mqtt_vbi_parser_ctx vbi_parser_ctx;
  138. size_t bytes_consumed;
  139. int str_idx;
  140. };
  141. struct mqtt_connack {
  142. uint8_t flags;
  143. uint8_t reason_code;
  144. };
  145. struct mqtt_puback {
  146. uint16_t packet_id;
  147. uint8_t reason_code;
  148. };
  149. struct mqtt_suback {
  150. uint16_t packet_id;
  151. uint8_t *reason_codes;
  152. uint8_t reason_code_count;
  153. uint8_t reason_codes_pending;
  154. };
  155. struct mqtt_publish {
  156. uint16_t topic_len;
  157. char *topic;
  158. uint16_t packet_id;
  159. size_t data_len;
  160. char *data;
  161. uint8_t qos;
  162. };
  163. struct mqtt_disconnect {
  164. uint8_t reason_code;
  165. };
  166. struct mqtt_ng_parser {
  167. rbuf_t received_data;
  168. uint8_t mqtt_control_packet_type;
  169. uint32_t mqtt_fixed_hdr_remaining_length;
  170. size_t mqtt_parsed_len;
  171. struct mqtt_vbi_parser_ctx vbi_parser;
  172. struct mqtt_properties_parser_ctx properties_parser;
  173. enum parser_state state;
  174. enum varhdr_parser_state varhdr_state;
  175. struct mqtt_property *varhdr_properties;
  176. union {
  177. struct mqtt_connack connack;
  178. struct mqtt_puback puback;
  179. struct mqtt_suback suback;
  180. struct mqtt_publish publish;
  181. struct mqtt_disconnect disconnect;
  182. } mqtt_packet;
  183. };
  184. struct topic_alias_data {
  185. uint16_t idx;
  186. uint32_t usage_count;
  187. };
  188. struct topic_aliases_data {
  189. c_rhash stoi_dict;
  190. uint32_t idx_max;
  191. uint32_t idx_assigned;
  192. pthread_rwlock_t rwlock;
  193. };
  194. struct mqtt_ng_client {
  195. struct transaction_buffer main_buffer;
  196. enum mqtt_client_state client_state;
  197. mqtt_msg_data connect_msg;
  198. mqtt_wss_log_ctx_t log;
  199. mqtt_ng_send_fnc_t send_fnc_ptr;
  200. void *user_ctx;
  201. // time when last fragment of MQTT message was sent
  202. time_t time_of_last_send;
  203. struct mqtt_ng_parser parser;
  204. size_t max_mem_bytes;
  205. void (*puback_callback)(uint16_t packet_id);
  206. void (*connack_callback)(void* user_ctx, int connack_reply);
  207. void (*msg_callback)(const char *topic, const void *msg, size_t msglen, int qos);
  208. unsigned int ping_pending:1;
  209. struct mqtt_ng_stats stats;
  210. pthread_mutex_t stats_mutex;
  211. struct topic_aliases_data tx_topic_aliases;
  212. c_rhash rx_aliases;
  213. size_t max_msg_size;
  214. };
  215. unsigned char pingreq[] = { MQTT_CPT_PINGREQ << 4, 0x00 };
  216. struct buffer_fragment ping_frag = {
  217. .data = pingreq,
  218. .flags = BUFFER_FRAG_MQTT_PACKET_HEAD | BUFFER_FRAG_MQTT_PACKET_TAIL,
  219. .free_fnc = NULL,
  220. .len = sizeof(pingreq),
  221. .next = NULL,
  222. .sent = 0,
  223. .packet_id = 0
  224. };
  225. int uint32_to_mqtt_vbi(uint32_t input, unsigned char *output) {
  226. int i = 1;
  227. *output = 0;
  228. /* MQTT 5 specs allows max 4 bytes of output
  229. making it 0xFF, 0xFF, 0xFF, 0x7F
  230. representing number 268435455 decimal
  231. see 1.5.5. Variable Byte Integer */
  232. if(input >= 256 * 1024 * 1024)
  233. return 0;
  234. if(!input) {
  235. *output = 0;
  236. return 1;
  237. }
  238. while(input) {
  239. output[i-1] = input & MQTT_VBI_DATA_MASK;
  240. input >>= 7;
  241. if (input)
  242. output[i-1] |= MQTT_VBI_CONTINUATION_FLAG;
  243. i++;
  244. }
  245. return i - 1;
  246. }
  247. int mqtt_vbi_to_uint32(char *input, uint32_t *output) {
  248. // dont want to operate directly on output
  249. // as I want it to be possible for input and output
  250. // pointer to be the same
  251. uint32_t result = 0;
  252. uint32_t multiplier = 1;
  253. do {
  254. result += (uint32_t)(*input & MQTT_VBI_DATA_MASK) * multiplier;
  255. if (multiplier > 128*128*128)
  256. return 1;
  257. multiplier <<= 7;
  258. } while (*input++ & MQTT_VBI_CONTINUATION_FLAG);
  259. *output = result;
  260. return 0;
  261. }
  262. #ifdef TESTS
  263. #include <stdio.h>
  264. #define MQTT_VBI_MAXLEN 4
  265. // we add extra byte to check we dont write out of bounds
  266. // in case where 4 bytes are supposed to be written
  267. static const char _mqtt_vbi_0[MQTT_VBI_MAXLEN + 1] = { 0x00, 0x00, 0x00, 0x00, 0x00 };
  268. static const char _mqtt_vbi_127[MQTT_VBI_MAXLEN + 1] = { 0x7F, 0x00, 0x00, 0x00, 0x00 };
  269. static const char _mqtt_vbi_128[MQTT_VBI_MAXLEN + 1] = { 0x80, 0x01, 0x00, 0x00, 0x00 };
  270. static const char _mqtt_vbi_16383[MQTT_VBI_MAXLEN + 1] = { 0xFF, 0x7F, 0x00, 0x00, 0x00 };
  271. static const char _mqtt_vbi_16384[MQTT_VBI_MAXLEN + 1] = { 0x80, 0x80, 0x01, 0x00, 0x00 };
  272. static const char _mqtt_vbi_2097151[MQTT_VBI_MAXLEN + 1] = { 0xFF, 0xFF, 0x7F, 0x00, 0x00 };
  273. static const char _mqtt_vbi_2097152[MQTT_VBI_MAXLEN + 1] = { 0x80, 0x80, 0x80, 0x01, 0x00 };
  274. static const char _mqtt_vbi_268435455[MQTT_VBI_MAXLEN + 1] = { 0xFF, 0xFF, 0xFF, 0x7F, 0x00 };
  275. static const char _mqtt_vbi_999999999[MQTT_VBI_MAXLEN + 1] = { 0x80, 0x80, 0x80, 0x80, 0x01 };
  276. #define MQTT_VBI_TESTCASE(case, expected_len) \
  277. { \
  278. memset(buf, 0, MQTT_VBI_MAXLEN + 1); \
  279. int len; \
  280. if ((len=uint32_to_mqtt_vbi(case, buf)) != expected_len) { \
  281. fprintf(stderr, "uint32_to_mqtt_vbi(case:%d, line:%d): Incorrect length returned. Expected %d, Got %d\n", case, __LINE__, expected_len, len); \
  282. return 1; \
  283. } \
  284. if (memcmp(buf, _mqtt_vbi_ ## case, MQTT_VBI_MAXLEN + 1 )) { \
  285. fprintf(stderr, "uint32_to_mqtt_vbi(case:%d, line:%d): Wrong output\n", case, __LINE__); \
  286. return 1; \
  287. } }
  288. int test_uint32_mqtt_vbi() {
  289. char buf[MQTT_VBI_MAXLEN + 1];
  290. MQTT_VBI_TESTCASE(0, 1)
  291. MQTT_VBI_TESTCASE(127, 1)
  292. MQTT_VBI_TESTCASE(128, 2)
  293. MQTT_VBI_TESTCASE(16383, 2)
  294. MQTT_VBI_TESTCASE(16384, 3)
  295. MQTT_VBI_TESTCASE(2097151, 3)
  296. MQTT_VBI_TESTCASE(2097152, 4)
  297. MQTT_VBI_TESTCASE(268435455, 4)
  298. memset(buf, 0, MQTT_VBI_MAXLEN + 1);
  299. int len;
  300. if ((len=uint32_to_mqtt_vbi(268435456, buf)) != 0) {
  301. fprintf(stderr, "uint32_to_mqtt_vbi(case:268435456, line:%d): Incorrect length returned. Expected 0, Got %d\n", __LINE__, len);
  302. return 1;
  303. }
  304. return 0;
  305. }
  306. #define MQTT_VBI2UINT_TESTCASE(case, expected_error) \
  307. { \
  308. uint32_t result; \
  309. int ret = mqtt_vbi_to_uint32(_mqtt_vbi_ ## case, &result); \
  310. if (ret && !(expected_error)) { \
  311. fprintf(stderr, "mqtt_vbi_to_uint(case:%d, line:%d): Unexpectedly Errored\n", (case), __LINE__); \
  312. return 1; \
  313. } \
  314. if (!ret && (expected_error)) { \
  315. fprintf(stderr, "mqtt_vbi_to_uint(case:%d, line:%d): Should return error but didnt\n", (case), __LINE__); \
  316. return 1; \
  317. } \
  318. if (!ret && result != (case)) { \
  319. fprintf(stderr, "mqtt_vbi_to_uint(case:%d, line:%d): Returned wrong result %d\n", (case), __LINE__, result); \
  320. return 1; \
  321. }}
  322. int test_mqtt_vbi_to_uint32() {
  323. MQTT_VBI2UINT_TESTCASE(0, 0)
  324. MQTT_VBI2UINT_TESTCASE(127, 0)
  325. MQTT_VBI2UINT_TESTCASE(128, 0)
  326. MQTT_VBI2UINT_TESTCASE(16383, 0)
  327. MQTT_VBI2UINT_TESTCASE(16384, 0)
  328. MQTT_VBI2UINT_TESTCASE(2097151, 0)
  329. MQTT_VBI2UINT_TESTCASE(2097152, 0)
  330. MQTT_VBI2UINT_TESTCASE(268435455, 0)
  331. MQTT_VBI2UINT_TESTCASE(999999999, 1)
  332. return 0;
  333. }
  334. #endif /* TESTS */
  335. // this helps with switch statements
  336. // as they have to use integer type (not pointer)
  337. enum memory_mode {
  338. MEMCPY,
  339. EXTERNAL_FREE_AFTER_USE,
  340. CALLER_RESPONSIBLE
  341. };
  342. static inline enum memory_mode ptr2memory_mode(void * ptr) {
  343. if (ptr == NULL)
  344. return MEMCPY;
  345. if (ptr == CALLER_RESPONSIBILITY)
  346. return CALLER_RESPONSIBLE;
  347. return EXTERNAL_FREE_AFTER_USE;
  348. }
  349. #define frag_is_marked_for_gc(frag) ((frag->flags & BUFFER_FRAG_GARBAGE_COLLECT) || ((frag->flags & BUFFER_FRAG_GARBAGE_COLLECT_ON_SEND) && frag->sent == frag->len))
  350. #define FRAG_SIZE_IN_BUFFER(frag) (sizeof(struct buffer_fragment) + ((frag->flags & BUFFER_FRAG_DATA_EXTERNAL) ? 0 : frag->len))
  351. static void buffer_frag_free_data(struct buffer_fragment *frag)
  352. {
  353. if ( frag->flags & BUFFER_FRAG_DATA_EXTERNAL && frag->data != NULL) {
  354. switch (ptr2memory_mode(frag->free_fnc)) {
  355. case MEMCPY:
  356. mw_free(frag->data);
  357. break;
  358. case EXTERNAL_FREE_AFTER_USE:
  359. frag->free_fnc(frag->data);
  360. break;
  361. case CALLER_RESPONSIBLE:
  362. break;
  363. }
  364. frag->data = NULL;
  365. }
  366. }
  367. #define HEADER_BUFFER_SIZE 1024*1024
  368. #define GROWTH_FACTOR 1.25
  369. #define BUFFER_BYTES_USED(buf) ((size_t)((buf)->tail - (buf)->data))
  370. #define BUFFER_BYTES_AVAILABLE(buf) ((buf)->size - BUFFER_BYTES_USED(buf))
  371. #define BUFFER_FIRST_FRAG(buf) ((struct buffer_fragment *)((buf)->tail_frag ? (buf)->data : NULL))
  372. static void buffer_purge(struct header_buffer *buf) {
  373. struct buffer_fragment *frag = BUFFER_FIRST_FRAG(buf);
  374. while (frag) {
  375. buffer_frag_free_data(frag);
  376. frag = frag->next;
  377. }
  378. buf->tail = buf->data;
  379. buf->tail_frag = NULL;
  380. }
  381. #define FRAG_PADDING(addr) ((MQTT_WSS_FRAG_MEMALIGN - ((uintptr_t)addr % MQTT_WSS_FRAG_MEMALIGN)) % MQTT_WSS_FRAG_MEMALIGN)
  382. static struct buffer_fragment *buffer_new_frag(struct header_buffer *buf, buffer_frag_flag_t flags)
  383. {
  384. uint8_t padding = FRAG_PADDING(buf->tail);
  385. if (BUFFER_BYTES_AVAILABLE(buf) < sizeof(struct buffer_fragment) + padding)
  386. return NULL;
  387. struct buffer_fragment *frag = (struct buffer_fragment *)(buf->tail + padding);
  388. memset(frag, 0, sizeof(*frag));
  389. buf->tail += sizeof(*frag) + padding;
  390. if (/*!((frag)->flags & BUFFER_FRAG_MQTT_PACKET_HEAD) &&*/ buf->tail_frag)
  391. buf->tail_frag->next = frag;
  392. buf->tail_frag = frag;
  393. frag->data = buf->tail;
  394. frag->flags = flags;
  395. return frag;
  396. }
  397. static void buffer_rebuild(struct header_buffer *buf)
  398. {
  399. struct buffer_fragment *frag = (struct buffer_fragment*)buf->data;
  400. do {
  401. buf->tail = (unsigned char *) frag + sizeof(struct buffer_fragment);
  402. buf->tail_frag = frag;
  403. if (!(frag->flags & BUFFER_FRAG_DATA_EXTERNAL)) {
  404. buf->tail_frag->data = buf->tail;
  405. buf->tail += frag->len;
  406. }
  407. if (frag->next != NULL)
  408. frag->next = (struct buffer_fragment*)(buf->tail + FRAG_PADDING(buf->tail));
  409. frag = frag->next;
  410. } while(frag);
  411. }
  412. static void buffer_garbage_collect(struct header_buffer *buf, mqtt_wss_log_ctx_t log_ctx)
  413. {
  414. #if !defined(MQTT_DEBUG_VERBOSE) && !defined(ADDITIONAL_CHECKS)
  415. (void) log_ctx;
  416. #endif
  417. #ifdef MQTT_DEBUG_VERBOSE
  418. mws_debug(log_ctx, "Buffer Garbage Collection!");
  419. #endif
  420. struct buffer_fragment *frag = BUFFER_FIRST_FRAG(buf);
  421. while (frag) {
  422. if (!frag_is_marked_for_gc(frag))
  423. break;
  424. buffer_frag_free_data(frag);
  425. frag = frag->next;
  426. }
  427. if (frag == BUFFER_FIRST_FRAG(buf)) {
  428. #ifdef MQTT_DEBUG_VERBOSE
  429. mws_debug(log_ctx, "Buffer Garbage Collection! No Space Reclaimed!");
  430. #endif
  431. return;
  432. }
  433. if (!frag) {
  434. buf->tail_frag = NULL;
  435. buf->tail = buf->data;
  436. return;
  437. }
  438. #ifdef ADDITIONAL_CHECKS
  439. if (!(frag->flags & BUFFER_FRAG_MQTT_PACKET_HEAD)) {
  440. mws_error(log_ctx, "Expected to find end of buffer (NULL) or next packet head!");
  441. return;
  442. }
  443. #endif
  444. memmove(buf->data, frag, buf->tail - (unsigned char *) frag);
  445. buffer_rebuild(buf);
  446. }
  447. static void transaction_buffer_garbage_collect(struct transaction_buffer *buf, mqtt_wss_log_ctx_t log_ctx)
  448. {
  449. #ifdef MQTT_DEBUG_VERBOSE
  450. mws_debug(log_ctx, "Transaction Buffer Garbage Collection! %s", buf->sending_frag == NULL ? "NULL" : "in flight message");
  451. #endif
  452. // Invalidate the cached sending fragment
  453. // as we will move data around
  454. if (buf->sending_frag != &ping_frag)
  455. buf->sending_frag = NULL;
  456. buffer_garbage_collect(&buf->hdr_buffer, log_ctx);
  457. }
  458. static int transaction_buffer_grow(struct transaction_buffer *buf, mqtt_wss_log_ctx_t log_ctx, float rate, size_t max)
  459. {
  460. if (buf->hdr_buffer.size >= max)
  461. return 0;
  462. // Invalidate the cached sending fragment
  463. // as we will move data around
  464. if (buf->sending_frag != &ping_frag)
  465. buf->sending_frag = NULL;
  466. buf->hdr_buffer.size *= rate;
  467. if (buf->hdr_buffer.size > max)
  468. buf->hdr_buffer.size = max;
  469. void *ret = mw_realloc(buf->hdr_buffer.data, buf->hdr_buffer.size);
  470. if (ret == NULL) {
  471. mws_warn(log_ctx, "Buffer growth failed (realloc)");
  472. return 1;
  473. }
  474. mws_debug(log_ctx, "Message metadata buffer was grown");
  475. buf->hdr_buffer.data = ret;
  476. buffer_rebuild(&buf->hdr_buffer);
  477. return 0;
  478. }
  479. inline static int transaction_buffer_init(struct transaction_buffer *to_init, size_t size)
  480. {
  481. pthread_mutex_init(&to_init->mutex, NULL);
  482. to_init->hdr_buffer.size = size;
  483. to_init->hdr_buffer.data = mw_malloc(size);
  484. if (to_init->hdr_buffer.data == NULL)
  485. return 1;
  486. to_init->hdr_buffer.tail = to_init->hdr_buffer.data;
  487. to_init->hdr_buffer.tail_frag = NULL;
  488. return 0;
  489. }
  490. static void transaction_buffer_destroy(struct transaction_buffer *to_init)
  491. {
  492. buffer_purge(&to_init->hdr_buffer);
  493. pthread_mutex_destroy(&to_init->mutex);
  494. mw_free(to_init->hdr_buffer.data);
  495. }
  496. // Creates transaction
  497. // saves state of buffer before any operation was done
  498. // allowing for rollback if things go wrong
  499. #define transaction_buffer_transaction_start(buf) \
  500. { LOCK_HDR_BUFFER(buf); \
  501. memcpy(&(buf)->state_backup, &(buf)->hdr_buffer, sizeof((buf)->hdr_buffer)); }
  502. #define transaction_buffer_transaction_commit(buf) UNLOCK_HDR_BUFFER(buf);
  503. void transaction_buffer_transaction_rollback(struct transaction_buffer *buf, struct buffer_fragment *frag)
  504. {
  505. memcpy(&buf->hdr_buffer, &buf->state_backup, sizeof(buf->hdr_buffer));
  506. if (buf->hdr_buffer.tail_frag != NULL)
  507. buf->hdr_buffer.tail_frag->next = NULL;
  508. while(frag) {
  509. buffer_frag_free_data(frag);
  510. // we are not actually freeing the structure itself
  511. // just the data it manages
  512. // structure itself is in permanent buffer
  513. // which is locked by HDR_BUFFER lock
  514. frag = frag->next;
  515. }
  516. UNLOCK_HDR_BUFFER(buf);
  517. }
  518. #define TX_ALIASES_INITIALIZE() c_rhash_new(0)
  519. #define RX_ALIASES_INITIALIZE() c_rhash_new(UINT16_MAX >> 8)
  520. struct mqtt_ng_client *mqtt_ng_init(struct mqtt_ng_init *settings)
  521. {
  522. struct mqtt_ng_client *client = mw_calloc(1, sizeof(struct mqtt_ng_client));
  523. if (client == NULL)
  524. return NULL;
  525. if (transaction_buffer_init(&client->main_buffer, HEADER_BUFFER_SIZE))
  526. goto err_free_client;
  527. client->rx_aliases = RX_ALIASES_INITIALIZE();
  528. if (client->rx_aliases == NULL)
  529. goto err_free_trx_buf;
  530. if (pthread_mutex_init(&client->stats_mutex, NULL))
  531. goto err_free_rx_alias;
  532. client->tx_topic_aliases.stoi_dict = TX_ALIASES_INITIALIZE();
  533. if (client->tx_topic_aliases.stoi_dict == NULL)
  534. goto err_free_stats_mutex;
  535. client->tx_topic_aliases.idx_max = UINT16_MAX;
  536. if (pthread_rwlock_init(&client->tx_topic_aliases.rwlock, NULL))
  537. goto err_free_tx_alias;
  538. // TODO just embed the struct into mqtt_ng_client
  539. client->parser.received_data = settings->data_in;
  540. client->send_fnc_ptr = settings->data_out_fnc;
  541. client->user_ctx = settings->user_ctx;
  542. client->log = settings->log;
  543. client->puback_callback = settings->puback_callback;
  544. client->connack_callback = settings->connack_callback;
  545. client->msg_callback = settings->msg_callback;
  546. return client;
  547. err_free_tx_alias:
  548. c_rhash_destroy(client->tx_topic_aliases.stoi_dict);
  549. err_free_stats_mutex:
  550. pthread_mutex_destroy(&client->stats_mutex);
  551. err_free_rx_alias:
  552. c_rhash_destroy(client->rx_aliases);
  553. err_free_trx_buf:
  554. transaction_buffer_destroy(&client->main_buffer);
  555. err_free_client:
  556. mw_free(client);
  557. return NULL;
  558. }
  559. static inline uint8_t get_control_packet_type(uint8_t first_hdr_byte)
  560. {
  561. return first_hdr_byte >> 4;
  562. }
  563. static void mqtt_ng_destroy_rx_alias_hash(c_rhash hash)
  564. {
  565. c_rhash_iter_t i = C_RHASH_ITER_T_INITIALIZER;
  566. uint64_t stored_key;
  567. void *to_free;
  568. while(!c_rhash_iter_uint64_keys(hash, &i, &stored_key)) {
  569. c_rhash_get_ptr_by_uint64(hash, stored_key, &to_free);
  570. mw_free(to_free);
  571. }
  572. c_rhash_destroy(hash);
  573. }
  574. static void mqtt_ng_destroy_tx_alias_hash(c_rhash hash)
  575. {
  576. c_rhash_iter_t i = C_RHASH_ITER_T_INITIALIZER;
  577. const char *stored_key;
  578. void *to_free;
  579. while(!c_rhash_iter_str_keys(hash, &i, &stored_key)) {
  580. c_rhash_get_ptr_by_str(hash, stored_key, &to_free);
  581. mw_free(to_free);
  582. }
  583. c_rhash_destroy(hash);
  584. }
  585. void mqtt_ng_destroy(struct mqtt_ng_client *client)
  586. {
  587. transaction_buffer_destroy(&client->main_buffer);
  588. pthread_mutex_destroy(&client->stats_mutex);
  589. mqtt_ng_destroy_tx_alias_hash(client->tx_topic_aliases.stoi_dict);
  590. pthread_rwlock_destroy(&client->tx_topic_aliases.rwlock);
  591. mqtt_ng_destroy_rx_alias_hash(client->rx_aliases);
  592. mw_free(client);
  593. }
  594. int frag_set_external_data(mqtt_wss_log_ctx_t log, struct buffer_fragment *frag, void *data, size_t data_len, free_fnc_t data_free_fnc)
  595. {
  596. if (frag->len) {
  597. // TODO?: This could potentially be done in future if we set rule
  598. // external data always follows in buffer data
  599. // could help reduce fragmentation in some messages but
  600. // currently not worth it considering time is tight
  601. mws_fatal(log, UNIT_LOG_PREFIX "INTERNAL ERROR: Cannot set external data to fragment already containing in buffer data!");
  602. return 1;
  603. }
  604. switch (ptr2memory_mode(data_free_fnc)) {
  605. case MEMCPY:
  606. frag->data = mw_malloc(data_len);
  607. if (frag->data == NULL) {
  608. mws_error(log, UNIT_LOG_PREFIX "OOM while malloc @_optimized_add");
  609. return 1;
  610. }
  611. memcpy(frag->data, data, data_len);
  612. break;
  613. case EXTERNAL_FREE_AFTER_USE:
  614. case CALLER_RESPONSIBLE:
  615. frag->data = data;
  616. break;
  617. }
  618. frag->free_fnc = data_free_fnc;
  619. frag->len = data_len;
  620. frag->flags |= BUFFER_FRAG_DATA_EXTERNAL;
  621. return 0;
  622. }
  623. // this is fixed part of variable header for connect packet
  624. // mqtt-v5.0-cs1, 3.1.2.1, 2.1.2.2
  625. static const char mqtt_protocol_name_frag[] =
  626. { 0x00, 0x04, 'M', 'Q', 'T', 'T', MQTT_VERSION_5_0 };
  627. #define MQTT_UTF8_STRING_SIZE(string) (2 + strlen(string))
  628. // see 1.5.5
  629. #define MQTT_VARSIZE_INT_BYTES(value) ( value > 2097152 ? 4 : ( value > 16384 ? 3 : ( value > 128 ? 2 : 1 ) ) )
  630. static size_t mqtt_ng_connect_size(struct mqtt_auth_properties *auth,
  631. struct mqtt_lwt_properties *lwt)
  632. {
  633. // First get the size of payload + variable header
  634. size_t size =
  635. + sizeof(mqtt_protocol_name_frag) /* Proto Name and Version */
  636. + 1 /* Connect Flags */
  637. + 2 /* Keep Alive */
  638. + 4 /* 3.1.2.11.1 Property Length - for now fixed to only Topic Alias Maximum, TODO TODO*/;
  639. // CONNECT payload. 3.1.3
  640. if (auth->client_id)
  641. size += MQTT_UTF8_STRING_SIZE(auth->client_id);
  642. if (lwt) {
  643. // 3.1.3.2 will properties TODO TODO
  644. size += 1;
  645. // 3.1.3.3
  646. if (lwt->will_topic)
  647. size += MQTT_UTF8_STRING_SIZE(lwt->will_topic);
  648. // 3.1.3.4 will payload
  649. if (lwt->will_message) {
  650. size += 2 + lwt->will_message_size;
  651. }
  652. }
  653. // 3.1.3.5
  654. if (auth->username)
  655. size += MQTT_UTF8_STRING_SIZE(auth->username);
  656. // 3.1.3.6
  657. if (auth->password)
  658. size += MQTT_UTF8_STRING_SIZE(auth->password);
  659. return size;
  660. }
  661. #define BUFFER_TRANSACTION_NEW_FRAG(buf, flags, frag, on_fail) \
  662. { if(frag==NULL) { \
  663. frag = buffer_new_frag(buf, (flags)); } \
  664. if(frag==NULL) { on_fail; }}
  665. #define CHECK_BYTES_AVAILABLE(buf, needed, fail) \
  666. { if (BUFFER_BYTES_AVAILABLE(buf) < (size_t)needed) { \
  667. fail; } }
  668. #define DATA_ADVANCE(buf, bytes, frag) { size_t b = (bytes); (buf)->tail += b; (frag)->len += b; }
  669. // TODO maybe just user client->buf.tail?
  670. #define WRITE_POS(frag) (&(frag->data[frag->len]))
  671. // [MQTT-1.5.2] Two Byte Integer
  672. #define PACK_2B_INT(buffer, integer, frag) { *(uint16_t *)WRITE_POS(frag) = htobe16((integer)); \
  673. DATA_ADVANCE(buffer, sizeof(uint16_t), frag); }
  674. static int _optimized_add(struct header_buffer *buf, mqtt_wss_log_ctx_t log_ctx, void *data, size_t data_len, free_fnc_t data_free_fnc, struct buffer_fragment **frag)
  675. {
  676. if (data_len > SMALL_STRING_DONT_FRAGMENT_LIMIT) {
  677. buffer_frag_flag_t flags = BUFFER_FRAG_DATA_EXTERNAL;
  678. if ((*frag)->flags & BUFFER_FRAG_GARBAGE_COLLECT_ON_SEND)
  679. flags |= BUFFER_FRAG_GARBAGE_COLLECT_ON_SEND;
  680. if( (*frag = buffer_new_frag(buf, flags)) == NULL ) {
  681. mws_error(log_ctx, "Out of buffer space while generating the message");
  682. return 1;
  683. }
  684. if (frag_set_external_data(log_ctx, *frag, data, data_len, data_free_fnc)) {
  685. mws_error(log_ctx, "Error adding external data to newly created fragment");
  686. return 1;
  687. }
  688. // we dont want to write to this fragment anymore
  689. *frag = NULL;
  690. } else if (data_len) {
  691. // if the data are small dont bother creating new fragments
  692. // store in buffer directly
  693. CHECK_BYTES_AVAILABLE(buf, data_len, return 1);
  694. memcpy(buf->tail, data, data_len);
  695. DATA_ADVANCE(buf, data_len, *frag);
  696. }
  697. return 0;
  698. }
  699. #define TRY_GENERATE_MESSAGE(generator_function, client, ...) \
  700. int rc = generator_function(&client->main_buffer, client->log, ##__VA_ARGS__); \
  701. if (rc == MQTT_NG_MSGGEN_BUFFER_OOM) { \
  702. LOCK_HDR_BUFFER(&client->main_buffer); \
  703. transaction_buffer_garbage_collect((&client->main_buffer), client->log); \
  704. UNLOCK_HDR_BUFFER(&client->main_buffer); \
  705. rc = generator_function(&client->main_buffer, client->log, ##__VA_ARGS__); \
  706. if (rc == MQTT_NG_MSGGEN_BUFFER_OOM && client->max_mem_bytes) { \
  707. LOCK_HDR_BUFFER(&client->main_buffer); \
  708. transaction_buffer_grow((&client->main_buffer), client->log, GROWTH_FACTOR, client->max_mem_bytes); \
  709. UNLOCK_HDR_BUFFER(&client->main_buffer); \
  710. rc = generator_function(&client->main_buffer, client->log, ##__VA_ARGS__); \
  711. } \
  712. if (rc == MQTT_NG_MSGGEN_BUFFER_OOM) \
  713. mws_error(client->log, "%s failed to generate message due to insufficient buffer space (line %d)", __FUNCTION__, __LINE__); \
  714. } \
  715. if (rc == MQTT_NG_MSGGEN_OK) { \
  716. pthread_mutex_lock(&client->stats_mutex); \
  717. client->stats.tx_messages_queued++; \
  718. pthread_mutex_unlock(&client->stats_mutex); \
  719. } \
  720. return rc;
  721. mqtt_msg_data mqtt_ng_generate_connect(struct transaction_buffer *trx_buf,
  722. mqtt_wss_log_ctx_t log_ctx,
  723. struct mqtt_auth_properties *auth,
  724. struct mqtt_lwt_properties *lwt,
  725. uint8_t clean_start,
  726. uint16_t keep_alive)
  727. {
  728. // Sanity Checks First (are given parameters correct and up to MQTT spec)
  729. if (!auth->client_id) {
  730. mws_error(log_ctx, "ClientID must be set. [MQTT-3.1.3-3]");
  731. return NULL;
  732. }
  733. size_t len = strlen(auth->client_id);
  734. if (!len) {
  735. // [MQTT-3.1.3-6] server MAY allow empty client_id and treat it
  736. // as specific client_id (not same as client_id not given)
  737. // however server MUST allow ClientIDs between 1-23 bytes [MQTT-3.1.3-5]
  738. // so we will warn client server might not like this and he is using it
  739. // at his own risk!
  740. mws_warn(log_ctx, "client_id provided is empty string. This might not be allowed by server [MQTT-3.1.3-6]");
  741. }
  742. if(len > MQTT_MAX_CLIENT_ID) {
  743. // [MQTT-3.1.3-5] server MUST allow client_id length 1-32
  744. // server MAY allow longer client_id, if user provides longer client_id
  745. // warn them he is doing so at his own risk!
  746. mws_warn(log_ctx, "client_id provided is longer than 23 bytes, server might not allow that [MQTT-3.1.3-5]");
  747. }
  748. if (lwt) {
  749. if (lwt->will_message && lwt->will_message_size > 65535) {
  750. mws_error(log_ctx, "Will message cannot be longer than 65535 bytes due to MQTT protocol limitations [MQTT-3.1.3-4] and [MQTT-1.5.6]");
  751. return NULL;
  752. }
  753. if (!lwt->will_topic) { //TODO topic given with strlen==0 ? check specs
  754. mws_error(log_ctx, "If will message is given will topic must also be given [MQTT-3.1.3.3]");
  755. return NULL;
  756. }
  757. if (lwt->will_qos > MQTT_MAX_QOS) {
  758. // refer to [MQTT-3-1.2-12]
  759. mws_error(log_ctx, "QOS for LWT message is bigger than max");
  760. return NULL;
  761. }
  762. }
  763. // >> START THE RODEO <<
  764. transaction_buffer_transaction_start(trx_buf);
  765. // Calculate the resulting message size sans fixed MQTT header
  766. size_t size = mqtt_ng_connect_size(auth, lwt);
  767. // Start generating the message
  768. struct buffer_fragment *frag = NULL;
  769. mqtt_msg_data ret = NULL;
  770. BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, BUFFER_FRAG_MQTT_PACKET_HEAD, frag, goto fail_rollback );
  771. ret = frag;
  772. // MQTT Fixed Header
  773. 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*/;
  774. CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, needed_bytes, goto fail_rollback);
  775. *WRITE_POS(frag) = MQTT_CPT_CONNECT << 4;
  776. DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
  777. DATA_ADVANCE(&trx_buf->hdr_buffer, uint32_to_mqtt_vbi(size, WRITE_POS(frag)), frag);
  778. memcpy(WRITE_POS(frag), mqtt_protocol_name_frag, sizeof(mqtt_protocol_name_frag));
  779. DATA_ADVANCE(&trx_buf->hdr_buffer, sizeof(mqtt_protocol_name_frag), frag);
  780. // [MQTT-3.1.2.3] Connect flags
  781. unsigned char *connect_flags = WRITE_POS(frag);
  782. *connect_flags = 0;
  783. if (auth->username)
  784. *connect_flags |= MQTT_CONNECT_FLAG_USERNAME;
  785. if (auth->password)
  786. *connect_flags |= MQTT_CONNECT_FLAG_PASSWORD;
  787. if (lwt) {
  788. *connect_flags |= MQTT_CONNECT_FLAG_LWT;
  789. *connect_flags |= lwt->will_qos << MQTT_CONNECT_FLAG_QOS_BITSHIFT;
  790. if (lwt->will_retain)
  791. *connect_flags |= MQTT_CONNECT_FLAG_LWT_RETAIN;
  792. }
  793. if (clean_start)
  794. *connect_flags |= MQTT_CONNECT_FLAG_CLEAN_START;
  795. DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
  796. PACK_2B_INT(&trx_buf->hdr_buffer, keep_alive, frag);
  797. // TODO Property Length [MQTT-3.1.3.2.1] temporary fixed to 3 (one property topic alias max)
  798. DATA_ADVANCE(&trx_buf->hdr_buffer, uint32_to_mqtt_vbi(3, WRITE_POS(frag)), frag);
  799. *WRITE_POS(frag) = MQTT_PROP_TOPIC_ALIAS_MAX;
  800. DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
  801. PACK_2B_INT(&trx_buf->hdr_buffer, 65535, frag);
  802. // [MQTT-3.1.3.1] Client identifier
  803. CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, 2, goto fail_rollback);
  804. PACK_2B_INT(&trx_buf->hdr_buffer, strlen(auth->client_id), frag);
  805. if (_optimized_add(&trx_buf->hdr_buffer, log_ctx, auth->client_id, strlen(auth->client_id), auth->client_id_free, &frag))
  806. goto fail_rollback;
  807. if (lwt != NULL) {
  808. // Will Properties [MQTT-3.1.3.2]
  809. // TODO for now fixed 0
  810. BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, 0, frag, goto fail_rollback);
  811. CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, 1, goto fail_rollback);
  812. *WRITE_POS(frag) = 0;
  813. DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
  814. // Will Topic [MQTT-3.1.3.3]
  815. CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, 2, goto fail_rollback);
  816. PACK_2B_INT(&trx_buf->hdr_buffer, strlen(lwt->will_topic), frag);
  817. if (_optimized_add(&trx_buf->hdr_buffer, log_ctx, lwt->will_topic, strlen(lwt->will_topic), lwt->will_topic_free, &frag))
  818. goto fail_rollback;
  819. // Will Payload [MQTT-3.1.3.4]
  820. if (lwt->will_message_size) {
  821. BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, 0, frag, goto fail_rollback);
  822. CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, 2, goto fail_rollback);
  823. PACK_2B_INT(&trx_buf->hdr_buffer, lwt->will_message_size, frag);
  824. if (_optimized_add(&trx_buf->hdr_buffer, log_ctx, lwt->will_message, lwt->will_message_size, lwt->will_topic_free, &frag))
  825. goto fail_rollback;
  826. }
  827. }
  828. // [MQTT-3.1.3.5]
  829. if (auth->username) {
  830. BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, 0, frag, goto fail_rollback);
  831. CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, 2, goto fail_rollback);
  832. PACK_2B_INT(&trx_buf->hdr_buffer, strlen(auth->username), frag);
  833. if (_optimized_add(&trx_buf->hdr_buffer, log_ctx, auth->username, strlen(auth->username), auth->username_free, &frag))
  834. goto fail_rollback;
  835. }
  836. // [MQTT-3.1.3.6]
  837. if (auth->password) {
  838. BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, 0, frag, goto fail_rollback);
  839. CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, 2, goto fail_rollback);
  840. PACK_2B_INT(&trx_buf->hdr_buffer, strlen(auth->password), frag);
  841. if (_optimized_add(&trx_buf->hdr_buffer, log_ctx, auth->password, strlen(auth->password), auth->password_free, &frag))
  842. goto fail_rollback;
  843. }
  844. trx_buf->hdr_buffer.tail_frag->flags |= BUFFER_FRAG_MQTT_PACKET_TAIL;
  845. transaction_buffer_transaction_commit(trx_buf);
  846. return ret;
  847. fail_rollback:
  848. transaction_buffer_transaction_rollback(trx_buf, ret);
  849. return NULL;
  850. }
  851. int mqtt_ng_connect(struct mqtt_ng_client *client,
  852. struct mqtt_auth_properties *auth,
  853. struct mqtt_lwt_properties *lwt,
  854. uint8_t clean_start,
  855. uint16_t keep_alive)
  856. {
  857. client->client_state = RAW;
  858. client->parser.state = MQTT_PARSE_FIXED_HEADER_PACKET_TYPE;
  859. LOCK_HDR_BUFFER(&client->main_buffer);
  860. client->main_buffer.sending_frag = NULL;
  861. if (clean_start)
  862. buffer_purge(&client->main_buffer.hdr_buffer);
  863. UNLOCK_HDR_BUFFER(&client->main_buffer);
  864. pthread_rwlock_wrlock(&client->tx_topic_aliases.rwlock);
  865. // according to MQTT spec topic aliases should not be persisted
  866. // even if clean session is true
  867. mqtt_ng_destroy_tx_alias_hash(client->tx_topic_aliases.stoi_dict);
  868. client->tx_topic_aliases.stoi_dict = TX_ALIASES_INITIALIZE();
  869. if (client->tx_topic_aliases.stoi_dict == NULL) {
  870. pthread_rwlock_unlock(&client->tx_topic_aliases.rwlock);
  871. return 1;
  872. }
  873. client->tx_topic_aliases.idx_assigned = 0;
  874. pthread_rwlock_unlock(&client->tx_topic_aliases.rwlock);
  875. mqtt_ng_destroy_rx_alias_hash(client->rx_aliases);
  876. client->rx_aliases = RX_ALIASES_INITIALIZE();
  877. if (client->rx_aliases == NULL)
  878. return 1;
  879. client->connect_msg = mqtt_ng_generate_connect(&client->main_buffer, client->log, auth, lwt, clean_start, keep_alive);
  880. if (client->connect_msg == NULL)
  881. return 1;
  882. pthread_mutex_lock(&client->stats_mutex);
  883. if (clean_start)
  884. client->stats.tx_messages_queued = 1;
  885. else
  886. client->stats.tx_messages_queued++;
  887. client->stats.tx_messages_sent = 0;
  888. client->stats.rx_messages_rcvd = 0;
  889. pthread_mutex_unlock(&client->stats_mutex);
  890. client->client_state = CONNECT_PENDING;
  891. return 0;
  892. }
  893. uint16_t get_unused_packet_id() {
  894. static uint16_t packet_id = 0;
  895. packet_id++;
  896. return packet_id ? packet_id : ++packet_id;
  897. }
  898. static inline size_t mqtt_ng_publish_size(const char *topic,
  899. size_t msg_len,
  900. uint16_t topic_id)
  901. {
  902. size_t retval = 2 /* Topic Name Length */
  903. + (topic == NULL ? 0 : strlen(topic))
  904. + 2 /* Packet identifier */
  905. + 1 /* Properties Length TODO for now fixed to 1 property */
  906. + msg_len;
  907. if (topic_id)
  908. retval += 3;
  909. return retval;
  910. }
  911. int mqtt_ng_generate_publish(struct transaction_buffer *trx_buf,
  912. mqtt_wss_log_ctx_t log_ctx,
  913. char *topic,
  914. free_fnc_t topic_free,
  915. void *msg,
  916. free_fnc_t msg_free,
  917. size_t msg_len,
  918. uint8_t publish_flags,
  919. uint16_t *packet_id,
  920. uint16_t topic_alias)
  921. {
  922. // >> START THE RODEO <<
  923. transaction_buffer_transaction_start(trx_buf);
  924. // Calculate the resulting message size sans fixed MQTT header
  925. size_t size = mqtt_ng_publish_size(topic, msg_len, topic_alias);
  926. // Start generating the message
  927. struct buffer_fragment *frag = NULL;
  928. mqtt_msg_data mqtt_msg = NULL;
  929. BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, BUFFER_FRAG_MQTT_PACKET_HEAD, frag, goto fail_rollback );
  930. // in case of QOS 0 we can garbage collect immediatelly after sending
  931. uint8_t qos = (publish_flags >> 1) & 0x03;
  932. if (!qos)
  933. frag->flags |= BUFFER_FRAG_GARBAGE_COLLECT_ON_SEND;
  934. mqtt_msg = frag;
  935. // MQTT Fixed Header
  936. size_t needed_bytes = 1 /* Packet type */ + MQTT_VARSIZE_INT_BYTES(size) + size - msg_len;
  937. CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, needed_bytes, goto fail_rollback);
  938. *WRITE_POS(frag) = (MQTT_CPT_PUBLISH << 4) | (publish_flags & 0xF);
  939. DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
  940. DATA_ADVANCE(&trx_buf->hdr_buffer, uint32_to_mqtt_vbi(size, WRITE_POS(frag)), frag);
  941. // MQTT Variable Header
  942. // [MQTT-3.3.2.1]
  943. PACK_2B_INT(&trx_buf->hdr_buffer, topic == NULL ? 0 : strlen(topic), frag);
  944. if (topic != NULL) {
  945. if (_optimized_add(&trx_buf->hdr_buffer, log_ctx, topic, strlen(topic), topic_free, &frag))
  946. goto fail_rollback;
  947. BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, 0, frag, goto fail_rollback);
  948. }
  949. // [MQTT-3.3.2.2]
  950. mqtt_msg->packet_id = get_unused_packet_id();
  951. *packet_id = mqtt_msg->packet_id;
  952. PACK_2B_INT(&trx_buf->hdr_buffer, mqtt_msg->packet_id, frag);
  953. // [MQTT-3.3.2.3.1] TODO Property Length for now fixed 0
  954. *WRITE_POS(frag) = topic_alias ? 3 : 0;
  955. DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
  956. if(topic_alias) {
  957. *WRITE_POS(frag) = MQTT_PROP_TOPIC_ALIAS;
  958. DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
  959. PACK_2B_INT(&trx_buf->hdr_buffer, topic_alias, frag);
  960. }
  961. if( (frag = buffer_new_frag(&trx_buf->hdr_buffer, BUFFER_FRAG_DATA_EXTERNAL)) == NULL )
  962. goto fail_rollback;
  963. if (frag_set_external_data(log_ctx, frag, msg, msg_len, msg_free))
  964. goto fail_rollback;
  965. trx_buf->hdr_buffer.tail_frag->flags |= BUFFER_FRAG_MQTT_PACKET_TAIL;
  966. if (!qos)
  967. trx_buf->hdr_buffer.tail_frag->flags |= BUFFER_FRAG_GARBAGE_COLLECT_ON_SEND;
  968. transaction_buffer_transaction_commit(trx_buf);
  969. return MQTT_NG_MSGGEN_OK;
  970. fail_rollback:
  971. transaction_buffer_transaction_rollback(trx_buf, mqtt_msg);
  972. return MQTT_NG_MSGGEN_BUFFER_OOM;
  973. }
  974. #define PUBLISH_SP_SIZE 64
  975. int mqtt_ng_publish(struct mqtt_ng_client *client,
  976. char *topic,
  977. free_fnc_t topic_free,
  978. void *msg,
  979. free_fnc_t msg_free,
  980. size_t msg_len,
  981. uint8_t publish_flags,
  982. uint16_t *packet_id)
  983. {
  984. struct topic_alias_data *alias = NULL;
  985. pthread_rwlock_rdlock(&client->tx_topic_aliases.rwlock);
  986. c_rhash_get_ptr_by_str(client->tx_topic_aliases.stoi_dict, topic, (void**)&alias);
  987. pthread_rwlock_unlock(&client->tx_topic_aliases.rwlock);
  988. uint16_t topic_id = 0;
  989. if (alias != NULL) {
  990. topic_id = alias->idx;
  991. uint32_t cnt = __atomic_fetch_add(&alias->usage_count, 1, __ATOMIC_SEQ_CST);
  992. if (cnt) {
  993. topic = NULL;
  994. topic_free = NULL;
  995. }
  996. }
  997. if (client->max_msg_size && PUBLISH_SP_SIZE + mqtt_ng_publish_size(topic, msg_len, topic_id) > client->max_msg_size) {
  998. mws_error(client->log, "Message too big for server: %zu", msg_len);
  999. return MQTT_NG_MSGGEN_MSG_TOO_BIG;
  1000. }
  1001. TRY_GENERATE_MESSAGE(mqtt_ng_generate_publish, client, topic, topic_free, msg, msg_free, msg_len, publish_flags, packet_id, topic_id);
  1002. }
  1003. static inline size_t mqtt_ng_subscribe_size(struct mqtt_sub *subs, size_t sub_count)
  1004. {
  1005. size_t len = 2 /* Packet Identifier */ + 1 /* Properties Length TODO for now fixed 0 */;
  1006. len += sub_count * (2 /* topic filter string length */ + 1 /* [MQTT-3.8.3.1] Subscription Options Byte */);
  1007. for (size_t i = 0; i < sub_count; i++) {
  1008. len += strlen(subs[i].topic);
  1009. }
  1010. return len;
  1011. }
  1012. int mqtt_ng_generate_subscribe(struct transaction_buffer *trx_buf, mqtt_wss_log_ctx_t log_ctx, struct mqtt_sub *subs, size_t sub_count)
  1013. {
  1014. // >> START THE RODEO <<
  1015. transaction_buffer_transaction_start(trx_buf);
  1016. // Calculate the resulting message size sans fixed MQTT header
  1017. size_t size = mqtt_ng_subscribe_size(subs, sub_count);
  1018. // Start generating the message
  1019. struct buffer_fragment *frag = NULL;
  1020. mqtt_msg_data ret = NULL;
  1021. BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, BUFFER_FRAG_MQTT_PACKET_HEAD, frag, goto fail_rollback);
  1022. ret = frag;
  1023. // MQTT Fixed Header
  1024. size_t needed_bytes = 1 /* Packet type */ + MQTT_VARSIZE_INT_BYTES(size) + 3 /*Packet ID + Property Length*/;
  1025. CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, needed_bytes, goto fail_rollback);
  1026. *WRITE_POS(frag) = (MQTT_CPT_SUBSCRIBE << 4) | 0x2 /* [MQTT-3.8.1-1] */;
  1027. DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
  1028. DATA_ADVANCE(&trx_buf->hdr_buffer, uint32_to_mqtt_vbi(size, WRITE_POS(frag)), frag);
  1029. // MQTT Variable Header
  1030. // [MQTT-3.8.2] PacketID
  1031. ret->packet_id = get_unused_packet_id();
  1032. PACK_2B_INT(&trx_buf->hdr_buffer, ret->packet_id, frag);
  1033. // [MQTT-3.8.2.1.1] Property Length // TODO for now fixed 0
  1034. *WRITE_POS(frag) = 0;
  1035. DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
  1036. for (size_t i = 0; i < sub_count; i++) {
  1037. BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, 0, frag, goto fail_rollback);
  1038. PACK_2B_INT(&trx_buf->hdr_buffer, strlen(subs[i].topic), frag);
  1039. if (_optimized_add(&trx_buf->hdr_buffer, log_ctx, subs[i].topic, strlen(subs[i].topic), subs[i].topic_free, &frag))
  1040. goto fail_rollback;
  1041. BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, 0, frag, goto fail_rollback);
  1042. *WRITE_POS(frag) = subs[i].options;
  1043. DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
  1044. }
  1045. trx_buf->hdr_buffer.tail_frag->flags |= BUFFER_FRAG_MQTT_PACKET_TAIL;
  1046. transaction_buffer_transaction_commit(trx_buf);
  1047. return MQTT_NG_MSGGEN_OK;
  1048. fail_rollback:
  1049. transaction_buffer_transaction_rollback(trx_buf, ret);
  1050. return MQTT_NG_MSGGEN_BUFFER_OOM;
  1051. }
  1052. int mqtt_ng_subscribe(struct mqtt_ng_client *client, struct mqtt_sub *subs, size_t sub_count)
  1053. {
  1054. TRY_GENERATE_MESSAGE(mqtt_ng_generate_subscribe, client, subs, sub_count);
  1055. }
  1056. int mqtt_ng_generate_disconnect(struct transaction_buffer *trx_buf, mqtt_wss_log_ctx_t log_ctx, uint8_t reason_code)
  1057. {
  1058. (void) log_ctx;
  1059. // >> START THE RODEO <<
  1060. transaction_buffer_transaction_start(trx_buf);
  1061. // Calculate the resulting message size sans fixed MQTT header
  1062. size_t size = reason_code ? 1 : 0;
  1063. // Start generating the message
  1064. struct buffer_fragment *frag = NULL;
  1065. mqtt_msg_data ret = NULL;
  1066. BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, BUFFER_FRAG_MQTT_PACKET_HEAD, frag, goto fail_rollback);
  1067. ret = frag;
  1068. // MQTT Fixed Header
  1069. size_t needed_bytes = 1 /* Packet type */ + MQTT_VARSIZE_INT_BYTES(size) + (reason_code ? 1 : 0);
  1070. CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, needed_bytes, goto fail_rollback);
  1071. *WRITE_POS(frag) = MQTT_CPT_DISCONNECT << 4;
  1072. DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
  1073. DATA_ADVANCE(&trx_buf->hdr_buffer, uint32_to_mqtt_vbi(size, WRITE_POS(frag)), frag);
  1074. if (reason_code) {
  1075. // MQTT Variable Header
  1076. // [MQTT-3.14.2.1] PacketID
  1077. *WRITE_POS(frag) = reason_code;
  1078. DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
  1079. }
  1080. trx_buf->hdr_buffer.tail_frag->flags |= BUFFER_FRAG_MQTT_PACKET_TAIL;
  1081. transaction_buffer_transaction_commit(trx_buf);
  1082. return MQTT_NG_MSGGEN_OK;
  1083. fail_rollback:
  1084. transaction_buffer_transaction_rollback(trx_buf, ret);
  1085. return MQTT_NG_MSGGEN_BUFFER_OOM;
  1086. }
  1087. int mqtt_ng_disconnect(struct mqtt_ng_client *client, uint8_t reason_code)
  1088. {
  1089. TRY_GENERATE_MESSAGE(mqtt_ng_generate_disconnect, client, reason_code);
  1090. }
  1091. static int mqtt_generate_puback(struct transaction_buffer *trx_buf, mqtt_wss_log_ctx_t log_ctx, uint16_t packet_id, uint8_t reason_code)
  1092. {
  1093. (void) log_ctx;
  1094. // >> START THE RODEO <<
  1095. transaction_buffer_transaction_start(trx_buf);
  1096. // Calculate the resulting message size sans fixed MQTT header
  1097. size_t size = 2 /* Packet ID */ + (reason_code ? 1 : 0) /* reason code */;
  1098. // Start generating the message
  1099. struct buffer_fragment *frag = NULL;
  1100. BUFFER_TRANSACTION_NEW_FRAG(&trx_buf->hdr_buffer, BUFFER_FRAG_MQTT_PACKET_HEAD | BUFFER_FRAG_GARBAGE_COLLECT_ON_SEND, frag, goto fail_rollback);
  1101. // MQTT Fixed Header
  1102. size_t needed_bytes = 1 /* Packet type */ + MQTT_VARSIZE_INT_BYTES(size) + size;
  1103. CHECK_BYTES_AVAILABLE(&trx_buf->hdr_buffer, needed_bytes, goto fail_rollback);
  1104. *WRITE_POS(frag) = MQTT_CPT_PUBACK << 4;
  1105. DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
  1106. DATA_ADVANCE(&trx_buf->hdr_buffer, uint32_to_mqtt_vbi(size, WRITE_POS(frag)), frag);
  1107. // MQTT Variable Header
  1108. PACK_2B_INT(&trx_buf->hdr_buffer, packet_id, frag);
  1109. if (reason_code) {
  1110. // MQTT Variable Header
  1111. // [MQTT-3.14.2.1] PacketID
  1112. *WRITE_POS(frag) = reason_code;
  1113. DATA_ADVANCE(&trx_buf->hdr_buffer, 1, frag);
  1114. }
  1115. trx_buf->hdr_buffer.tail_frag->flags |= BUFFER_FRAG_MQTT_PACKET_TAIL;
  1116. transaction_buffer_transaction_commit(trx_buf);
  1117. return MQTT_NG_MSGGEN_OK;
  1118. fail_rollback:
  1119. transaction_buffer_transaction_rollback(trx_buf, frag);
  1120. return MQTT_NG_MSGGEN_BUFFER_OOM;
  1121. }
  1122. static int mqtt_ng_puback(struct mqtt_ng_client *client, uint16_t packet_id, uint8_t reason_code)
  1123. {
  1124. TRY_GENERATE_MESSAGE(mqtt_generate_puback, client, packet_id, reason_code);
  1125. }
  1126. int mqtt_ng_ping(struct mqtt_ng_client *client)
  1127. {
  1128. client->ping_pending = 1;
  1129. return MQTT_NG_MSGGEN_OK;
  1130. }
  1131. #define MQTT_NG_CLIENT_NEED_MORE_BYTES 0x10
  1132. #define MQTT_NG_CLIENT_MQTT_PACKET_DONE 0x11
  1133. #define MQTT_NG_CLIENT_PARSE_DONE 0x12
  1134. #define MQTT_NG_CLIENT_WANT_WRITE 0x13
  1135. #define MQTT_NG_CLIENT_OK_CALL_AGAIN 0
  1136. #define MQTT_NG_CLIENT_PROTOCOL_ERROR -1
  1137. #define MQTT_NG_CLIENT_SERVER_RETURNED_ERROR -2
  1138. #define MQTT_NG_CLIENT_NOT_IMPL_YET -3
  1139. #define MQTT_NG_CLIENT_OOM -4
  1140. #define MQTT_NG_CLIENT_INTERNAL_ERROR -5
  1141. #define BUF_READ_CHECK_AT_LEAST(buf, x) \
  1142. if (rbuf_bytes_available(buf) < (x)) \
  1143. return MQTT_NG_CLIENT_NEED_MORE_BYTES;
  1144. #define vbi_parser_reset_ctx(ctx) memset(ctx, 0, sizeof(struct mqtt_vbi_parser_ctx))
  1145. static int vbi_parser_parse(struct mqtt_vbi_parser_ctx *ctx, rbuf_t data, mqtt_wss_log_ctx_t log)
  1146. {
  1147. if (ctx->bytes > MQTT_VBI_MAXBYTES - 1) {
  1148. mws_error(log, "MQTT Variable Byte Integer can't be longer than %d bytes", MQTT_VBI_MAXBYTES);
  1149. return MQTT_NG_CLIENT_PROTOCOL_ERROR;
  1150. }
  1151. if (!ctx->bytes || ctx->data[ctx->bytes-1] & MQTT_VBI_CONTINUATION_FLAG) {
  1152. BUF_READ_CHECK_AT_LEAST(data, 1);
  1153. ctx->bytes++;
  1154. rbuf_pop(data, &ctx->data[ctx->bytes-1], 1);
  1155. if ( ctx->data[ctx->bytes-1] & MQTT_VBI_CONTINUATION_FLAG )
  1156. return MQTT_NG_CLIENT_OK_CALL_AGAIN;
  1157. }
  1158. if (mqtt_vbi_to_uint32(ctx->data, &ctx->result)) {
  1159. mws_error(log, "MQTT Variable Byte Integer failed to be parsed.");
  1160. return MQTT_NG_CLIENT_PROTOCOL_ERROR;
  1161. }
  1162. return MQTT_NG_CLIENT_PARSE_DONE;
  1163. }
  1164. static void mqtt_properties_parser_ctx_reset(struct mqtt_properties_parser_ctx *ctx)
  1165. {
  1166. ctx->state = PROPERTIES_LENGTH;
  1167. while (ctx->head) {
  1168. struct mqtt_property *f = ctx->head;
  1169. ctx->head = ctx->head->next;
  1170. if (f->type == MQTT_TYPE_STR || f->type == MQTT_TYPE_STR_PAIR)
  1171. mw_free(f->data.strings[0]);
  1172. if (f->type == MQTT_TYPE_STR_PAIR)
  1173. mw_free(f->data.strings[1]);
  1174. if (f->type == MQTT_TYPE_BIN)
  1175. mw_free(f->data.bindata);
  1176. mw_free(f);
  1177. }
  1178. ctx->tail = NULL;
  1179. ctx->properties_length = 0;
  1180. ctx->bytes_consumed = 0;
  1181. vbi_parser_reset_ctx(&ctx->vbi_parser_ctx);
  1182. }
  1183. struct mqtt_property_type {
  1184. uint8_t id;
  1185. enum mqtt_datatype datatype;
  1186. const char* name;
  1187. };
  1188. const struct mqtt_property_type mqtt_property_types[] = {
  1189. { .id = MQTT_PROP_TOPIC_ALIAS, .name = MQTT_PROP_TOPIC_ALIAS_NAME, .datatype = MQTT_TYPE_UINT_16 },
  1190. { .id = MQTT_PROP_PAYLOAD_FMT_INDICATOR, .name = MQTT_PROP_PAYLOAD_FMT_INDICATOR_NAME, .datatype = MQTT_TYPE_UINT_8 },
  1191. { .id = MQTT_PROP_MSG_EXPIRY_INTERVAL, .name = MQTT_PROP_MSG_EXPIRY_INTERVAL_NAME, .datatype = MQTT_TYPE_UINT_32 },
  1192. { .id = MQTT_PROP_CONTENT_TYPE, .name = MQTT_PROP_CONTENT_TYPE_NAME, .datatype = MQTT_TYPE_STR },
  1193. { .id = MQTT_PROP_RESPONSE_TOPIC, .name = MQTT_PROP_RESPONSE_TOPIC_NAME, .datatype = MQTT_TYPE_STR },
  1194. { .id = MQTT_PROP_CORRELATION_DATA, .name = MQTT_PROP_CORRELATION_DATA_NAME, .datatype = MQTT_TYPE_BIN },
  1195. { .id = MQTT_PROP_SUB_IDENTIFIER, .name = MQTT_PROP_SUB_IDENTIFIER_NAME, .datatype = MQTT_TYPE_VBI },
  1196. { .id = MQTT_PROP_SESSION_EXPIRY_INTERVAL, .name = MQTT_PROP_SESSION_EXPIRY_INTERVAL_NAME, .datatype = MQTT_TYPE_UINT_32 },
  1197. { .id = MQTT_PROP_ASSIGNED_CLIENT_ID, .name = MQTT_PROP_ASSIGNED_CLIENT_ID_NAME, .datatype = MQTT_TYPE_STR },
  1198. { .id = MQTT_PROP_SERVER_KEEP_ALIVE, .name = MQTT_PROP_SERVER_KEEP_ALIVE_NAME, .datatype = MQTT_TYPE_UINT_16 },
  1199. { .id = MQTT_PROP_AUTH_METHOD, .name = MQTT_PROP_AUTH_METHOD_NAME, .datatype = MQTT_TYPE_STR },
  1200. { .id = MQTT_PROP_AUTH_DATA, .name = MQTT_PROP_AUTH_DATA_NAME, .datatype = MQTT_TYPE_BIN },
  1201. { .id = MQTT_PROP_REQ_PROBLEM_INFO, .name = MQTT_PROP_REQ_PROBLEM_INFO_NAME, .datatype = MQTT_TYPE_UINT_8 },
  1202. { .id = MQTT_PROP_WILL_DELAY_INTERVAL, .name = MQTT_PROP_WIIL_DELAY_INTERVAL_NAME, .datatype = MQTT_TYPE_UINT_32 },
  1203. { .id = MQTT_PROP_REQ_RESP_INFORMATION, .name = MQTT_PROP_REQ_RESP_INFORMATION_NAME, .datatype = MQTT_TYPE_UINT_8 },
  1204. { .id = MQTT_PROP_RESP_INFORMATION, .name = MQTT_PROP_RESP_INFORMATION_NAME, .datatype = MQTT_TYPE_STR },
  1205. { .id = MQTT_PROP_SERVER_REF, .name = MQTT_PROP_SERVER_REF_NAME, .datatype = MQTT_TYPE_STR },
  1206. { .id = MQTT_PROP_REASON_STR, .name = MQTT_PROP_REASON_STR_NAME, .datatype = MQTT_TYPE_STR },
  1207. { .id = MQTT_PROP_RECEIVE_MAX, .name = MQTT_PROP_RECEIVE_MAX_NAME, .datatype = MQTT_TYPE_UINT_16 },
  1208. { .id = MQTT_PROP_TOPIC_ALIAS_MAX, .name = MQTT_PROP_TOPIC_ALIAS_MAX_NAME, .datatype = MQTT_TYPE_UINT_16 },
  1209. // MQTT_PROP_TOPIC_ALIAS is first as it is most often used
  1210. { .id = MQTT_PROP_MAX_QOS, .name = MQTT_PROP_MAX_QOS_NAME, .datatype = MQTT_TYPE_UINT_8 },
  1211. { .id = MQTT_PROP_RETAIN_AVAIL, .name = MQTT_PROP_RETAIN_AVAIL_NAME, .datatype = MQTT_TYPE_UINT_8 },
  1212. { .id = MQTT_PROP_USR, .name = MQTT_PROP_USR_NAME, .datatype = MQTT_TYPE_STR_PAIR },
  1213. { .id = MQTT_PROP_MAX_PKT_SIZE, .name = MQTT_PROP_MAX_PKT_SIZE_NAME, .datatype = MQTT_TYPE_UINT_32 },
  1214. { .id = MQTT_PROP_WILDCARD_SUB_AVAIL, .name = MQTT_PROP_WILDCARD_SUB_AVAIL_NAME, .datatype = MQTT_TYPE_UINT_8 },
  1215. { .id = MQTT_PROP_SUB_ID_AVAIL, .name = MQTT_PROP_SUB_ID_AVAIL_NAME, .datatype = MQTT_TYPE_UINT_8 },
  1216. { .id = MQTT_PROP_SHARED_SUB_AVAIL, .name = MQTT_PROP_SHARED_SUB_AVAIL_NAME, .datatype = MQTT_TYPE_UINT_8 },
  1217. { .id = 0, .name = NULL, .datatype = MQTT_TYPE_UNKNOWN }
  1218. };
  1219. static int get_property_type_by_id(uint8_t property_id) {
  1220. for (int i = 0; mqtt_property_types[i].datatype != MQTT_TYPE_UNKNOWN; i++) {
  1221. if (mqtt_property_types[i].id == property_id)
  1222. return mqtt_property_types[i].datatype;
  1223. }
  1224. return MQTT_TYPE_UNKNOWN;
  1225. }
  1226. struct mqtt_property *get_property_by_id(struct mqtt_property *props, uint8_t property_id)
  1227. {
  1228. while (props) {
  1229. if (props->id == property_id) {
  1230. return props;
  1231. }
  1232. props = props->next;
  1233. }
  1234. return NULL;
  1235. }
  1236. // Parses [MQTT-2.2.2]
  1237. static int parse_properties_array(struct mqtt_properties_parser_ctx *ctx, rbuf_t data, mqtt_wss_log_ctx_t log)
  1238. {
  1239. int rc;
  1240. switch (ctx->state) {
  1241. case PROPERTIES_LENGTH:
  1242. rc = vbi_parser_parse(&ctx->vbi_parser_ctx, data, log);
  1243. if (rc == MQTT_NG_CLIENT_PARSE_DONE) {
  1244. ctx->properties_length = ctx->vbi_parser_ctx.result;
  1245. ctx->bytes_consumed += ctx->vbi_parser_ctx.bytes;
  1246. ctx->vbi_length = ctx->vbi_parser_ctx.bytes;
  1247. if (!ctx->properties_length)
  1248. return MQTT_NG_CLIENT_PARSE_DONE;
  1249. ctx->state = PROPERTY_CREATE;
  1250. break;
  1251. }
  1252. return rc;
  1253. case PROPERTY_CREATE:
  1254. BUF_READ_CHECK_AT_LEAST(data, 1);
  1255. struct mqtt_property *prop = mw_calloc(1, sizeof(struct mqtt_property));
  1256. if (ctx->head == NULL) {
  1257. ctx->head = prop;
  1258. ctx->tail = prop;
  1259. } else {
  1260. ctx->tail->next = prop;
  1261. ctx->tail = ctx->tail->next;
  1262. }
  1263. ctx->state = PROPERTY_ID;
  1264. /* FALLTHROUGH */
  1265. case PROPERTY_ID:
  1266. rbuf_pop(data, (char*)&ctx->tail->id, 1);
  1267. ctx->bytes_consumed += 1;
  1268. ctx->tail->type = get_property_type_by_id(ctx->tail->id);
  1269. switch (ctx->tail->type) {
  1270. case MQTT_TYPE_UINT_16:
  1271. ctx->state = PROPERTY_TYPE_UINT16;
  1272. break;
  1273. case MQTT_TYPE_UINT_32:
  1274. ctx->state = PROPERTY_TYPE_UINT32;
  1275. break;
  1276. case MQTT_TYPE_UINT_8:
  1277. ctx->state = PROPERTY_TYPE_UINT8;
  1278. break;
  1279. case MQTT_TYPE_VBI:
  1280. ctx->state = PROPERTY_TYPE_VBI;
  1281. vbi_parser_reset_ctx(&ctx->vbi_parser_ctx);
  1282. break;
  1283. case MQTT_TYPE_STR:
  1284. case MQTT_TYPE_STR_PAIR:
  1285. ctx->str_idx = 0;
  1286. /* FALLTHROUGH */
  1287. case MQTT_TYPE_BIN:
  1288. ctx->state = PROPERTY_TYPE_STR_BIN_LEN;
  1289. break;
  1290. default:
  1291. mws_error(log, "Unsupported property type %d for property id %d.", (int)ctx->tail->type, (int)ctx->tail->id);
  1292. return MQTT_NG_CLIENT_PROTOCOL_ERROR;
  1293. }
  1294. break;
  1295. case PROPERTY_TYPE_STR_BIN_LEN:
  1296. BUF_READ_CHECK_AT_LEAST(data, sizeof(uint16_t));
  1297. rbuf_pop(data, (char*)&ctx->tail->bindata_len, sizeof(uint16_t));
  1298. ctx->tail->bindata_len = be16toh(ctx->tail->bindata_len);
  1299. ctx->bytes_consumed += 2;
  1300. switch (ctx->tail->type) {
  1301. case MQTT_TYPE_BIN:
  1302. ctx->state = PROPERTY_TYPE_BIN;
  1303. break;
  1304. case MQTT_TYPE_STR:
  1305. case MQTT_TYPE_STR_PAIR:
  1306. ctx->state = PROPERTY_TYPE_STR;
  1307. break;
  1308. default:
  1309. mws_error(log, "Unexpected datatype in PROPERTY_TYPE_STR_BIN_LEN %d", (int)ctx->tail->type);
  1310. return MQTT_NG_CLIENT_INTERNAL_ERROR;
  1311. }
  1312. break;
  1313. case PROPERTY_TYPE_STR:
  1314. BUF_READ_CHECK_AT_LEAST(data, ctx->tail->bindata_len);
  1315. ctx->tail->data.strings[ctx->str_idx] = mw_malloc(ctx->tail->bindata_len + 1);
  1316. rbuf_pop(data, ctx->tail->data.strings[ctx->str_idx], ctx->tail->bindata_len);
  1317. ctx->tail->data.strings[ctx->str_idx][ctx->tail->bindata_len] = 0;
  1318. ctx->str_idx++;
  1319. ctx->bytes_consumed += ctx->tail->bindata_len;
  1320. if (ctx->tail->type == MQTT_TYPE_STR_PAIR && ctx->str_idx < 2) {
  1321. ctx->state = PROPERTY_TYPE_STR_BIN_LEN;
  1322. break;
  1323. }
  1324. ctx->state = PROPERTY_NEXT;
  1325. break;
  1326. case PROPERTY_TYPE_BIN:
  1327. BUF_READ_CHECK_AT_LEAST(data, ctx->tail->bindata_len);
  1328. ctx->tail->data.bindata = mw_malloc(ctx->tail->bindata_len);
  1329. rbuf_pop(data, ctx->tail->data.bindata, ctx->tail->bindata_len);
  1330. ctx->bytes_consumed += ctx->tail->bindata_len;
  1331. ctx->state = PROPERTY_NEXT;
  1332. break;
  1333. case PROPERTY_TYPE_VBI:
  1334. rc = vbi_parser_parse(&ctx->vbi_parser_ctx, data, log);
  1335. if (rc == MQTT_NG_CLIENT_PARSE_DONE) {
  1336. ctx->tail->data.uint32 = ctx->vbi_parser_ctx.result;
  1337. ctx->bytes_consumed += ctx->vbi_parser_ctx.bytes;
  1338. ctx->state = PROPERTY_NEXT;
  1339. break;
  1340. }
  1341. return rc;
  1342. case PROPERTY_TYPE_UINT8:
  1343. BUF_READ_CHECK_AT_LEAST(data, sizeof(uint8_t));
  1344. rbuf_pop(data, (char*)&ctx->tail->data.uint8, sizeof(uint8_t));
  1345. ctx->bytes_consumed += sizeof(uint8_t);
  1346. ctx->state = PROPERTY_NEXT;
  1347. break;
  1348. case PROPERTY_TYPE_UINT32:
  1349. BUF_READ_CHECK_AT_LEAST(data, sizeof(uint32_t));
  1350. rbuf_pop(data, (char*)&ctx->tail->data.uint32, sizeof(uint32_t));
  1351. ctx->tail->data.uint32 = be32toh(ctx->tail->data.uint32);
  1352. ctx->bytes_consumed += sizeof(uint32_t);
  1353. ctx->state = PROPERTY_NEXT;
  1354. break;
  1355. case PROPERTY_TYPE_UINT16:
  1356. BUF_READ_CHECK_AT_LEAST(data, sizeof(uint16_t));
  1357. rbuf_pop(data, (char*)&ctx->tail->data.uint16, sizeof(uint16_t));
  1358. ctx->tail->data.uint16 = be16toh(ctx->tail->data.uint16);
  1359. ctx->bytes_consumed += sizeof(uint16_t);
  1360. ctx->state = PROPERTY_NEXT;
  1361. /* FALLTHROUGH */
  1362. case PROPERTY_NEXT:
  1363. if (ctx->properties_length > ctx->bytes_consumed - ctx->vbi_length) {
  1364. ctx->state = PROPERTY_CREATE;
  1365. break;
  1366. } else
  1367. return MQTT_NG_CLIENT_PARSE_DONE;
  1368. }
  1369. return MQTT_NG_CLIENT_OK_CALL_AGAIN;
  1370. }
  1371. static int parse_connack_varhdr(struct mqtt_ng_client *client)
  1372. {
  1373. struct mqtt_ng_parser *parser = &client->parser;
  1374. switch (parser->varhdr_state) {
  1375. case MQTT_PARSE_VARHDR_INITIAL:
  1376. BUF_READ_CHECK_AT_LEAST(parser->received_data, 2);
  1377. rbuf_pop(parser->received_data, (char*)&parser->mqtt_packet.connack.flags, 1);
  1378. rbuf_pop(parser->received_data, (char*)&parser->mqtt_packet.connack.reason_code, 1);
  1379. parser->varhdr_state = MQTT_PARSE_VARHDR_PROPS;
  1380. mqtt_properties_parser_ctx_reset(&parser->properties_parser);
  1381. break;
  1382. case MQTT_PARSE_VARHDR_PROPS:
  1383. return parse_properties_array(&parser->properties_parser, parser->received_data, client->log);
  1384. default:
  1385. ERROR("invalid state for connack varhdr parser");
  1386. return MQTT_NG_CLIENT_INTERNAL_ERROR;
  1387. }
  1388. return MQTT_NG_CLIENT_OK_CALL_AGAIN;
  1389. }
  1390. static int parse_disconnect_varhdr(struct mqtt_ng_client *client)
  1391. {
  1392. struct mqtt_ng_parser *parser = &client->parser;
  1393. switch (parser->varhdr_state) {
  1394. case MQTT_PARSE_VARHDR_INITIAL:
  1395. if (!parser->mqtt_fixed_hdr_remaining_length) {
  1396. // [MQTT-3.14.2.1] if reason code omitted act same as == 0
  1397. parser->mqtt_packet.disconnect.reason_code = 0;
  1398. return MQTT_NG_CLIENT_PARSE_DONE;
  1399. }
  1400. BUF_READ_CHECK_AT_LEAST(parser->received_data, 1);
  1401. rbuf_pop(parser->received_data, (char*)&parser->mqtt_packet.connack.reason_code, 1);
  1402. if (parser->mqtt_fixed_hdr_remaining_length == 1)
  1403. return MQTT_NG_CLIENT_PARSE_DONE;
  1404. parser->varhdr_state = MQTT_PARSE_VARHDR_PROPS;
  1405. mqtt_properties_parser_ctx_reset(&parser->properties_parser);
  1406. break;
  1407. case MQTT_PARSE_VARHDR_PROPS:
  1408. return parse_properties_array(&parser->properties_parser, parser->received_data, client->log);
  1409. default:
  1410. ERROR("invalid state for connack varhdr parser");
  1411. return MQTT_NG_CLIENT_INTERNAL_ERROR;
  1412. }
  1413. return MQTT_NG_CLIENT_OK_CALL_AGAIN;
  1414. }
  1415. static int parse_puback_varhdr(struct mqtt_ng_client *client)
  1416. {
  1417. struct mqtt_ng_parser *parser = &client->parser;
  1418. switch (parser->varhdr_state) {
  1419. case MQTT_PARSE_VARHDR_INITIAL:
  1420. BUF_READ_CHECK_AT_LEAST(parser->received_data, 2);
  1421. rbuf_pop(parser->received_data, (char*)&parser->mqtt_packet.puback.packet_id, 2);
  1422. parser->mqtt_packet.puback.packet_id = be16toh(parser->mqtt_packet.puback.packet_id);
  1423. if (parser->mqtt_fixed_hdr_remaining_length < 3) {
  1424. // [MQTT-3.4.2.1] if length is not big enough for reason code
  1425. // it is omitted and handled same as if it was present and == 0
  1426. // initially missed this detail and was wondering WTF is going on (sigh)
  1427. parser->mqtt_packet.puback.reason_code = 0;
  1428. return MQTT_NG_CLIENT_PARSE_DONE;
  1429. }
  1430. parser->varhdr_state = MQTT_PARSE_VARHDR_OPTIONAL_REASON_CODE;
  1431. /* FALLTHROUGH */
  1432. case MQTT_PARSE_VARHDR_OPTIONAL_REASON_CODE:
  1433. BUF_READ_CHECK_AT_LEAST(parser->received_data, 1);
  1434. rbuf_pop(parser->received_data, (char*)&parser->mqtt_packet.puback.reason_code, 1);
  1435. // LOL so in CONNACK you have to have 0 byte to
  1436. // signify empty properties list
  1437. // but in PUBACK it can be omitted if remaining length doesn't allow it (sigh)
  1438. if (parser->mqtt_fixed_hdr_remaining_length < 4)
  1439. return MQTT_NG_CLIENT_PARSE_DONE;
  1440. parser->varhdr_state = MQTT_PARSE_VARHDR_PROPS;
  1441. mqtt_properties_parser_ctx_reset(&parser->properties_parser);
  1442. /* FALLTHROUGH */
  1443. case MQTT_PARSE_VARHDR_PROPS:
  1444. return parse_properties_array(&parser->properties_parser, parser->received_data, client->log);
  1445. default:
  1446. ERROR("invalid state for puback varhdr parser");
  1447. return MQTT_NG_CLIENT_INTERNAL_ERROR;
  1448. }
  1449. return MQTT_NG_CLIENT_OK_CALL_AGAIN;
  1450. }
  1451. static int parse_suback_varhdr(struct mqtt_ng_client *client)
  1452. {
  1453. int rc;
  1454. size_t avail;
  1455. struct mqtt_ng_parser *parser = &client->parser;
  1456. struct mqtt_suback *suback = &client->parser.mqtt_packet.suback;
  1457. switch (parser->varhdr_state) {
  1458. case MQTT_PARSE_VARHDR_INITIAL:
  1459. suback->reason_codes = NULL;
  1460. BUF_READ_CHECK_AT_LEAST(parser->received_data, 2);
  1461. rbuf_pop(parser->received_data, (char*)&suback->packet_id, 2);
  1462. suback->packet_id = be16toh(suback->packet_id);
  1463. parser->varhdr_state = MQTT_PARSE_VARHDR_PROPS;
  1464. parser->mqtt_parsed_len = 2;
  1465. mqtt_properties_parser_ctx_reset(&parser->properties_parser);
  1466. /* FALLTHROUGH */
  1467. case MQTT_PARSE_VARHDR_PROPS:
  1468. rc = parse_properties_array(&parser->properties_parser, parser->received_data, client->log);
  1469. if (rc != MQTT_NG_CLIENT_PARSE_DONE)
  1470. return rc;
  1471. parser->mqtt_parsed_len += parser->properties_parser.bytes_consumed;
  1472. suback->reason_code_count = parser->mqtt_fixed_hdr_remaining_length - parser->mqtt_parsed_len;
  1473. suback->reason_codes = mw_calloc(suback->reason_code_count, sizeof(*suback->reason_codes));
  1474. suback->reason_codes_pending = suback->reason_code_count;
  1475. parser->varhdr_state = MQTT_PARSE_REASONCODES;
  1476. /* FALLTHROUGH */
  1477. case MQTT_PARSE_REASONCODES:
  1478. avail = rbuf_bytes_available(parser->received_data);
  1479. if (avail < 1)
  1480. return MQTT_NG_CLIENT_NEED_MORE_BYTES;
  1481. suback->reason_codes_pending -= rbuf_pop(parser->received_data, (char*)suback->reason_codes, MIN(suback->reason_codes_pending, avail));
  1482. if (!suback->reason_codes_pending)
  1483. return MQTT_NG_CLIENT_PARSE_DONE;
  1484. return MQTT_NG_CLIENT_NEED_MORE_BYTES;
  1485. default:
  1486. ERROR("invalid state for suback varhdr parser");
  1487. return MQTT_NG_CLIENT_INTERNAL_ERROR;
  1488. }
  1489. return MQTT_NG_CLIENT_OK_CALL_AGAIN;
  1490. }
  1491. static int parse_publish_varhdr(struct mqtt_ng_client *client)
  1492. {
  1493. int rc;
  1494. struct mqtt_ng_parser *parser = &client->parser;
  1495. struct mqtt_publish *publish = &client->parser.mqtt_packet.publish;
  1496. switch (parser->varhdr_state) {
  1497. case MQTT_PARSE_VARHDR_INITIAL:
  1498. BUF_READ_CHECK_AT_LEAST(parser->received_data, 2);
  1499. publish->topic = NULL;
  1500. publish->qos = ((parser->mqtt_control_packet_type >> 1) & 0x03);
  1501. rbuf_pop(parser->received_data, (char*)&publish->topic_len, 2);
  1502. publish->topic_len = be16toh(publish->topic_len);
  1503. parser->mqtt_parsed_len = 2;
  1504. if (!publish->topic_len) {
  1505. parser->varhdr_state = MQTT_PARSE_VARHDR_POST_TOPICNAME;
  1506. break;
  1507. }
  1508. publish->topic = mw_calloc(1, publish->topic_len + 1 /* add 0x00 */);
  1509. if (publish->topic == NULL)
  1510. return MQTT_NG_CLIENT_OOM;
  1511. parser->varhdr_state = MQTT_PARSE_VARHDR_TOPICNAME;
  1512. /* FALLTHROUGH */
  1513. case MQTT_PARSE_VARHDR_TOPICNAME:
  1514. // TODO check empty topic can be valid? In which case we have to skip this step
  1515. BUF_READ_CHECK_AT_LEAST(parser->received_data, publish->topic_len);
  1516. rbuf_pop(parser->received_data, publish->topic, publish->topic_len);
  1517. parser->mqtt_parsed_len += publish->topic_len;
  1518. parser->varhdr_state = MQTT_PARSE_VARHDR_POST_TOPICNAME;
  1519. /* FALLTHROUGH */
  1520. case MQTT_PARSE_VARHDR_POST_TOPICNAME:
  1521. mqtt_properties_parser_ctx_reset(&parser->properties_parser);
  1522. if (!publish->qos) { // PacketID present only for QOS > 0 [MQTT-3.3.2.2]
  1523. parser->varhdr_state = MQTT_PARSE_VARHDR_PROPS;
  1524. break;
  1525. }
  1526. parser->varhdr_state = MQTT_PARSE_VARHDR_PACKET_ID;
  1527. /* FALLTHROUGH */
  1528. case MQTT_PARSE_VARHDR_PACKET_ID:
  1529. BUF_READ_CHECK_AT_LEAST(parser->received_data, 2);
  1530. rbuf_pop(parser->received_data, (char*)&publish->packet_id, 2);
  1531. publish->packet_id = be16toh(publish->packet_id);
  1532. parser->varhdr_state = MQTT_PARSE_VARHDR_PROPS;
  1533. parser->mqtt_parsed_len += 2;
  1534. /* FALLTHROUGH */
  1535. case MQTT_PARSE_VARHDR_PROPS:
  1536. rc = parse_properties_array(&parser->properties_parser, parser->received_data, client->log);
  1537. if (rc != MQTT_NG_CLIENT_PARSE_DONE)
  1538. return rc;
  1539. parser->mqtt_parsed_len += parser->properties_parser.bytes_consumed;
  1540. parser->varhdr_state = MQTT_PARSE_PAYLOAD;
  1541. /* FALLTHROUGH */
  1542. case MQTT_PARSE_PAYLOAD:
  1543. if (parser->mqtt_fixed_hdr_remaining_length < parser->mqtt_parsed_len) {
  1544. mw_free(publish->topic);
  1545. publish->topic = NULL;
  1546. ERROR("Error parsing PUBLISH message");
  1547. return MQTT_NG_CLIENT_PROTOCOL_ERROR;
  1548. }
  1549. publish->data_len = parser->mqtt_fixed_hdr_remaining_length - parser->mqtt_parsed_len;
  1550. if (!publish->data_len) {
  1551. publish->data = NULL;
  1552. return MQTT_NG_CLIENT_PARSE_DONE; // 0 length payload is OK [MQTT-3.3.3]
  1553. }
  1554. BUF_READ_CHECK_AT_LEAST(parser->received_data, publish->data_len);
  1555. publish->data = mw_malloc(publish->data_len);
  1556. if (publish->data == NULL) {
  1557. mw_free(publish->topic);
  1558. publish->topic = NULL;
  1559. return MQTT_NG_CLIENT_OOM;
  1560. }
  1561. rbuf_pop(parser->received_data, publish->data, publish->data_len);
  1562. parser->mqtt_parsed_len += publish->data_len;
  1563. return MQTT_NG_CLIENT_PARSE_DONE;
  1564. default:
  1565. ERROR("invalid state for publish varhdr parser");
  1566. return MQTT_NG_CLIENT_INTERNAL_ERROR;
  1567. }
  1568. return MQTT_NG_CLIENT_OK_CALL_AGAIN;
  1569. }
  1570. // TODO move to separate file, dont send whole client pointer just to be able
  1571. // to access LOG context send parser only which should include log
  1572. static int parse_data(struct mqtt_ng_client *client)
  1573. {
  1574. int rc;
  1575. struct mqtt_ng_parser *parser = &client->parser;
  1576. switch(parser->state) {
  1577. case MQTT_PARSE_FIXED_HEADER_PACKET_TYPE:
  1578. BUF_READ_CHECK_AT_LEAST(parser->received_data, 1);
  1579. rbuf_pop(parser->received_data, (char*)&parser->mqtt_control_packet_type, 1);
  1580. vbi_parser_reset_ctx(&parser->vbi_parser);
  1581. parser->state = MQTT_PARSE_FIXED_HEADER_LEN;
  1582. break;
  1583. case MQTT_PARSE_FIXED_HEADER_LEN:
  1584. rc = vbi_parser_parse(&parser->vbi_parser, parser->received_data, client->log);
  1585. if (rc == MQTT_NG_CLIENT_PARSE_DONE) {
  1586. parser->mqtt_fixed_hdr_remaining_length = parser->vbi_parser.result;
  1587. parser->state = MQTT_PARSE_VARIABLE_HEADER;
  1588. parser->varhdr_state = MQTT_PARSE_VARHDR_INITIAL;
  1589. break;
  1590. }
  1591. return rc;
  1592. case MQTT_PARSE_VARIABLE_HEADER:
  1593. switch (get_control_packet_type(parser->mqtt_control_packet_type)) {
  1594. case MQTT_CPT_CONNACK:
  1595. rc = parse_connack_varhdr(client);
  1596. if (rc == MQTT_NG_CLIENT_PARSE_DONE) {
  1597. parser->state = MQTT_PARSE_MQTT_PACKET_DONE;
  1598. break;
  1599. }
  1600. return rc;
  1601. case MQTT_CPT_PUBACK:
  1602. rc = parse_puback_varhdr(client);
  1603. if (rc == MQTT_NG_CLIENT_PARSE_DONE) {
  1604. parser->state = MQTT_PARSE_MQTT_PACKET_DONE;
  1605. break;
  1606. }
  1607. return rc;
  1608. case MQTT_CPT_SUBACK:
  1609. rc = parse_suback_varhdr(client);
  1610. if (rc != MQTT_NG_CLIENT_NEED_MORE_BYTES && rc != MQTT_NG_CLIENT_OK_CALL_AGAIN) {
  1611. mw_free(parser->mqtt_packet.suback.reason_codes);
  1612. }
  1613. if (rc == MQTT_NG_CLIENT_PARSE_DONE) {
  1614. parser->state = MQTT_PARSE_MQTT_PACKET_DONE;
  1615. break;
  1616. }
  1617. return rc;
  1618. case MQTT_CPT_PUBLISH:
  1619. rc = parse_publish_varhdr(client);
  1620. if (rc == MQTT_NG_CLIENT_PARSE_DONE) {
  1621. parser->state = MQTT_PARSE_MQTT_PACKET_DONE;
  1622. break;
  1623. }
  1624. return rc;
  1625. case MQTT_CPT_PINGRESP:
  1626. if (parser->mqtt_fixed_hdr_remaining_length) {
  1627. ERROR ("PINGRESP has to be 0 Remaining Length."); // [MQTT-3.13.1]
  1628. return MQTT_NG_CLIENT_PROTOCOL_ERROR;
  1629. }
  1630. parser->state = MQTT_PARSE_MQTT_PACKET_DONE;
  1631. break;
  1632. case MQTT_CPT_DISCONNECT:
  1633. rc = parse_disconnect_varhdr(client);
  1634. if (rc == MQTT_NG_CLIENT_PARSE_DONE) {
  1635. parser->state = MQTT_PARSE_MQTT_PACKET_DONE;
  1636. break;
  1637. }
  1638. return rc;
  1639. default:
  1640. ERROR("Parsing Control Packet Type %" PRIu8 " not implemented yet.", get_control_packet_type(parser->mqtt_control_packet_type));
  1641. rbuf_bump_tail(parser->received_data, parser->mqtt_fixed_hdr_remaining_length);
  1642. parser->state = MQTT_PARSE_MQTT_PACKET_DONE;
  1643. return MQTT_NG_CLIENT_NOT_IMPL_YET;
  1644. }
  1645. // we could also return MQTT_NG_CLIENT_OK_CALL_AGAIN
  1646. // and be called again later
  1647. /* FALLTHROUGH */
  1648. case MQTT_PARSE_MQTT_PACKET_DONE:
  1649. parser->state = MQTT_PARSE_FIXED_HEADER_PACKET_TYPE;
  1650. return MQTT_NG_CLIENT_MQTT_PACKET_DONE;
  1651. }
  1652. return MQTT_NG_CLIENT_OK_CALL_AGAIN;
  1653. }
  1654. // set next MQTT fragment to send
  1655. // return 1 if nothing to send
  1656. // return -1 on error
  1657. // return 0 if there is fragment set
  1658. static int mqtt_ng_next_to_send(struct mqtt_ng_client *client) {
  1659. if (client->client_state == CONNECT_PENDING) {
  1660. client->main_buffer.sending_frag = client->connect_msg;
  1661. client->client_state = CONNECTING;
  1662. return 0;
  1663. }
  1664. if (client->client_state != CONNECTED)
  1665. return -1;
  1666. struct buffer_fragment *frag = BUFFER_FIRST_FRAG(&client->main_buffer.hdr_buffer);
  1667. while (frag) {
  1668. if ( frag->sent != frag->len )
  1669. break;
  1670. frag = frag->next;
  1671. }
  1672. if ( client->ping_pending && (!frag || (frag->flags & BUFFER_FRAG_MQTT_PACKET_HEAD && frag->sent == 0)) ) {
  1673. client->ping_pending = 0;
  1674. ping_frag.sent = 0;
  1675. client->main_buffer.sending_frag = &ping_frag;
  1676. return 0;
  1677. }
  1678. client->main_buffer.sending_frag = frag;
  1679. return frag == NULL ? 1 : 0;
  1680. }
  1681. // send current fragment
  1682. // return 0 if whole remaining length could be sent as a whole
  1683. // return -1 if send buffer was filled and
  1684. // nothing could be written anymore
  1685. // return 1 if last fragment of a message was fully sent
  1686. static int send_fragment(struct mqtt_ng_client *client) {
  1687. struct buffer_fragment *frag = client->main_buffer.sending_frag;
  1688. // for readability
  1689. unsigned char *ptr = frag->data + frag->sent;
  1690. size_t bytes = frag->len - frag->sent;
  1691. size_t processed = 0;
  1692. if (bytes)
  1693. processed = client->send_fnc_ptr(client->user_ctx, ptr, bytes);
  1694. else
  1695. WARN("This fragment was fully sent already. This should not happen!");
  1696. frag->sent += processed;
  1697. if (frag->sent != frag->len)
  1698. return -1;
  1699. if (frag->flags & BUFFER_FRAG_MQTT_PACKET_TAIL) {
  1700. client->time_of_last_send = time(NULL);
  1701. pthread_mutex_lock(&client->stats_mutex);
  1702. if (client->main_buffer.sending_frag != &ping_frag)
  1703. client->stats.tx_messages_queued--;
  1704. client->stats.tx_messages_sent++;
  1705. pthread_mutex_unlock(&client->stats_mutex);
  1706. client->main_buffer.sending_frag = NULL;
  1707. return 1;
  1708. }
  1709. client->main_buffer.sending_frag = frag->next;
  1710. return 0;
  1711. }
  1712. // attempt sending all fragments of current single MQTT packet
  1713. static int send_all_message_fragments(struct mqtt_ng_client *client) {
  1714. int rc;
  1715. while ( !(rc = send_fragment(client)) );
  1716. return rc;
  1717. }
  1718. static void try_send_all(struct mqtt_ng_client *client) {
  1719. do {
  1720. if (client->main_buffer.sending_frag == NULL && mqtt_ng_next_to_send(client))
  1721. return;
  1722. } while(send_all_message_fragments(client) >= 0);
  1723. }
  1724. static inline void mark_message_for_gc(struct buffer_fragment *frag)
  1725. {
  1726. while (frag) {
  1727. frag->flags |= BUFFER_FRAG_GARBAGE_COLLECT;
  1728. buffer_frag_free_data(frag);
  1729. if (frag->flags & BUFFER_FRAG_MQTT_PACKET_TAIL)
  1730. return;
  1731. frag = frag->next;
  1732. }
  1733. }
  1734. static int mark_packet_acked(struct mqtt_ng_client *client, uint16_t packet_id)
  1735. {
  1736. LOCK_HDR_BUFFER(&client->main_buffer);
  1737. struct buffer_fragment *frag = BUFFER_FIRST_FRAG(&client->main_buffer.hdr_buffer);
  1738. while (frag) {
  1739. if ( (frag->flags & BUFFER_FRAG_MQTT_PACKET_HEAD) && frag->packet_id == packet_id) {
  1740. if (!frag->sent) {
  1741. ERROR("Received packet_id (%" PRIu16 ") belongs to MQTT packet which was not yet sent!", packet_id);
  1742. UNLOCK_HDR_BUFFER(&client->main_buffer);
  1743. return 1;
  1744. }
  1745. mark_message_for_gc(frag);
  1746. UNLOCK_HDR_BUFFER(&client->main_buffer);
  1747. return 0;
  1748. }
  1749. frag = frag->next;
  1750. }
  1751. ERROR("Received packet_id (%" PRIu16 ") is unknown!", packet_id);
  1752. UNLOCK_HDR_BUFFER(&client->main_buffer);
  1753. return 1;
  1754. }
  1755. int handle_incoming_traffic(struct mqtt_ng_client *client)
  1756. {
  1757. int rc;
  1758. struct mqtt_publish *pub;
  1759. while( (rc = parse_data(client)) == MQTT_NG_CLIENT_OK_CALL_AGAIN );
  1760. if ( rc == MQTT_NG_CLIENT_MQTT_PACKET_DONE ) {
  1761. struct mqtt_property *prop;
  1762. #ifdef MQTT_DEBUG_VERBOSE
  1763. DEBUG("MQTT Packet Parsed Successfully!");
  1764. #endif
  1765. pthread_mutex_lock(&client->stats_mutex);
  1766. client->stats.rx_messages_rcvd++;
  1767. pthread_mutex_unlock(&client->stats_mutex);
  1768. switch (get_control_packet_type(client->parser.mqtt_control_packet_type)) {
  1769. case MQTT_CPT_CONNACK:
  1770. #ifdef MQTT_DEBUG_VERBOSE
  1771. DEBUG("Received CONNACK");
  1772. #endif
  1773. LOCK_HDR_BUFFER(&client->main_buffer);
  1774. mark_message_for_gc(client->connect_msg);
  1775. UNLOCK_HDR_BUFFER(&client->main_buffer);
  1776. client->connect_msg = NULL;
  1777. if (client->client_state != CONNECTING) {
  1778. ERROR("Received unexpected CONNACK");
  1779. client->client_state = ERROR;
  1780. return MQTT_NG_CLIENT_PROTOCOL_ERROR;
  1781. }
  1782. if ((prop = get_property_by_id(client->parser.properties_parser.head, MQTT_PROP_MAX_PKT_SIZE)) != NULL) {
  1783. INFO("MQTT server limits message size to %" PRIu32, prop->data.uint32);
  1784. client->max_msg_size = prop->data.uint32;
  1785. }
  1786. if (client->connack_callback)
  1787. client->connack_callback(client->user_ctx, client->parser.mqtt_packet.connack.reason_code);
  1788. if (!client->parser.mqtt_packet.connack.reason_code) {
  1789. INFO("MQTT Connection Accepted By Server");
  1790. client->client_state = CONNECTED;
  1791. break;
  1792. }
  1793. client->client_state = ERROR;
  1794. return MQTT_NG_CLIENT_SERVER_RETURNED_ERROR;
  1795. case MQTT_CPT_PUBACK:
  1796. #ifdef MQTT_DEBUG_VERBOSE
  1797. DEBUG("Received PUBACK %" PRIu16, client->parser.mqtt_packet.puback.packet_id);
  1798. #endif
  1799. if (mark_packet_acked(client, client->parser.mqtt_packet.puback.packet_id))
  1800. return MQTT_NG_CLIENT_PROTOCOL_ERROR;
  1801. if (client->puback_callback)
  1802. client->puback_callback(client->parser.mqtt_packet.puback.packet_id);
  1803. break;
  1804. case MQTT_CPT_PINGRESP:
  1805. #ifdef MQTT_DEBUG_VERBOSE
  1806. DEBUG("Received PINGRESP");
  1807. #endif
  1808. break;
  1809. case MQTT_CPT_SUBACK:
  1810. #ifdef MQTT_DEBUG_VERBOSE
  1811. DEBUG("Received SUBACK %" PRIu16, client->parser.mqtt_packet.suback.packet_id);
  1812. #endif
  1813. if (mark_packet_acked(client, client->parser.mqtt_packet.suback.packet_id))
  1814. return MQTT_NG_CLIENT_PROTOCOL_ERROR;
  1815. break;
  1816. case MQTT_CPT_PUBLISH:
  1817. #ifdef MQTT_DEBUG_VERBOSE
  1818. DEBUG("Recevied PUBLISH");
  1819. #endif
  1820. pub = &client->parser.mqtt_packet.publish;
  1821. if (pub->qos > 1) {
  1822. mw_free(pub->topic);
  1823. mw_free(pub->data);
  1824. return MQTT_NG_CLIENT_NOT_IMPL_YET;
  1825. }
  1826. if ( pub->qos == 1 && (rc = mqtt_ng_puback(client, pub->packet_id, 0)) ) {
  1827. client->client_state = ERROR;
  1828. ERROR("Error generating PUBACK reply for PUBLISH");
  1829. return rc;
  1830. }
  1831. if ( (prop = get_property_by_id(client->parser.properties_parser.head, MQTT_PROP_TOPIC_ALIAS)) != NULL ) {
  1832. // Topic Alias property was sent from server
  1833. void *topic_ptr;
  1834. if (!c_rhash_get_ptr_by_uint64(client->rx_aliases, prop->data.uint8, &topic_ptr)) {
  1835. if (pub->topic != NULL) {
  1836. ERROR("We do not yet support topic alias reassignment");
  1837. return MQTT_NG_CLIENT_NOT_IMPL_YET;
  1838. }
  1839. pub->topic = topic_ptr;
  1840. } else {
  1841. if (pub->topic == NULL) {
  1842. ERROR("Topic alias with id %d unknown and topic not set by server!", prop->data.uint8);
  1843. return MQTT_NG_CLIENT_PROTOCOL_ERROR;
  1844. }
  1845. c_rhash_insert_uint64_ptr(client->rx_aliases, prop->data.uint8, pub->topic);
  1846. }
  1847. }
  1848. if (client->msg_callback)
  1849. client->msg_callback(pub->topic, pub->data, pub->data_len, pub->qos);
  1850. // in case we have property topic alias and we have topic we take over the string
  1851. // and add pointer to it into topic alias list
  1852. if (prop == NULL)
  1853. mw_free(pub->topic);
  1854. mw_free(pub->data);
  1855. return MQTT_NG_CLIENT_WANT_WRITE;
  1856. case MQTT_CPT_DISCONNECT:
  1857. INFO ("Got MQTT DISCONNECT control packet from server. Reason code: %d", (int)client->parser.mqtt_packet.disconnect.reason_code);
  1858. client->client_state = DISCONNECTED;
  1859. break;
  1860. }
  1861. }
  1862. return rc;
  1863. }
  1864. int mqtt_ng_sync(struct mqtt_ng_client *client)
  1865. {
  1866. if (client->client_state == RAW || client->client_state == DISCONNECTED)
  1867. return 0;
  1868. if (client->client_state == ERROR)
  1869. return 1;
  1870. LOCK_HDR_BUFFER(&client->main_buffer);
  1871. try_send_all(client);
  1872. UNLOCK_HDR_BUFFER(&client->main_buffer);
  1873. int rc;
  1874. while ((rc = handle_incoming_traffic(client)) != MQTT_NG_CLIENT_NEED_MORE_BYTES) {
  1875. if (rc < 0)
  1876. break;
  1877. if (rc == MQTT_NG_CLIENT_WANT_WRITE) {
  1878. LOCK_HDR_BUFFER(&client->main_buffer);
  1879. try_send_all(client);
  1880. UNLOCK_HDR_BUFFER(&client->main_buffer);
  1881. }
  1882. }
  1883. if (rc < 0)
  1884. return rc;
  1885. return 0;
  1886. }
  1887. time_t mqtt_ng_last_send_time(struct mqtt_ng_client *client)
  1888. {
  1889. return client->time_of_last_send;
  1890. }
  1891. void mqtt_ng_set_max_mem(struct mqtt_ng_client *client, size_t bytes)
  1892. {
  1893. client->max_mem_bytes = bytes;
  1894. }
  1895. void mqtt_ng_get_stats(struct mqtt_ng_client *client, struct mqtt_ng_stats *stats)
  1896. {
  1897. pthread_mutex_lock(&client->stats_mutex);
  1898. memcpy(stats, &client->stats, sizeof(struct mqtt_ng_stats));
  1899. pthread_mutex_unlock(&client->stats_mutex);
  1900. stats->tx_bytes_queued = 0;
  1901. stats->tx_buffer_reclaimable = 0;
  1902. LOCK_HDR_BUFFER(&client->main_buffer);
  1903. stats->tx_buffer_used = BUFFER_BYTES_USED(&client->main_buffer.hdr_buffer);
  1904. stats->tx_buffer_free = BUFFER_BYTES_AVAILABLE(&client->main_buffer.hdr_buffer);
  1905. stats->tx_buffer_size = client->main_buffer.hdr_buffer.size;
  1906. struct buffer_fragment *frag = BUFFER_FIRST_FRAG(&client->main_buffer.hdr_buffer);
  1907. while (frag) {
  1908. stats->tx_bytes_queued += frag->len - frag->sent;
  1909. if (frag_is_marked_for_gc(frag))
  1910. stats->tx_buffer_reclaimable += FRAG_SIZE_IN_BUFFER(frag);
  1911. frag = frag->next;
  1912. }
  1913. UNLOCK_HDR_BUFFER(&client->main_buffer);
  1914. }
  1915. int mqtt_ng_set_topic_alias(struct mqtt_ng_client *client, const char *topic)
  1916. {
  1917. uint16_t idx;
  1918. pthread_rwlock_wrlock(&client->tx_topic_aliases.rwlock);
  1919. if (client->tx_topic_aliases.idx_assigned >= client->tx_topic_aliases.idx_max) {
  1920. pthread_rwlock_unlock(&client->tx_topic_aliases.rwlock);
  1921. mws_error(client->log, "Tx topic alias indexes were exhausted (current version of the library doesn't support reassigning yet. Feel free to contribute.");
  1922. return 0; //0 is not a valid topic alias
  1923. }
  1924. struct topic_alias_data *alias;
  1925. if (!c_rhash_get_ptr_by_str(client->tx_topic_aliases.stoi_dict, topic, (void**)&alias)) {
  1926. // this is not a problem for library but might be helpful to warn user
  1927. // as it might indicate bug in their program (but also might be expected)
  1928. idx = alias->idx;
  1929. pthread_rwlock_unlock(&client->tx_topic_aliases.rwlock);
  1930. mws_debug(client->log, "%s topic \"%s\" already has alias set. Ignoring.", __FUNCTION__, topic);
  1931. return idx;
  1932. }
  1933. alias = mw_malloc(sizeof(struct topic_alias_data));
  1934. idx = ++client->tx_topic_aliases.idx_assigned;
  1935. alias->idx = idx;
  1936. __atomic_store_n(&alias->usage_count, 0, __ATOMIC_SEQ_CST);
  1937. c_rhash_insert_str_ptr(client->tx_topic_aliases.stoi_dict, topic, (void*)alias);
  1938. pthread_rwlock_unlock(&client->tx_topic_aliases.rwlock);
  1939. return idx;
  1940. }