ym310_mqtt.c 34 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067
  1. /*
  2. * YM310 MQTT Dual-Instance Management - FIXED v7.1.0
  3. *
  4. * CRITICAL FIXES:
  5. * - Fixed buffer overflow in urc_qmtrecv_handler
  6. * - Added strict bounds checking on all string operations
  7. * - Added memory alignment for buffer structures
  8. * - Added extensive debug logging
  9. * - Fixed race condition in data_ready flag handling
  10. */
  11. #define DBG_TAG "Y-MQ "
  12. #define DBG_LVL DBG_INFO
  13. #include <rtdbg.h>
  14. #include "ym310_mqtt.h"
  15. #include <string.h>
  16. #include <stdio.h>
  17. #include <rtthread.h>
  18. #include <rthw.h> /* Provides rt_hw_interrupt_disable/enable */
  19. #include <ctype.h> /* Provides isdigit */
  20. #include <easyflash.h> /* Provides ef_get_env_blob */
  21. #include "ym310_mqtt.h"
  22. #include "ym310_protecl.h" /* Provides SYS_INFO and save_sys_info declarations */
  23. /*=============================================================================
  24. * Global Variables
  25. *=============================================================================*/
  26. mqtt_instance_t g_mqtt_instances[MQTT_INSTANCE_MAX];
  27. rt_mutex_t g_mqtt_mutex = RT_NULL;
  28. /* Receive buffers - populated by URC handler */
  29. mqtt_recv_buf_t g_mqtt_recv_bufs[MQTT_INSTANCE_MAX];
  30. /* Callback handlers for each instance */
  31. static mqtt_recv_handler_t g_mqtt_callbacks[MQTT_INSTANCE_MAX] = {RT_NULL, RT_NULL};
  32. rt_thread_t g_mqtt_thread = RT_NULL;
  33. static volatile rt_bool_t s_mqtt_thread_running = RT_FALSE;
  34. static rt_uint8_t s_msg_id = 1;
  35. volatile rt_tick_t g_last_mqtt_success_tick = 0;
  36. /* Connection retry timing */
  37. static rt_tick_t s_last_connect_attempt[MQTT_INSTANCE_MAX] = {0};
  38. uint8_t read_platform_flag(void);
  39. #define CONNECT_RETRY_BASE_MS 5000
  40. #define CONNECT_RETRY_MAX_MS 60000
  41. /*=============================================================================
  42. * DEBUG: IPC Object Integrity Check
  43. *=============================================================================*/
  44. extern rt_mutex_t s_at_mutex; // Added at file top
  45. // static void check_ipc_integrity(const char *location)
  46. //{
  47. // extern rt_sem_t s_cmd_sem;
  48. // extern rt_sem_t s_rx_sem;
  49. //
  50. // // LOG_D("IPC_CHECK [%s]: Checking IPC objects...", location);
  51. //
  52. // if (g_mqtt_mutex != RT_NULL)
  53. // {
  54. // rt_object_t obj = (rt_object_t)g_mqtt_mutex;
  55. // rt_uint8_t type = obj->type & ~RT_Object_Class_Static;
  56. // if (type != RT_Object_Class_Mutex)
  57. // {
  58. // LOG_E("IPC_CHECK [%s]: g_mqtt_mutex CORRUPTED! type=%d", location, type);
  59. // }
  60. // }
  61. //
  62. // if (s_at_mutex != RT_NULL)
  63. // {
  64. // rt_object_t obj = (rt_object_t)s_at_mutex;
  65. // rt_uint8_t type = obj->type & ~RT_Object_Class_Static;
  66. // if (type != RT_Object_Class_Mutex)
  67. // {
  68. // LOG_E("IPC_CHECK [%s]: s_at_mutex CORRUPTED! type=%d", location, type);
  69. // }
  70. // }
  71. //
  72. // if (s_cmd_sem != RT_NULL)
  73. // {
  74. // rt_object_t obj = (rt_object_t)s_cmd_sem;
  75. // rt_uint8_t type = obj->type & ~RT_Object_Class_Static;
  76. // if (type != RT_Object_Class_Semaphore)
  77. // {
  78. // LOG_E("IPC_CHECK [%s]: s_cmd_sem CORRUPTED! type=%d", location, type);
  79. // }
  80. // }
  81. //
  82. // if (s_rx_sem != RT_NULL)
  83. // {
  84. // rt_object_t obj = (rt_object_t)s_rx_sem;
  85. // rt_uint8_t type = obj->type & ~RT_Object_Class_Static;
  86. // if (type != RT_Object_Class_Semaphore)
  87. // {
  88. // LOG_E("IPC_CHECK [%s]: s_rx_sem CORRUPTED! type=%d", location, type);
  89. // }
  90. // }
  91. // }
  92. static void urc_qmtrecv_handler(const char *urc)
  93. {
  94. int client_idx = -1, msgid = 0;
  95. char topic[MQTT_MAX_TOPIC_LEN] = {0};
  96. int topic_len = 0;
  97. const char *payload_start = NULL;
  98. int payload_len = 0;
  99. int declared_len = 0;
  100. if (urc == RT_NULL)
  101. return;
  102. /* Find +QMTRECV: prefix */
  103. const char *p = rt_strstr(urc, "+QMTRECV:");
  104. if (p == RT_NULL)
  105. return;
  106. p += 9;
  107. /* Skip whitespace */
  108. while (*p == ' ' || *p == '\t')
  109. p++;
  110. /* Parse client_idx and msgid */
  111. if (sscanf(p, "%d,%d", &client_idx, &msgid) != 2)
  112. return;
  113. if (client_idx < 0 || client_idx >= MQTT_INSTANCE_MAX)
  114. return;
  115. /* Locate topic start quote */
  116. const char *topic_start = rt_strstr(p, "\"");
  117. if (topic_start == RT_NULL)
  118. return;
  119. topic_start++;
  120. const char *topic_end = rt_strstr(topic_start, "\"");
  121. if (topic_end == RT_NULL || topic_end == topic_start)
  122. return;
  123. topic_len = topic_end - topic_start;
  124. if (topic_len >= MQTT_MAX_TOPIC_LEN)
  125. topic_len = MQTT_MAX_TOPIC_LEN - 1;
  126. rt_memcpy(topic, topic_start, topic_len);
  127. topic[topic_len] = '\0';
  128. /* Locate payload */
  129. const char *after_topic = topic_end + 1;
  130. while (*after_topic == ' ' || *after_topic == '\t')
  131. after_topic++;
  132. if (*after_topic != ',')
  133. return;
  134. after_topic++;
  135. while (*after_topic == ' ' || *after_topic == '\t')
  136. after_topic++;
  137. if (*after_topic == '"')
  138. {
  139. /* No length field format: payload directly after comma */
  140. payload_start = after_topic; // Point to start quote
  141. declared_len = 0; // No length, use last quote to locate
  142. }
  143. else if (isdigit(*after_topic))
  144. {
  145. /* Has length field format */
  146. declared_len = atoi(after_topic);
  147. while (*after_topic && *after_topic != ',')
  148. after_topic++;
  149. if (*after_topic != ',')
  150. return;
  151. after_topic++;
  152. while (*after_topic == ' ' || *after_topic == '\t')
  153. after_topic++;
  154. if (*after_topic != '"')
  155. return;
  156. payload_start = after_topic;
  157. }
  158. else
  159. {
  160. return;
  161. }
  162. /* Determine actual payload length */
  163. payload_start++; // Skip start quote, point to first char
  164. if (declared_len > 0)
  165. {
  166. /* Has length field: use length field value, not rely on quotes */
  167. payload_len = declared_len;
  168. if (payload_len >= MQTT_MAX_PAYLOAD_SIZE)
  169. {
  170. LOG_W("URC: Payload too large (declared=%d >= %d), truncating",
  171. payload_len, MQTT_MAX_PAYLOAD_SIZE);
  172. payload_len = MQTT_MAX_PAYLOAD_SIZE - 1;
  173. }
  174. }
  175. else
  176. {
  177. /* No length field: find last quote as end */
  178. const char *payload_end = strrchr(payload_start, '"');
  179. if (payload_end == RT_NULL)
  180. {
  181. LOG_E("URC: No closing quote found for payload");
  182. return;
  183. }
  184. payload_len = payload_end - payload_start;
  185. if (payload_len >= MQTT_MAX_PAYLOAD_SIZE)
  186. {
  187. LOG_W("URC: Payload too large (%d >= %d), truncating",
  188. payload_len, MQTT_MAX_PAYLOAD_SIZE);
  189. payload_len = MQTT_MAX_PAYLOAD_SIZE - 1;
  190. }
  191. }
  192. LOG_D("URC: MQTT[%d] RX: topic=%s, len=%d", client_idx, topic, payload_len);
  193. /* ========== Disable interrupts to write recv buffer ========== */
  194. // rt_base_t //level = rt_hw_interrupt_disable();
  195. if (g_mqtt_recv_bufs[client_idx].data_ready)
  196. {
  197. // rt_hw_interrupt_enable(level);
  198. LOG_W("MQTT[%d] RX buffer busy, dropping msg", client_idx);
  199. return;
  200. }
  201. /* Copy topic */
  202. rt_strncpy(g_mqtt_recv_bufs[client_idx].topic, topic, MQTT_MAX_TOPIC_LEN - 1);
  203. g_mqtt_recv_bufs[client_idx].topic[MQTT_MAX_TOPIC_LEN - 1] = '\0';
  204. /* Copy payload */
  205. if (payload_len > 0)
  206. {
  207. rt_memcpy(g_mqtt_recv_bufs[client_idx].payload, payload_start, payload_len);
  208. }
  209. g_mqtt_recv_bufs[client_idx].payload[payload_len] = '\0';
  210. g_mqtt_recv_bufs[client_idx].payload_len = payload_len;
  211. g_mqtt_recv_bufs[client_idx].data_ready = RT_TRUE;
  212. /* Update statistics */
  213. g_mqtt_instances[client_idx].rx_count++;
  214. g_mqtt_instances[client_idx].last_activity = rt_tick_get();
  215. // rt_hw_interrupt_enable(level);
  216. LOG_D("URC: Data stored in recv buffer[%d]", client_idx);
  217. }
  218. static void urc_qmtstat_handler(const char *urc)
  219. {
  220. int client_idx, err_code;
  221. sscanf(urc, "+QMTSTAT: %d,%d", &client_idx, &err_code);
  222. if (client_idx >= 0 && client_idx < MQTT_INSTANCE_MAX)
  223. {
  224. LOG_W("MQTT[%d] status: err=%d", client_idx, err_code);
  225. if (err_code != 0)
  226. {
  227. // rt_base_t //level = rt_hw_interrupt_disable();
  228. g_mqtt_instances[client_idx].state = MQTT_STATE_DISCONNECTED;
  229. g_mqtt_instances[client_idx].err_count++;
  230. // rt_hw_interrupt_enable(level);
  231. }
  232. }
  233. }
  234. static void urc_qmtconn_handler(const char *urc)
  235. {
  236. int client_idx, result, ret_code;
  237. sscanf(urc, "+QMTCONN: %d,%d,%d", &client_idx, &result, &ret_code);
  238. LOG_D("MQTT[%d] conn: result=%d, ret=%d", client_idx, result, ret_code);
  239. }
  240. static void urc_qmtopen_handler(const char *urc)
  241. {
  242. int client_idx, result;
  243. sscanf(urc, "+QMTOPEN: %d,%d", &client_idx, &result);
  244. LOG_D("MQTT[%d] open: result=%d", client_idx, result);
  245. }
  246. /*=============================================================================
  247. * Default Callbacks
  248. *=============================================================================*/
  249. static void mqtt0_recv_callback(rt_uint8_t client_idx, const char *topic,
  250. const char *payload, rt_uint16_t len)
  251. {
  252. if (strstr(payload, "\"cmd\": \"sensor\", \"ext\":") == NULL)
  253. {
  254. LOG_I("mqtt0_recv_callback========================================");
  255. LOG_I("MQTT[%d] RX MSG:", client_idx);
  256. LOG_I(" Topic: %s", topic);
  257. LOG_I(" Len: %d", len);
  258. LOG_I(" Payload: %s", payload);
  259. LOG_I("mqtt0_recv_callback========================================");
  260. }
  261. // rt_kprintf("recv data start >>>");
  262. // for (int a = 0; a < len; a++)
  263. // {
  264. // rt_kprintf("%02X ",payload[a]);
  265. // if(a > 0 && a %32 == 0 )
  266. // {
  267. // rt_kprintf("\n");
  268. // }
  269. // }
  270. // rt_kprintf("<<<recv data finsh\n\n");
  271. ym310_protocol_process(client_idx, topic, payload, len);
  272. }
  273. static void mqtt1_recv_callback(rt_uint8_t client_idx, const char *topic,
  274. const char *payload, rt_uint16_t len)
  275. {
  276. LOG_I("mqtt1_recv_callback========================================");
  277. LOG_I("MQTT[%d] RX MSG:", client_idx);
  278. LOG_I(" Topic: %s", topic);
  279. LOG_I(" Len: %d", len);
  280. LOG_I(" Payload: %s", payload);
  281. LOG_I("mqtt1_recv_callback========================================");
  282. ym310_protocol_process(client_idx, topic, payload, len);
  283. }
  284. static void format_time(char *buf, size_t size, const struct rtc_time *tm)
  285. {
  286. rt_snprintf(buf, size, "%04d%02d%02d%02d%02d%02d",
  287. tm->year, tm->month, tm->date,
  288. tm->hour, tm->minute, tm->second);
  289. }
  290. uint8_t read_platform_flag(void)
  291. {
  292. uint8_t flag = 0;
  293. size_t len = ef_get_env_blob(EF_PLATFORM_FLAG, &flag, sizeof(flag), NULL);
  294. if (len == sizeof(flag) && flag == PLATFORM_VICE_MAGIC)
  295. {
  296. return 1; /* Backup platform */
  297. }
  298. return 0; /* Main platform (default) */
  299. }
  300. #include "ym310_protecl.h" // Provides SYS_INFO and sys_info declarations
  301. #include "main.h"
  302. extern struct rtc_time time_now;
  303. /*=============================================================================
  304. * MQTT Instance Configuration
  305. *=============================================================================*/
  306. static void ym310_mqtt_init_config(void)
  307. {
  308. rt_memset(g_mqtt_instances, 0, sizeof(g_mqtt_instances));
  309. rt_memset(g_mqtt_recv_bufs, 0, sizeof(g_mqtt_recv_bufs));
  310. uint8_t is_vice = read_platform_flag();
  311. LOG_I("Platform flag: %s", is_vice ? "VICE" : "MAIN");
  312. /* Get IMEI, ICCID, CSQ and update to sys_info (only save on first get or change) */
  313. char resp[128];
  314. at_resp_type_t at_ret;
  315. rt_bool_t need_save = RT_FALSE;
  316. /* Get IMEI (AT+CGSN) */
  317. at_ret = ym310_at_exec_cmd("AT+CGSN", resp, sizeof(resp), 5000);
  318. if (at_ret == AT_RESP_OK)
  319. {
  320. char *p = resp;
  321. while (*p && !isdigit(*p))
  322. p++; /* Skip non-digit prefix */
  323. int len = 0;
  324. while (p[len] && isdigit(p[len]))
  325. len++;
  326. if (len > 0 && len < (int)sizeof(sys_info.imei))
  327. {
  328. char new_imei[sizeof(sys_info.imei)] = {0};
  329. rt_memcpy(new_imei, p, len);
  330. if (strcmp(sys_info.imei, new_imei) != 0)
  331. {
  332. rt_memcpy(sys_info.imei, new_imei, len);
  333. sys_info.imei[len] = '\0';
  334. LOG_D("IMEI updated: %s", sys_info.imei);
  335. need_save = RT_TRUE;
  336. }
  337. else
  338. {
  339. LOG_D("IMEI unchanged: %s", sys_info.imei);
  340. }
  341. }
  342. }
  343. else
  344. {
  345. LOG_W("Failed to get IMEI");
  346. }
  347. /* Get ICCID (AT+ICCID) */
  348. at_ret = ym310_at_exec_cmd("AT+ICCID", resp, sizeof(resp), 5000);
  349. if (at_ret == AT_RESP_OK)
  350. {
  351. char *p = strstr(resp, "+ICCID:");
  352. if (p)
  353. {
  354. p += 7;
  355. while (*p == ' ' || *p == ':' || *p == '\r' || *p == '\n')
  356. p++;
  357. char *start = p;
  358. char *end = p;
  359. // Move to end of digit string (stop at space or newline)
  360. while (*end && *end != '\r' && *end != '\n' && *end != ' ')
  361. end++;
  362. int len = end - start;
  363. if (len > 0 && len < (int)sizeof(sys_info.iccid))
  364. {
  365. char new_iccid[sizeof(sys_info.iccid)] = {0};
  366. rt_memcpy(new_iccid, start, len);
  367. if (strcmp(sys_info.iccid, new_iccid) != 0)
  368. {
  369. rt_memcpy(sys_info.iccid, new_iccid, len);
  370. sys_info.iccid[len] = '\0';
  371. LOG_D("ICCID updated: %s", sys_info.iccid);
  372. save_sys_info();
  373. }
  374. else
  375. {
  376. LOG_D("ICCID unchanged: %s", sys_info.iccid);
  377. }
  378. }
  379. }
  380. }
  381. else
  382. {
  383. LOG_W("Failed to get ICCID");
  384. }
  385. /* Get CSQ */
  386. at_ret = ym310_at_exec_cmd("AT+CSQ", resp, sizeof(resp), 3000);
  387. if (at_ret == AT_RESP_OK)
  388. {
  389. int rssi, ber;
  390. if (ym310_parse_csq_response(resp, &rssi, &ber) == 0)
  391. {
  392. if (sys_info.csq != rssi)
  393. {
  394. sys_info.csq = rssi;
  395. LOG_D("CSQ updated: %d", sys_info.csq);
  396. // CSQ changes frequently, not saved on power up, call save_sys_info() if needed
  397. }
  398. else
  399. {
  400. LOG_D("CSQ unchanged: %d", sys_info.csq);
  401. }
  402. }
  403. else
  404. {
  405. LOG_W("Failed to parse CSQ response");
  406. }
  407. }
  408. else
  409. {
  410. LOG_W("Failed to get CSQ");
  411. }
  412. /* Save sys_info to Flash if changed */
  413. if (need_save)
  414. {
  415. save_sys_info();
  416. }
  417. /* Generate IMEI suffix */
  418. char imei_suffix[32];
  419. if (strlen(sys_info.imei) > 0)
  420. {
  421. rt_snprintf(imei_suffix, sizeof(imei_suffix), "/%s", sys_info.imei);
  422. }
  423. else
  424. {
  425. imei_suffix[0] = '\0';
  426. }
  427. /* Helper: Safe topic concatenation, avoid double slash */
  428. #define BUILD_TOPIC(dest, base, suffix) \
  429. do \
  430. { \
  431. int base_len = strlen(base); \
  432. const char *suf = (suffix); \
  433. if (base_len > 0 && base[base_len - 1] == '/') \
  434. { \
  435. if (suf[0] == '/') \
  436. suf++; \
  437. } \
  438. rt_snprintf((dest), MQTT_MAX_TOPIC_LEN, "%s%s", base, suf); \
  439. } while (0)
  440. /* Configure instance 0 (primary MQTT) */
  441. g_mqtt_instances[0].config.enabled = 1;
  442. g_mqtt_instances[0].config.client_idx = 0;
  443. char timestampe[32] = {0};
  444. format_time(timestampe, sizeof(timestampe), &time_now);
  445. char temp_topic[MQTT_MAX_TOPIC_LEN]; // Temp buffer
  446. if (!is_vice)
  447. {
  448. /* Main platform */
  449. rt_strncpy(g_mqtt_instances[0].config.host, MQTT0CFG_MAIN_HOST, MQTT_MAX_HOST_LEN - 1);
  450. g_mqtt_instances[0].config.port = MQTT0CFG_MAIN_PORT;
  451. /* Concatenate subscribe topic */
  452. BUILD_TOPIC(temp_topic, MQTT0CFG_MAIN_TSUB, imei_suffix);
  453. rt_strncpy(g_mqtt_instances[0].config.sub_topic, temp_topic, MQTT_MAX_TOPIC_LEN - 1);
  454. /* Concatenate publish topic */
  455. BUILD_TOPIC(temp_topic, MQTT0CFG_MAIN_TPUB, imei_suffix);
  456. rt_strncpy(g_mqtt_instances[0].config.pub_topic, temp_topic, MQTT_MAX_TOPIC_LEN - 1);
  457. rt_strncpy(g_mqtt_instances[0].config.client_id,
  458. (strlen(sys_info.imei) < 6) ? timestampe : sys_info.imei,
  459. MQTT_MAX_CLIENT_ID_LEN - 1);
  460. rt_strncpy(g_mqtt_instances[0].config.username, MQTT0CFG_MAIN_USR, MQTT_MAX_USERNAME_LEN - 1);
  461. rt_strncpy(g_mqtt_instances[0].config.password, MQTT0CFG_MAIN_PWD, MQTT_MAX_PASSWORD_LEN - 1);
  462. g_mqtt_instances[0].config.qos = MQTT0CFG_MAIN_QOS;
  463. }
  464. else
  465. {
  466. /* Backup platform */
  467. rt_strncpy(g_mqtt_instances[0].config.host, MQTT0CFG_VICE_HOST, MQTT_MAX_HOST_LEN - 1);
  468. g_mqtt_instances[0].config.port = MQTT0CFG_VICE_PORT;
  469. BUILD_TOPIC(temp_topic, MQTT0CFG_VICE_TSUB, imei_suffix);
  470. rt_strncpy(g_mqtt_instances[0].config.sub_topic, temp_topic, MQTT_MAX_TOPIC_LEN - 1);
  471. BUILD_TOPIC(temp_topic, MQTT0CFG_VICE_TPUB, imei_suffix);
  472. rt_strncpy(g_mqtt_instances[0].config.pub_topic, temp_topic, MQTT_MAX_TOPIC_LEN - 1);
  473. rt_strncpy(g_mqtt_instances[0].config.client_id,
  474. (strlen(sys_info.imei) < 6) ? timestampe : sys_info.imei,
  475. MQTT_MAX_CLIENT_ID_LEN - 1);
  476. rt_strncpy(g_mqtt_instances[0].config.username, MQTT0CFG_VICE_USR, MQTT_MAX_USERNAME_LEN - 1);
  477. rt_strncpy(g_mqtt_instances[0].config.password, MQTT0CFG_VICE_PWD, MQTT_MAX_PASSWORD_LEN - 1);
  478. g_mqtt_instances[0].config.qos = MQTT0CFG_VICE_QOS;
  479. }
  480. /* Instance 1 is customer config */
  481. SERVER_CFG *cfg = &sys_info.customer_cfg;
  482. g_mqtt_instances[1].config.enabled = (cfg->enable == VALID) ? 1 : 0;
  483. g_mqtt_instances[1].config.client_idx = 1;
  484. rt_strncpy(g_mqtt_instances[1].config.host, cfg->host, MQTT_MAX_HOST_LEN - 1);
  485. g_mqtt_instances[1].config.port = cfg->port;
  486. if (strstr(sys_info.customer_cfg.client_id, "USE_IMEI") != NULL)
  487. {
  488. rt_strncpy(g_mqtt_instances[1].config.client_id,
  489. (strlen(sys_info.imei) < 6) ? timestampe : sys_info.imei,
  490. MQTT_MAX_CLIENT_ID_LEN - 1);
  491. }
  492. else
  493. {
  494. rt_strncpy(g_mqtt_instances[1].config.client_id, cfg->client_id, MQTT_MAX_CLIENT_ID_LEN - 1);
  495. }
  496. rt_strncpy(g_mqtt_instances[1].config.username, cfg->username, MQTT_MAX_USERNAME_LEN - 1);
  497. rt_strncpy(g_mqtt_instances[1].config.password, cfg->password, MQTT_MAX_PASSWORD_LEN - 1);
  498. /* 拼接订阅主题 */
  499. BUILD_TOPIC(temp_topic, cfg->sub_topic, imei_suffix);
  500. rt_strncpy(g_mqtt_instances[1].config.sub_topic, temp_topic, MQTT_MAX_TOPIC_LEN - 1);
  501. /* 拼接发布主题 */
  502. BUILD_TOPIC(temp_topic, cfg->pub_topic, imei_suffix);
  503. rt_strncpy(g_mqtt_instances[1].config.pub_topic, temp_topic, MQTT_MAX_TOPIC_LEN - 1);
  504. g_mqtt_instances[1].config.qos = cfg->qos;
  505. g_mqtt_instances[1].state = MQTT_STATE_DISCONNECTED;
  506. g_mqtt_instances[1].recv_callback = RT_NULL;
  507. LOG_I("MQTT config initialized");
  508. }
  509. /*=============================================================================
  510. * Network Operations
  511. *=============================================================================*/
  512. int ym310_mqtt_open_network(rt_uint8_t client_idx)
  513. {
  514. at_resp_type_t resp;
  515. char cmd[256];
  516. char urc_buf[64];
  517. mqtt_instance_t *inst = &g_mqtt_instances[client_idx];
  518. if (!inst->config.enabled)
  519. {
  520. return -RT_ERROR;
  521. }
  522. LOG_I("MQTT[%d] Opening network to %s:%d...",
  523. client_idx, inst->config.host, inst->config.port);
  524. resp = ym310_at_exec_cmd_with_retry(
  525. "AT+QMTCFG=\"recv/mode\",0,0,1", RT_NULL, 0,
  526. YM310_AT_CMD_TIMEOUT_MS, 2);
  527. if (resp != AT_RESP_OK)
  528. {
  529. LOG_W("MQTT[%d] Failed to configure recv mode", client_idx);
  530. }
  531. rt_snprintf(cmd, sizeof(cmd), "AT+QMTOPEN=%d,\"%s\",%d",
  532. client_idx, inst->config.host, inst->config.port);
  533. resp = ym310_at_exec_async(cmd, "+QMTOPEN:", urc_buf, sizeof(urc_buf), 30000);
  534. if (resp != AT_RESP_OK)
  535. {
  536. LOG_E("MQTT[%d] Failed to open network", client_idx);
  537. return -RT_ERROR;
  538. }
  539. int result_code = ym310_str_extract_nth_int(urc_buf, 1);
  540. if (result_code != 0 && result_code != 1)
  541. {
  542. LOG_E("MQTT[%d] Network open failed: result=%d", client_idx, result_code);
  543. return -RT_ERROR;
  544. }
  545. LOG_I("MQTT[%d] Network opened (result=%d)", client_idx, result_code);
  546. rt_thread_mdelay(6000);
  547. return RT_EOK;
  548. }
  549. int ym310_mqtt_close_network(rt_uint8_t client_idx)
  550. {
  551. at_resp_type_t resp;
  552. char cmd[64];
  553. LOG_I("MQTT[%d] Closing network...", client_idx);
  554. rt_snprintf(cmd, sizeof(cmd), "AT+QMTDISC=%d", client_idx);
  555. resp = ym310_at_exec_async(cmd, "+QMTDISC:", RT_NULL, 0, 10000);
  556. if (resp != AT_RESP_OK)
  557. {
  558. LOG_W("MQTT[%d] Network close timeout", client_idx);
  559. return -RT_ERROR;
  560. }
  561. LOG_I("MQTT[%d] Network closed", client_idx);
  562. return RT_EOK;
  563. }
  564. int ym310_mqtt_connect(rt_uint8_t client_idx)
  565. {
  566. at_resp_type_t resp;
  567. char cmd[512];
  568. char urc_buf[64];
  569. mqtt_instance_t *inst = &g_mqtt_instances[client_idx];
  570. if (!inst->config.enabled)
  571. {
  572. return -RT_ERROR;
  573. }
  574. LOG_I("MQTT[%d] Connecting to broker...", client_idx);
  575. inst->state = MQTT_STATE_CONNECTING;
  576. if (inst->config.username[0] != '\0' && inst->config.password[0] != '\0')
  577. {
  578. rt_snprintf(cmd, sizeof(cmd), "AT+QMTCONN=%d,\"%s\",\"%s\",\"%s\"",
  579. client_idx, inst->config.client_id,
  580. inst->config.username, inst->config.password);
  581. }
  582. else
  583. {
  584. rt_snprintf(cmd, sizeof(cmd), "AT+QMTCONN=%d,\"%s\"",
  585. client_idx, inst->config.client_id);
  586. }
  587. resp = ym310_at_exec_async(cmd, "+QMTCONN:", urc_buf, sizeof(urc_buf), 30000);
  588. if (resp != AT_RESP_OK)
  589. {
  590. LOG_E("MQTT[%d] Connect failed", client_idx);
  591. inst->state = MQTT_STATE_DISCONNECTED;
  592. return -RT_ERROR;
  593. }
  594. int result = ym310_str_extract_nth_int(urc_buf, 1);
  595. int ret_code = ym310_str_extract_nth_int(urc_buf, 2);
  596. if (result != 0)
  597. {
  598. LOG_E("MQTT[%d] Connection failed: result=%d, ret=%d", client_idx, result, ret_code);
  599. inst->state = MQTT_STATE_DISCONNECTED;
  600. return -RT_ERROR;
  601. }
  602. inst->state = MQTT_STATE_CONNECTED;
  603. inst->last_activity = rt_tick_get();
  604. LOG_I("MQTT[%d] Connected successfully", client_idx);
  605. return RT_EOK;
  606. }
  607. int ym310_mqtt_disconnect(rt_uint8_t client_idx)
  608. {
  609. mqtt_instance_t *inst = &g_mqtt_instances[client_idx];
  610. if (inst->state == MQTT_STATE_DISCONNECTED)
  611. {
  612. return RT_EOK;
  613. }
  614. inst->state = MQTT_STATE_DISCONNECTING;
  615. ym310_mqtt_close_network(client_idx);
  616. inst->state = MQTT_STATE_DISCONNECTED;
  617. LOG_I("MQTT[%d] Disconnected", client_idx);
  618. return RT_EOK;
  619. }
  620. int ym310_mqtt_subscribe(rt_uint8_t client_idx, const char *topic, rt_uint8_t qos)
  621. {
  622. at_resp_type_t resp;
  623. char cmd[256];
  624. char urc_buf[64];
  625. mqtt_instance_t *inst = &g_mqtt_instances[client_idx];
  626. rt_uint8_t msg_id = s_msg_id++;
  627. if (!inst->config.enabled || inst->state != MQTT_STATE_CONNECTED)
  628. {
  629. return -RT_ERROR;
  630. }
  631. // LOG_I("MQTT[%d] Subscribing to %s (qos=%d)...", client_idx, topic, qos);
  632. rt_snprintf(cmd, sizeof(cmd), "AT+QMTSUB=%d,%d,\"%s\",%d",
  633. client_idx, msg_id, topic, qos);
  634. resp = ym310_at_exec_async(cmd, "+QMTSUB:", urc_buf, sizeof(urc_buf), 30000);
  635. if (resp != AT_RESP_OK)
  636. {
  637. LOG_E("MQTT[%d] Subscribe failed", client_idx);
  638. return -RT_ERROR;
  639. }
  640. int result = -1, granted_qos = -1;
  641. sscanf(urc_buf, "+QMTSUB: %*d,%*d,%d,%d", &result, &granted_qos);
  642. if (result != 0)
  643. {
  644. LOG_E("MQTT[%d] Subscribe result=%d", client_idx, result);
  645. return -RT_ERROR;
  646. }
  647. rt_kprintf("\n MQTT[%d] Subscribed to %s (qos=%d) SUCCESSFULLY\n\n", client_idx, topic, granted_qos);
  648. return RT_EOK;
  649. }
  650. int ym310_mqtt_unsubscribe(rt_uint8_t client_idx, const char *topic)
  651. {
  652. at_resp_type_t resp;
  653. char cmd[256];
  654. mqtt_instance_t *inst = &g_mqtt_instances[client_idx];
  655. rt_uint8_t msg_id = s_msg_id++;
  656. if (!inst->config.enabled || inst->state != MQTT_STATE_CONNECTED)
  657. {
  658. return -RT_ERROR;
  659. }
  660. rt_snprintf(cmd, sizeof(cmd), "AT+QMTUNS=%d,%d,\"%s\"", client_idx, msg_id, topic);
  661. resp = ym310_at_exec_cmd_with_retry(cmd, RT_NULL, 0, 30000, 2);
  662. if (resp != AT_RESP_OK)
  663. {
  664. LOG_W("MQTT[%d] Unsubscribe failed", client_idx);
  665. return -RT_ERROR;
  666. }
  667. LOG_I("MQTT[%d] Unsubscribed from %s", client_idx, topic);
  668. return RT_EOK;
  669. }
  670. static void ym310_mqtt_thread_entry(void *param)
  671. {
  672. int ret;
  673. rt_tick_t last_check[MQTT_INSTANCE_MAX] = {0};
  674. LOG_D("========================================");
  675. LOG_D("MQTT THREAD STARTED (Simplified v1.0)");
  676. LOG_D("========================================");
  677. /* Init MQTT config (keep as-is) */
  678. ym310_mqtt_init_config();
  679. /* Register default callbacks (example print) */
  680. g_mqtt_callbacks[0] = mqtt0_recv_callback; // Defined at file top
  681. g_mqtt_callbacks[1] = mqtt1_recv_callback;
  682. /* Register MQTT URC handlers (keep as-is) */
  683. ym310_at_register_urc("+QMTRECV:", urc_qmtrecv_handler);
  684. ym310_at_register_urc("+QMTSTAT:", urc_qmtstat_handler);
  685. ym310_at_register_urc("+QMTCONN:", urc_qmtconn_handler);
  686. ym310_at_register_urc("+QMTOPEN:", urc_qmtopen_handler);
  687. s_mqtt_thread_running = RT_TRUE;
  688. while (s_mqtt_thread_running)
  689. {
  690. /****************** 1. Process recv data (simplified) ******************/
  691. for (int i = 0; i < MQTT_INSTANCE_MAX; i++)
  692. {
  693. rt_bool_t has_data = RT_FALSE;
  694. char topic[MQTT_MAX_TOPIC_LEN] = {0};
  695. char payload[MQTT_MAX_PAYLOAD_SIZE] = {0};
  696. rt_uint16_t payload_len = 0;
  697. /* Disable interrupts to safely read from recv buffer and clear flag */
  698. // rt_base_t //level = rt_hw_interrupt_disable();
  699. if (g_mqtt_recv_bufs[i].data_ready)
  700. {
  701. /* Copy to local variable */
  702. rt_strncpy(topic, g_mqtt_recv_bufs[i].topic, MQTT_MAX_TOPIC_LEN - 1);
  703. rt_strncpy(payload, g_mqtt_recv_bufs[i].payload, MQTT_MAX_PAYLOAD_SIZE - 1);
  704. payload_len = g_mqtt_recv_bufs[i].payload_len;
  705. /* Clear buffer (only clear flag, no need to clear data) */
  706. g_mqtt_recv_bufs[i].data_ready = RT_FALSE;
  707. g_mqtt_recv_bufs[i].payload_len = 0;
  708. g_mqtt_recv_bufs[i].topic[0] = '\0';
  709. g_mqtt_recv_bufs[i].payload[0] = '\0';
  710. has_data = RT_TRUE;
  711. }
  712. // rt_hw_interrupt_enable(level);
  713. if (has_data)
  714. {
  715. if (ym310_at_mutex_lock(1000) == RT_EOK)
  716. {
  717. if (g_mqtt_callbacks[i] != RT_NULL)
  718. {
  719. g_mqtt_callbacks[i](i, topic, payload, payload_len);
  720. }
  721. else
  722. {
  723. LOG_D("MQTT[%d] RX: topic=%s, payload=%.*s",
  724. i, topic, payload_len, payload);
  725. }
  726. ym310_at_mutex_unlock();
  727. }
  728. else
  729. {
  730. LOG_W("MQTT[%d] Failed to lock AT mutex for processing", i);
  731. }
  732. }
  733. }
  734. /****************** 2. Connection management (keep as-is) ******************/
  735. for (int i = 0; i < MQTT_INSTANCE_MAX; i++)
  736. {
  737. mqtt_instance_t *inst = &g_mqtt_instances[i];
  738. if (!inst->config.enabled)
  739. {
  740. continue;
  741. }
  742. // rt_base_t //level = rt_hw_interrupt_disable();
  743. mqtt_state_t current_state = inst->state;
  744. // rt_hw_interrupt_enable(level);
  745. switch (current_state)
  746. {
  747. case MQTT_STATE_DISCONNECTED:
  748. {
  749. rt_tick_t current_tick = rt_tick_get();
  750. rt_tick_t retry_interval;
  751. rt_uint32_t backoff = 1;
  752. for (int j = 0; j < inst->err_count && j < 4; j++)
  753. {
  754. backoff *= 2;
  755. }
  756. retry_interval = CONNECT_RETRY_BASE_MS * backoff;
  757. if (retry_interval > CONNECT_RETRY_MAX_MS)
  758. {
  759. retry_interval = CONNECT_RETRY_MAX_MS;
  760. }
  761. if ((current_tick - s_last_connect_attempt[i]) > rt_tick_from_millisecond(retry_interval))
  762. {
  763. s_last_connect_attempt[i] = current_tick;
  764. // LOG_I("MQTT[%d] Connecting (err=%u, interval=%ums)...",
  765. // i, inst->err_count, retry_interval);
  766. ret = ym310_mqtt_open_network(i);
  767. if (ret == RT_EOK)
  768. {
  769. ret = ym310_mqtt_connect(i);
  770. if (ret == RT_EOK)
  771. {
  772. ym310_mqtt_subscribe(i, inst->config.sub_topic, inst->config.qos);
  773. }
  774. else
  775. {
  776. ym310_mqtt_close_network(i);
  777. }
  778. }
  779. if (ret != RT_EOK)
  780. {
  781. LOG_W("MQTT[%d] Connection failed, will retry", i);
  782. }
  783. }
  784. break;
  785. }
  786. case MQTT_STATE_CONNECTED:
  787. if ((rt_tick_get() - last_check[i]) > rt_tick_from_millisecond(60000))
  788. {
  789. last_check[i] = rt_tick_get();
  790. LOG_D("MQTT[%d] TX:%u RX:%u ERR:%u",
  791. i, inst->tx_count, inst->rx_count, inst->err_count);
  792. }
  793. break;
  794. case MQTT_STATE_RECONNECTING:
  795. LOG_I("MQTT[%d] Reconnecting...", i);
  796. ym310_mqtt_disconnect(i);
  797. // level = rt_hw_interrupt_disable();
  798. inst->state = MQTT_STATE_DISCONNECTED;
  799. // rt_hw_interrupt_enable(level);
  800. break;
  801. default:
  802. break;
  803. }
  804. }
  805. rt_thread_mdelay(50); // Main loop interval
  806. }
  807. /****************** 3. Cleanup (keep as-is) ******************/
  808. for (int i = 0; i < MQTT_INSTANCE_MAX; i++)
  809. {
  810. if (g_mqtt_instances[i].config.enabled)
  811. {
  812. ym310_mqtt_disconnect(i);
  813. }
  814. }
  815. ym310_at_unregister_urc("+QMTRECV:");
  816. ym310_at_unregister_urc("+QMTSTAT:");
  817. ym310_at_unregister_urc("+QMTCONN:");
  818. ym310_at_unregister_urc("+QMTOPEN:");
  819. LOG_I("MQTT thread stopped");
  820. }
  821. /*=============================================================================
  822. * Public API
  823. *=============================================================================*/
  824. int ym310_mqtt_init(void)
  825. {
  826. if (g_mqtt_mutex == RT_NULL)
  827. {
  828. g_mqtt_mutex = rt_mutex_create("mqtt_mtx", RT_IPC_FLAG_FIFO);
  829. if (g_mqtt_mutex == RT_NULL)
  830. {
  831. LOG_E("Failed to create MQTT mutex!");
  832. return -RT_ERROR;
  833. }
  834. // LOG_D("g_mqtt_mutex created at 0x%p", g_mqtt_mutex);
  835. }
  836. // LOG_I("MQTT framework initialized");
  837. return RT_EOK;
  838. }
  839. void ym310_mqtt_deinit(void)
  840. {
  841. /* Cleanup handled by thread */
  842. }
  843. rt_thread_t ym310_start_mqtt_thread(void)
  844. {
  845. if (ym310_mqtt_init() != RT_EOK)
  846. {
  847. return RT_NULL;
  848. }
  849. g_mqtt_thread = rt_thread_create("ym310_mqtt", ym310_mqtt_thread_entry, RT_NULL,
  850. MQTT_THREAD_STACK_SIZE, MQTT_THREAD_PRIORITY, 10);
  851. if (g_mqtt_thread == RT_NULL)
  852. {
  853. LOG_E("Failed to create MQTT thread!");
  854. return RT_NULL;
  855. }
  856. rt_thread_startup(g_mqtt_thread);
  857. // LOG_I("MQTT thread started");
  858. return g_mqtt_thread;
  859. }
  860. mqtt_state_t ym310_mqtt_get_state(rt_uint8_t client_idx)
  861. {
  862. if (client_idx >= MQTT_INSTANCE_MAX)
  863. {
  864. return MQTT_STATE_DISCONNECTED;
  865. }
  866. return g_mqtt_instances[client_idx].state;
  867. }
  868. rt_bool_t ym310_mqtt_is_connected(rt_uint8_t client_idx)
  869. {
  870. if (client_idx >= MQTT_INSTANCE_MAX)
  871. {
  872. return RT_FALSE;
  873. }
  874. return (g_mqtt_instances[client_idx].state == MQTT_STATE_CONNECTED);
  875. }
  876. int ym310_mqtt_set_recv_callback(rt_uint8_t client_idx,
  877. void (*callback)(rt_uint8_t, const char *,
  878. const char *, rt_uint16_t))
  879. {
  880. if (client_idx >= MQTT_INSTANCE_MAX)
  881. {
  882. return -RT_ERROR;
  883. }
  884. // rt_base_t //level = rt_hw_interrupt_disable();
  885. g_mqtt_callbacks[client_idx] = (mqtt_recv_handler_t)callback;
  886. // rt_hw_interrupt_enable(level);
  887. LOG_I("MQTT[%d] Callback registered", client_idx);
  888. return RT_EOK;
  889. }
  890. void ym310_mqtt_force_reconnect(rt_uint8_t client_idx)
  891. {
  892. if (client_idx >= MQTT_INSTANCE_MAX)
  893. {
  894. return;
  895. }
  896. mqtt_instance_t *inst = &g_mqtt_instances[client_idx];
  897. LOG_I("MQTT[%d] Forcing reconnection...", client_idx);
  898. ym310_mqtt_disconnect(client_idx);
  899. // rt_base_t //level = rt_hw_interrupt_disable();
  900. inst->state = MQTT_STATE_DISCONNECTED;
  901. inst->err_count = 0;
  902. // rt_hw_interrupt_enable(level);
  903. s_last_connect_attempt[client_idx] = 0;
  904. LOG_I("MQTT[%d] Reconnection triggered", client_idx);
  905. }
  906. void ym310_mqtt_force_reconnect_all(void)
  907. {
  908. LOG_I("Forcing reconnection for all MQTT instances...");
  909. for (int i = 0; i < MQTT_INSTANCE_MAX; i++)
  910. {
  911. if (g_mqtt_instances[i].config.enabled)
  912. {
  913. ym310_mqtt_force_reconnect(i);
  914. }
  915. }
  916. }