| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067 |
- /*
- * YM310 MQTT Dual-Instance Management - FIXED v7.1.0
- *
- * CRITICAL FIXES:
- * - Fixed buffer overflow in urc_qmtrecv_handler
- * - Added strict bounds checking on all string operations
- * - Added memory alignment for buffer structures
- * - Added extensive debug logging
- * - Fixed race condition in data_ready flag handling
- */
- #define DBG_TAG "Y-MQ "
- #define DBG_LVL DBG_INFO
- #include <rtdbg.h>
- #include "ym310_mqtt.h"
- #include <string.h>
- #include <stdio.h>
- #include <rtthread.h>
- #include <rthw.h> /* Provides rt_hw_interrupt_disable/enable */
- #include <ctype.h> /* Provides isdigit */
- #include <easyflash.h> /* Provides ef_get_env_blob */
- #include "ym310_mqtt.h"
- #include "ym310_protecl.h" /* Provides SYS_INFO and save_sys_info declarations */
- /*=============================================================================
- * Global Variables
- *=============================================================================*/
- mqtt_instance_t g_mqtt_instances[MQTT_INSTANCE_MAX];
- rt_mutex_t g_mqtt_mutex = RT_NULL;
- /* Receive buffers - populated by URC handler */
- mqtt_recv_buf_t g_mqtt_recv_bufs[MQTT_INSTANCE_MAX];
- /* Callback handlers for each instance */
- static mqtt_recv_handler_t g_mqtt_callbacks[MQTT_INSTANCE_MAX] = {RT_NULL, RT_NULL};
- rt_thread_t g_mqtt_thread = RT_NULL;
- static volatile rt_bool_t s_mqtt_thread_running = RT_FALSE;
- static rt_uint8_t s_msg_id = 1;
- volatile rt_tick_t g_last_mqtt_success_tick = 0;
- /* Connection retry timing */
- static rt_tick_t s_last_connect_attempt[MQTT_INSTANCE_MAX] = {0};
- uint8_t read_platform_flag(void);
- #define CONNECT_RETRY_BASE_MS 5000
- #define CONNECT_RETRY_MAX_MS 60000
- /*=============================================================================
- * DEBUG: IPC Object Integrity Check
- *=============================================================================*/
- extern rt_mutex_t s_at_mutex; // Added at file top
- // static void check_ipc_integrity(const char *location)
- //{
- // extern rt_sem_t s_cmd_sem;
- // extern rt_sem_t s_rx_sem;
- //
- // // LOG_D("IPC_CHECK [%s]: Checking IPC objects...", location);
- //
- // if (g_mqtt_mutex != RT_NULL)
- // {
- // rt_object_t obj = (rt_object_t)g_mqtt_mutex;
- // rt_uint8_t type = obj->type & ~RT_Object_Class_Static;
- // if (type != RT_Object_Class_Mutex)
- // {
- // LOG_E("IPC_CHECK [%s]: g_mqtt_mutex CORRUPTED! type=%d", location, type);
- // }
- // }
- //
- // if (s_at_mutex != RT_NULL)
- // {
- // rt_object_t obj = (rt_object_t)s_at_mutex;
- // rt_uint8_t type = obj->type & ~RT_Object_Class_Static;
- // if (type != RT_Object_Class_Mutex)
- // {
- // LOG_E("IPC_CHECK [%s]: s_at_mutex CORRUPTED! type=%d", location, type);
- // }
- // }
- //
- // if (s_cmd_sem != RT_NULL)
- // {
- // rt_object_t obj = (rt_object_t)s_cmd_sem;
- // rt_uint8_t type = obj->type & ~RT_Object_Class_Static;
- // if (type != RT_Object_Class_Semaphore)
- // {
- // LOG_E("IPC_CHECK [%s]: s_cmd_sem CORRUPTED! type=%d", location, type);
- // }
- // }
- //
- // if (s_rx_sem != RT_NULL)
- // {
- // rt_object_t obj = (rt_object_t)s_rx_sem;
- // rt_uint8_t type = obj->type & ~RT_Object_Class_Static;
- // if (type != RT_Object_Class_Semaphore)
- // {
- // LOG_E("IPC_CHECK [%s]: s_rx_sem CORRUPTED! type=%d", location, type);
- // }
- // }
- // }
- static void urc_qmtrecv_handler(const char *urc)
- {
- int client_idx = -1, msgid = 0;
- char topic[MQTT_MAX_TOPIC_LEN] = {0};
- int topic_len = 0;
- const char *payload_start = NULL;
- int payload_len = 0;
- int declared_len = 0;
- if (urc == RT_NULL)
- return;
- /* Find +QMTRECV: prefix */
- const char *p = rt_strstr(urc, "+QMTRECV:");
- if (p == RT_NULL)
- return;
- p += 9;
- /* Skip whitespace */
- while (*p == ' ' || *p == '\t')
- p++;
- /* Parse client_idx and msgid */
- if (sscanf(p, "%d,%d", &client_idx, &msgid) != 2)
- return;
- if (client_idx < 0 || client_idx >= MQTT_INSTANCE_MAX)
- return;
- /* Locate topic start quote */
- const char *topic_start = rt_strstr(p, "\"");
- if (topic_start == RT_NULL)
- return;
- topic_start++;
- const char *topic_end = rt_strstr(topic_start, "\"");
- if (topic_end == RT_NULL || topic_end == topic_start)
- return;
- topic_len = topic_end - topic_start;
- if (topic_len >= MQTT_MAX_TOPIC_LEN)
- topic_len = MQTT_MAX_TOPIC_LEN - 1;
- rt_memcpy(topic, topic_start, topic_len);
- topic[topic_len] = '\0';
- /* Locate payload */
- const char *after_topic = topic_end + 1;
- while (*after_topic == ' ' || *after_topic == '\t')
- after_topic++;
- if (*after_topic != ',')
- return;
- after_topic++;
- while (*after_topic == ' ' || *after_topic == '\t')
- after_topic++;
- if (*after_topic == '"')
- {
- /* No length field format: payload directly after comma */
- payload_start = after_topic; // Point to start quote
- declared_len = 0; // No length, use last quote to locate
- }
- else if (isdigit(*after_topic))
- {
- /* Has length field format */
- declared_len = atoi(after_topic);
- while (*after_topic && *after_topic != ',')
- after_topic++;
- if (*after_topic != ',')
- return;
- after_topic++;
- while (*after_topic == ' ' || *after_topic == '\t')
- after_topic++;
- if (*after_topic != '"')
- return;
- payload_start = after_topic;
- }
- else
- {
- return;
- }
- /* Determine actual payload length */
- payload_start++; // Skip start quote, point to first char
- if (declared_len > 0)
- {
- /* Has length field: use length field value, not rely on quotes */
- payload_len = declared_len;
- if (payload_len >= MQTT_MAX_PAYLOAD_SIZE)
- {
- LOG_W("URC: Payload too large (declared=%d >= %d), truncating",
- payload_len, MQTT_MAX_PAYLOAD_SIZE);
- payload_len = MQTT_MAX_PAYLOAD_SIZE - 1;
- }
- }
- else
- {
- /* No length field: find last quote as end */
- const char *payload_end = strrchr(payload_start, '"');
- if (payload_end == RT_NULL)
- {
- LOG_E("URC: No closing quote found for payload");
- return;
- }
- payload_len = payload_end - payload_start;
- if (payload_len >= MQTT_MAX_PAYLOAD_SIZE)
- {
- LOG_W("URC: Payload too large (%d >= %d), truncating",
- payload_len, MQTT_MAX_PAYLOAD_SIZE);
- payload_len = MQTT_MAX_PAYLOAD_SIZE - 1;
- }
- }
- LOG_D("URC: MQTT[%d] RX: topic=%s, len=%d", client_idx, topic, payload_len);
- /* ========== Disable interrupts to write recv buffer ========== */
- // rt_base_t //level = rt_hw_interrupt_disable();
- if (g_mqtt_recv_bufs[client_idx].data_ready)
- {
- // rt_hw_interrupt_enable(level);
- LOG_W("MQTT[%d] RX buffer busy, dropping msg", client_idx);
- return;
- }
- /* Copy topic */
- rt_strncpy(g_mqtt_recv_bufs[client_idx].topic, topic, MQTT_MAX_TOPIC_LEN - 1);
- g_mqtt_recv_bufs[client_idx].topic[MQTT_MAX_TOPIC_LEN - 1] = '\0';
- /* Copy payload */
- if (payload_len > 0)
- {
- rt_memcpy(g_mqtt_recv_bufs[client_idx].payload, payload_start, payload_len);
- }
- g_mqtt_recv_bufs[client_idx].payload[payload_len] = '\0';
- g_mqtt_recv_bufs[client_idx].payload_len = payload_len;
- g_mqtt_recv_bufs[client_idx].data_ready = RT_TRUE;
- /* Update statistics */
- g_mqtt_instances[client_idx].rx_count++;
- g_mqtt_instances[client_idx].last_activity = rt_tick_get();
- // rt_hw_interrupt_enable(level);
- LOG_D("URC: Data stored in recv buffer[%d]", client_idx);
- }
- static void urc_qmtstat_handler(const char *urc)
- {
- int client_idx, err_code;
- sscanf(urc, "+QMTSTAT: %d,%d", &client_idx, &err_code);
- if (client_idx >= 0 && client_idx < MQTT_INSTANCE_MAX)
- {
- LOG_W("MQTT[%d] status: err=%d", client_idx, err_code);
- if (err_code != 0)
- {
- // rt_base_t //level = rt_hw_interrupt_disable();
- g_mqtt_instances[client_idx].state = MQTT_STATE_DISCONNECTED;
- g_mqtt_instances[client_idx].err_count++;
- // rt_hw_interrupt_enable(level);
- }
- }
- }
- static void urc_qmtconn_handler(const char *urc)
- {
- int client_idx, result, ret_code;
- sscanf(urc, "+QMTCONN: %d,%d,%d", &client_idx, &result, &ret_code);
- LOG_D("MQTT[%d] conn: result=%d, ret=%d", client_idx, result, ret_code);
- }
- static void urc_qmtopen_handler(const char *urc)
- {
- int client_idx, result;
- sscanf(urc, "+QMTOPEN: %d,%d", &client_idx, &result);
- LOG_D("MQTT[%d] open: result=%d", client_idx, result);
- }
- /*=============================================================================
- * Default Callbacks
- *=============================================================================*/
- static void mqtt0_recv_callback(rt_uint8_t client_idx, const char *topic,
- const char *payload, rt_uint16_t len)
- {
- if (strstr(payload, "\"cmd\": \"sensor\", \"ext\":") == NULL)
- {
- LOG_I("mqtt0_recv_callback========================================");
- LOG_I("MQTT[%d] RX MSG:", client_idx);
- LOG_I(" Topic: %s", topic);
- LOG_I(" Len: %d", len);
- LOG_I(" Payload: %s", payload);
- LOG_I("mqtt0_recv_callback========================================");
- }
- // rt_kprintf("recv data start >>>");
- // for (int a = 0; a < len; a++)
- // {
- // rt_kprintf("%02X ",payload[a]);
- // if(a > 0 && a %32 == 0 )
- // {
- // rt_kprintf("\n");
- // }
- // }
- // rt_kprintf("<<<recv data finsh\n\n");
- ym310_protocol_process(client_idx, topic, payload, len);
- }
- static void mqtt1_recv_callback(rt_uint8_t client_idx, const char *topic,
- const char *payload, rt_uint16_t len)
- {
- LOG_I("mqtt1_recv_callback========================================");
- LOG_I("MQTT[%d] RX MSG:", client_idx);
- LOG_I(" Topic: %s", topic);
- LOG_I(" Len: %d", len);
- LOG_I(" Payload: %s", payload);
- LOG_I("mqtt1_recv_callback========================================");
- ym310_protocol_process(client_idx, topic, payload, len);
- }
- static void format_time(char *buf, size_t size, const struct rtc_time *tm)
- {
- rt_snprintf(buf, size, "%04d%02d%02d%02d%02d%02d",
- tm->year, tm->month, tm->date,
- tm->hour, tm->minute, tm->second);
- }
- uint8_t read_platform_flag(void)
- {
- uint8_t flag = 0;
- size_t len = ef_get_env_blob(EF_PLATFORM_FLAG, &flag, sizeof(flag), NULL);
- if (len == sizeof(flag) && flag == PLATFORM_VICE_MAGIC)
- {
- return 1; /* Backup platform */
- }
- return 0; /* Main platform (default) */
- }
- #include "ym310_protecl.h" // Provides SYS_INFO and sys_info declarations
- #include "main.h"
- extern struct rtc_time time_now;
- /*=============================================================================
- * MQTT Instance Configuration
- *=============================================================================*/
- static void ym310_mqtt_init_config(void)
- {
- rt_memset(g_mqtt_instances, 0, sizeof(g_mqtt_instances));
- rt_memset(g_mqtt_recv_bufs, 0, sizeof(g_mqtt_recv_bufs));
- uint8_t is_vice = read_platform_flag();
- LOG_I("Platform flag: %s", is_vice ? "VICE" : "MAIN");
- /* Get IMEI, ICCID, CSQ and update to sys_info (only save on first get or change) */
- char resp[128];
- at_resp_type_t at_ret;
- rt_bool_t need_save = RT_FALSE;
- /* Get IMEI (AT+CGSN) */
- at_ret = ym310_at_exec_cmd("AT+CGSN", resp, sizeof(resp), 5000);
- if (at_ret == AT_RESP_OK)
- {
- char *p = resp;
- while (*p && !isdigit(*p))
- p++; /* Skip non-digit prefix */
- int len = 0;
- while (p[len] && isdigit(p[len]))
- len++;
- if (len > 0 && len < (int)sizeof(sys_info.imei))
- {
- char new_imei[sizeof(sys_info.imei)] = {0};
- rt_memcpy(new_imei, p, len);
- if (strcmp(sys_info.imei, new_imei) != 0)
- {
- rt_memcpy(sys_info.imei, new_imei, len);
- sys_info.imei[len] = '\0';
- LOG_D("IMEI updated: %s", sys_info.imei);
- need_save = RT_TRUE;
- }
- else
- {
- LOG_D("IMEI unchanged: %s", sys_info.imei);
- }
- }
- }
- else
- {
- LOG_W("Failed to get IMEI");
- }
- /* Get ICCID (AT+ICCID) */
- at_ret = ym310_at_exec_cmd("AT+ICCID", resp, sizeof(resp), 5000);
- if (at_ret == AT_RESP_OK)
- {
- char *p = strstr(resp, "+ICCID:");
- if (p)
- {
- p += 7;
- while (*p == ' ' || *p == ':' || *p == '\r' || *p == '\n')
- p++;
- char *start = p;
- char *end = p;
- // Move to end of digit string (stop at space or newline)
- while (*end && *end != '\r' && *end != '\n' && *end != ' ')
- end++;
- int len = end - start;
- if (len > 0 && len < (int)sizeof(sys_info.iccid))
- {
- char new_iccid[sizeof(sys_info.iccid)] = {0};
- rt_memcpy(new_iccid, start, len);
- if (strcmp(sys_info.iccid, new_iccid) != 0)
- {
- rt_memcpy(sys_info.iccid, new_iccid, len);
- sys_info.iccid[len] = '\0';
- LOG_D("ICCID updated: %s", sys_info.iccid);
- save_sys_info();
- }
- else
- {
- LOG_D("ICCID unchanged: %s", sys_info.iccid);
- }
- }
- }
- }
- else
- {
- LOG_W("Failed to get ICCID");
- }
- /* Get CSQ */
- at_ret = ym310_at_exec_cmd("AT+CSQ", resp, sizeof(resp), 3000);
- if (at_ret == AT_RESP_OK)
- {
- int rssi, ber;
- if (ym310_parse_csq_response(resp, &rssi, &ber) == 0)
- {
- if (sys_info.csq != rssi)
- {
- sys_info.csq = rssi;
- LOG_D("CSQ updated: %d", sys_info.csq);
- // CSQ changes frequently, not saved on power up, call save_sys_info() if needed
- }
- else
- {
- LOG_D("CSQ unchanged: %d", sys_info.csq);
- }
- }
- else
- {
- LOG_W("Failed to parse CSQ response");
- }
- }
- else
- {
- LOG_W("Failed to get CSQ");
- }
- /* Save sys_info to Flash if changed */
- if (need_save)
- {
- save_sys_info();
- }
- /* Generate IMEI suffix */
- char imei_suffix[32];
- if (strlen(sys_info.imei) > 0)
- {
- rt_snprintf(imei_suffix, sizeof(imei_suffix), "/%s", sys_info.imei);
- }
- else
- {
- imei_suffix[0] = '\0';
- }
- /* Helper: Safe topic concatenation, avoid double slash */
- #define BUILD_TOPIC(dest, base, suffix) \
- do \
- { \
- int base_len = strlen(base); \
- const char *suf = (suffix); \
- if (base_len > 0 && base[base_len - 1] == '/') \
- { \
- if (suf[0] == '/') \
- suf++; \
- } \
- rt_snprintf((dest), MQTT_MAX_TOPIC_LEN, "%s%s", base, suf); \
- } while (0)
- /* Configure instance 0 (primary MQTT) */
- g_mqtt_instances[0].config.enabled = 1;
- g_mqtt_instances[0].config.client_idx = 0;
- char timestampe[32] = {0};
- format_time(timestampe, sizeof(timestampe), &time_now);
- char temp_topic[MQTT_MAX_TOPIC_LEN]; // Temp buffer
- if (!is_vice)
- {
- /* Main platform */
- rt_strncpy(g_mqtt_instances[0].config.host, MQTT0CFG_MAIN_HOST, MQTT_MAX_HOST_LEN - 1);
- g_mqtt_instances[0].config.port = MQTT0CFG_MAIN_PORT;
- /* Concatenate subscribe topic */
- BUILD_TOPIC(temp_topic, MQTT0CFG_MAIN_TSUB, imei_suffix);
- rt_strncpy(g_mqtt_instances[0].config.sub_topic, temp_topic, MQTT_MAX_TOPIC_LEN - 1);
- /* Concatenate publish topic */
- BUILD_TOPIC(temp_topic, MQTT0CFG_MAIN_TPUB, imei_suffix);
- rt_strncpy(g_mqtt_instances[0].config.pub_topic, temp_topic, MQTT_MAX_TOPIC_LEN - 1);
- rt_strncpy(g_mqtt_instances[0].config.client_id,
- (strlen(sys_info.imei) < 6) ? timestampe : sys_info.imei,
- MQTT_MAX_CLIENT_ID_LEN - 1);
- rt_strncpy(g_mqtt_instances[0].config.username, MQTT0CFG_MAIN_USR, MQTT_MAX_USERNAME_LEN - 1);
- rt_strncpy(g_mqtt_instances[0].config.password, MQTT0CFG_MAIN_PWD, MQTT_MAX_PASSWORD_LEN - 1);
- g_mqtt_instances[0].config.qos = MQTT0CFG_MAIN_QOS;
- }
- else
- {
- /* Backup platform */
- rt_strncpy(g_mqtt_instances[0].config.host, MQTT0CFG_VICE_HOST, MQTT_MAX_HOST_LEN - 1);
- g_mqtt_instances[0].config.port = MQTT0CFG_VICE_PORT;
- BUILD_TOPIC(temp_topic, MQTT0CFG_VICE_TSUB, imei_suffix);
- rt_strncpy(g_mqtt_instances[0].config.sub_topic, temp_topic, MQTT_MAX_TOPIC_LEN - 1);
- BUILD_TOPIC(temp_topic, MQTT0CFG_VICE_TPUB, imei_suffix);
- rt_strncpy(g_mqtt_instances[0].config.pub_topic, temp_topic, MQTT_MAX_TOPIC_LEN - 1);
- rt_strncpy(g_mqtt_instances[0].config.client_id,
- (strlen(sys_info.imei) < 6) ? timestampe : sys_info.imei,
- MQTT_MAX_CLIENT_ID_LEN - 1);
- rt_strncpy(g_mqtt_instances[0].config.username, MQTT0CFG_VICE_USR, MQTT_MAX_USERNAME_LEN - 1);
- rt_strncpy(g_mqtt_instances[0].config.password, MQTT0CFG_VICE_PWD, MQTT_MAX_PASSWORD_LEN - 1);
- g_mqtt_instances[0].config.qos = MQTT0CFG_VICE_QOS;
- }
- /* Instance 1 is customer config */
- SERVER_CFG *cfg = &sys_info.customer_cfg;
- g_mqtt_instances[1].config.enabled = (cfg->enable == VALID) ? 1 : 0;
- g_mqtt_instances[1].config.client_idx = 1;
- rt_strncpy(g_mqtt_instances[1].config.host, cfg->host, MQTT_MAX_HOST_LEN - 1);
- g_mqtt_instances[1].config.port = cfg->port;
- if (strstr(sys_info.customer_cfg.client_id, "USE_IMEI") != NULL)
- {
- rt_strncpy(g_mqtt_instances[1].config.client_id,
- (strlen(sys_info.imei) < 6) ? timestampe : sys_info.imei,
- MQTT_MAX_CLIENT_ID_LEN - 1);
- }
- else
- {
- rt_strncpy(g_mqtt_instances[1].config.client_id, cfg->client_id, MQTT_MAX_CLIENT_ID_LEN - 1);
- }
- rt_strncpy(g_mqtt_instances[1].config.username, cfg->username, MQTT_MAX_USERNAME_LEN - 1);
- rt_strncpy(g_mqtt_instances[1].config.password, cfg->password, MQTT_MAX_PASSWORD_LEN - 1);
- /* 拼接订阅主题 */
- BUILD_TOPIC(temp_topic, cfg->sub_topic, imei_suffix);
- rt_strncpy(g_mqtt_instances[1].config.sub_topic, temp_topic, MQTT_MAX_TOPIC_LEN - 1);
- /* 拼接发布主题 */
- BUILD_TOPIC(temp_topic, cfg->pub_topic, imei_suffix);
- rt_strncpy(g_mqtt_instances[1].config.pub_topic, temp_topic, MQTT_MAX_TOPIC_LEN - 1);
- g_mqtt_instances[1].config.qos = cfg->qos;
- g_mqtt_instances[1].state = MQTT_STATE_DISCONNECTED;
- g_mqtt_instances[1].recv_callback = RT_NULL;
- LOG_I("MQTT config initialized");
- }
- /*=============================================================================
- * Network Operations
- *=============================================================================*/
- int ym310_mqtt_open_network(rt_uint8_t client_idx)
- {
- at_resp_type_t resp;
- char cmd[256];
- char urc_buf[64];
- mqtt_instance_t *inst = &g_mqtt_instances[client_idx];
- if (!inst->config.enabled)
- {
- return -RT_ERROR;
- }
- LOG_I("MQTT[%d] Opening network to %s:%d...",
- client_idx, inst->config.host, inst->config.port);
- resp = ym310_at_exec_cmd_with_retry(
- "AT+QMTCFG=\"recv/mode\",0,0,1", RT_NULL, 0,
- YM310_AT_CMD_TIMEOUT_MS, 2);
- if (resp != AT_RESP_OK)
- {
- LOG_W("MQTT[%d] Failed to configure recv mode", client_idx);
- }
- rt_snprintf(cmd, sizeof(cmd), "AT+QMTOPEN=%d,\"%s\",%d",
- client_idx, inst->config.host, inst->config.port);
- resp = ym310_at_exec_async(cmd, "+QMTOPEN:", urc_buf, sizeof(urc_buf), 30000);
- if (resp != AT_RESP_OK)
- {
- LOG_E("MQTT[%d] Failed to open network", client_idx);
- return -RT_ERROR;
- }
- int result_code = ym310_str_extract_nth_int(urc_buf, 1);
- if (result_code != 0 && result_code != 1)
- {
- LOG_E("MQTT[%d] Network open failed: result=%d", client_idx, result_code);
- return -RT_ERROR;
- }
- LOG_I("MQTT[%d] Network opened (result=%d)", client_idx, result_code);
- rt_thread_mdelay(6000);
- return RT_EOK;
- }
- int ym310_mqtt_close_network(rt_uint8_t client_idx)
- {
- at_resp_type_t resp;
- char cmd[64];
- LOG_I("MQTT[%d] Closing network...", client_idx);
- rt_snprintf(cmd, sizeof(cmd), "AT+QMTDISC=%d", client_idx);
- resp = ym310_at_exec_async(cmd, "+QMTDISC:", RT_NULL, 0, 10000);
- if (resp != AT_RESP_OK)
- {
- LOG_W("MQTT[%d] Network close timeout", client_idx);
- return -RT_ERROR;
- }
- LOG_I("MQTT[%d] Network closed", client_idx);
- return RT_EOK;
- }
- int ym310_mqtt_connect(rt_uint8_t client_idx)
- {
- at_resp_type_t resp;
- char cmd[512];
- char urc_buf[64];
- mqtt_instance_t *inst = &g_mqtt_instances[client_idx];
- if (!inst->config.enabled)
- {
- return -RT_ERROR;
- }
- LOG_I("MQTT[%d] Connecting to broker...", client_idx);
- inst->state = MQTT_STATE_CONNECTING;
- if (inst->config.username[0] != '\0' && inst->config.password[0] != '\0')
- {
- rt_snprintf(cmd, sizeof(cmd), "AT+QMTCONN=%d,\"%s\",\"%s\",\"%s\"",
- client_idx, inst->config.client_id,
- inst->config.username, inst->config.password);
- }
- else
- {
- rt_snprintf(cmd, sizeof(cmd), "AT+QMTCONN=%d,\"%s\"",
- client_idx, inst->config.client_id);
- }
- resp = ym310_at_exec_async(cmd, "+QMTCONN:", urc_buf, sizeof(urc_buf), 30000);
- if (resp != AT_RESP_OK)
- {
- LOG_E("MQTT[%d] Connect failed", client_idx);
- inst->state = MQTT_STATE_DISCONNECTED;
- return -RT_ERROR;
- }
- int result = ym310_str_extract_nth_int(urc_buf, 1);
- int ret_code = ym310_str_extract_nth_int(urc_buf, 2);
- if (result != 0)
- {
- LOG_E("MQTT[%d] Connection failed: result=%d, ret=%d", client_idx, result, ret_code);
- inst->state = MQTT_STATE_DISCONNECTED;
- return -RT_ERROR;
- }
- inst->state = MQTT_STATE_CONNECTED;
- inst->last_activity = rt_tick_get();
- LOG_I("MQTT[%d] Connected successfully", client_idx);
- return RT_EOK;
- }
- int ym310_mqtt_disconnect(rt_uint8_t client_idx)
- {
- mqtt_instance_t *inst = &g_mqtt_instances[client_idx];
- if (inst->state == MQTT_STATE_DISCONNECTED)
- {
- return RT_EOK;
- }
- inst->state = MQTT_STATE_DISCONNECTING;
- ym310_mqtt_close_network(client_idx);
- inst->state = MQTT_STATE_DISCONNECTED;
- LOG_I("MQTT[%d] Disconnected", client_idx);
- return RT_EOK;
- }
- int ym310_mqtt_subscribe(rt_uint8_t client_idx, const char *topic, rt_uint8_t qos)
- {
- at_resp_type_t resp;
- char cmd[256];
- char urc_buf[64];
- mqtt_instance_t *inst = &g_mqtt_instances[client_idx];
- rt_uint8_t msg_id = s_msg_id++;
- if (!inst->config.enabled || inst->state != MQTT_STATE_CONNECTED)
- {
- return -RT_ERROR;
- }
- // LOG_I("MQTT[%d] Subscribing to %s (qos=%d)...", client_idx, topic, qos);
- rt_snprintf(cmd, sizeof(cmd), "AT+QMTSUB=%d,%d,\"%s\",%d",
- client_idx, msg_id, topic, qos);
- resp = ym310_at_exec_async(cmd, "+QMTSUB:", urc_buf, sizeof(urc_buf), 30000);
- if (resp != AT_RESP_OK)
- {
- LOG_E("MQTT[%d] Subscribe failed", client_idx);
- return -RT_ERROR;
- }
- int result = -1, granted_qos = -1;
- sscanf(urc_buf, "+QMTSUB: %*d,%*d,%d,%d", &result, &granted_qos);
- if (result != 0)
- {
- LOG_E("MQTT[%d] Subscribe result=%d", client_idx, result);
- return -RT_ERROR;
- }
- rt_kprintf("\n MQTT[%d] Subscribed to %s (qos=%d) SUCCESSFULLY\n\n", client_idx, topic, granted_qos);
- return RT_EOK;
- }
- int ym310_mqtt_unsubscribe(rt_uint8_t client_idx, const char *topic)
- {
- at_resp_type_t resp;
- char cmd[256];
- mqtt_instance_t *inst = &g_mqtt_instances[client_idx];
- rt_uint8_t msg_id = s_msg_id++;
- if (!inst->config.enabled || inst->state != MQTT_STATE_CONNECTED)
- {
- return -RT_ERROR;
- }
- rt_snprintf(cmd, sizeof(cmd), "AT+QMTUNS=%d,%d,\"%s\"", client_idx, msg_id, topic);
- resp = ym310_at_exec_cmd_with_retry(cmd, RT_NULL, 0, 30000, 2);
- if (resp != AT_RESP_OK)
- {
- LOG_W("MQTT[%d] Unsubscribe failed", client_idx);
- return -RT_ERROR;
- }
- LOG_I("MQTT[%d] Unsubscribed from %s", client_idx, topic);
- return RT_EOK;
- }
- static void ym310_mqtt_thread_entry(void *param)
- {
- int ret;
- rt_tick_t last_check[MQTT_INSTANCE_MAX] = {0};
- LOG_D("========================================");
- LOG_D("MQTT THREAD STARTED (Simplified v1.0)");
- LOG_D("========================================");
- /* Init MQTT config (keep as-is) */
- ym310_mqtt_init_config();
- /* Register default callbacks (example print) */
- g_mqtt_callbacks[0] = mqtt0_recv_callback; // Defined at file top
- g_mqtt_callbacks[1] = mqtt1_recv_callback;
- /* Register MQTT URC handlers (keep as-is) */
- ym310_at_register_urc("+QMTRECV:", urc_qmtrecv_handler);
- ym310_at_register_urc("+QMTSTAT:", urc_qmtstat_handler);
- ym310_at_register_urc("+QMTCONN:", urc_qmtconn_handler);
- ym310_at_register_urc("+QMTOPEN:", urc_qmtopen_handler);
- s_mqtt_thread_running = RT_TRUE;
- while (s_mqtt_thread_running)
- {
- /****************** 1. Process recv data (simplified) ******************/
- for (int i = 0; i < MQTT_INSTANCE_MAX; i++)
- {
- rt_bool_t has_data = RT_FALSE;
- char topic[MQTT_MAX_TOPIC_LEN] = {0};
- char payload[MQTT_MAX_PAYLOAD_SIZE] = {0};
- rt_uint16_t payload_len = 0;
- /* Disable interrupts to safely read from recv buffer and clear flag */
- // rt_base_t //level = rt_hw_interrupt_disable();
- if (g_mqtt_recv_bufs[i].data_ready)
- {
- /* Copy to local variable */
- rt_strncpy(topic, g_mqtt_recv_bufs[i].topic, MQTT_MAX_TOPIC_LEN - 1);
- rt_strncpy(payload, g_mqtt_recv_bufs[i].payload, MQTT_MAX_PAYLOAD_SIZE - 1);
- payload_len = g_mqtt_recv_bufs[i].payload_len;
- /* Clear buffer (only clear flag, no need to clear data) */
- g_mqtt_recv_bufs[i].data_ready = RT_FALSE;
- g_mqtt_recv_bufs[i].payload_len = 0;
- g_mqtt_recv_bufs[i].topic[0] = '\0';
- g_mqtt_recv_bufs[i].payload[0] = '\0';
- has_data = RT_TRUE;
- }
- // rt_hw_interrupt_enable(level);
- if (has_data)
- {
- if (ym310_at_mutex_lock(1000) == RT_EOK)
- {
- if (g_mqtt_callbacks[i] != RT_NULL)
- {
- g_mqtt_callbacks[i](i, topic, payload, payload_len);
- }
- else
- {
- LOG_D("MQTT[%d] RX: topic=%s, payload=%.*s",
- i, topic, payload_len, payload);
- }
- ym310_at_mutex_unlock();
- }
- else
- {
- LOG_W("MQTT[%d] Failed to lock AT mutex for processing", i);
- }
- }
- }
- /****************** 2. Connection management (keep as-is) ******************/
- for (int i = 0; i < MQTT_INSTANCE_MAX; i++)
- {
- mqtt_instance_t *inst = &g_mqtt_instances[i];
- if (!inst->config.enabled)
- {
- continue;
- }
- // rt_base_t //level = rt_hw_interrupt_disable();
- mqtt_state_t current_state = inst->state;
- // rt_hw_interrupt_enable(level);
- switch (current_state)
- {
- case MQTT_STATE_DISCONNECTED:
- {
- rt_tick_t current_tick = rt_tick_get();
- rt_tick_t retry_interval;
- rt_uint32_t backoff = 1;
- for (int j = 0; j < inst->err_count && j < 4; j++)
- {
- backoff *= 2;
- }
- retry_interval = CONNECT_RETRY_BASE_MS * backoff;
- if (retry_interval > CONNECT_RETRY_MAX_MS)
- {
- retry_interval = CONNECT_RETRY_MAX_MS;
- }
- if ((current_tick - s_last_connect_attempt[i]) > rt_tick_from_millisecond(retry_interval))
- {
- s_last_connect_attempt[i] = current_tick;
- // LOG_I("MQTT[%d] Connecting (err=%u, interval=%ums)...",
- // i, inst->err_count, retry_interval);
- ret = ym310_mqtt_open_network(i);
- if (ret == RT_EOK)
- {
- ret = ym310_mqtt_connect(i);
- if (ret == RT_EOK)
- {
- ym310_mqtt_subscribe(i, inst->config.sub_topic, inst->config.qos);
- }
- else
- {
- ym310_mqtt_close_network(i);
- }
- }
- if (ret != RT_EOK)
- {
- LOG_W("MQTT[%d] Connection failed, will retry", i);
- }
- }
- break;
- }
- case MQTT_STATE_CONNECTED:
- if ((rt_tick_get() - last_check[i]) > rt_tick_from_millisecond(60000))
- {
- last_check[i] = rt_tick_get();
- LOG_D("MQTT[%d] TX:%u RX:%u ERR:%u",
- i, inst->tx_count, inst->rx_count, inst->err_count);
- }
- break;
- case MQTT_STATE_RECONNECTING:
- LOG_I("MQTT[%d] Reconnecting...", i);
- ym310_mqtt_disconnect(i);
- // level = rt_hw_interrupt_disable();
- inst->state = MQTT_STATE_DISCONNECTED;
- // rt_hw_interrupt_enable(level);
- break;
- default:
- break;
- }
- }
- rt_thread_mdelay(50); // Main loop interval
- }
- /****************** 3. Cleanup (keep as-is) ******************/
- for (int i = 0; i < MQTT_INSTANCE_MAX; i++)
- {
- if (g_mqtt_instances[i].config.enabled)
- {
- ym310_mqtt_disconnect(i);
- }
- }
- ym310_at_unregister_urc("+QMTRECV:");
- ym310_at_unregister_urc("+QMTSTAT:");
- ym310_at_unregister_urc("+QMTCONN:");
- ym310_at_unregister_urc("+QMTOPEN:");
- LOG_I("MQTT thread stopped");
- }
- /*=============================================================================
- * Public API
- *=============================================================================*/
- int ym310_mqtt_init(void)
- {
- if (g_mqtt_mutex == RT_NULL)
- {
- g_mqtt_mutex = rt_mutex_create("mqtt_mtx", RT_IPC_FLAG_FIFO);
- if (g_mqtt_mutex == RT_NULL)
- {
- LOG_E("Failed to create MQTT mutex!");
- return -RT_ERROR;
- }
- // LOG_D("g_mqtt_mutex created at 0x%p", g_mqtt_mutex);
- }
- // LOG_I("MQTT framework initialized");
- return RT_EOK;
- }
- void ym310_mqtt_deinit(void)
- {
- /* Cleanup handled by thread */
- }
- rt_thread_t ym310_start_mqtt_thread(void)
- {
- if (ym310_mqtt_init() != RT_EOK)
- {
- return RT_NULL;
- }
- g_mqtt_thread = rt_thread_create("ym310_mqtt", ym310_mqtt_thread_entry, RT_NULL,
- MQTT_THREAD_STACK_SIZE, MQTT_THREAD_PRIORITY, 10);
- if (g_mqtt_thread == RT_NULL)
- {
- LOG_E("Failed to create MQTT thread!");
- return RT_NULL;
- }
- rt_thread_startup(g_mqtt_thread);
- // LOG_I("MQTT thread started");
- return g_mqtt_thread;
- }
- mqtt_state_t ym310_mqtt_get_state(rt_uint8_t client_idx)
- {
- if (client_idx >= MQTT_INSTANCE_MAX)
- {
- return MQTT_STATE_DISCONNECTED;
- }
- return g_mqtt_instances[client_idx].state;
- }
- rt_bool_t ym310_mqtt_is_connected(rt_uint8_t client_idx)
- {
- if (client_idx >= MQTT_INSTANCE_MAX)
- {
- return RT_FALSE;
- }
- return (g_mqtt_instances[client_idx].state == MQTT_STATE_CONNECTED);
- }
- int ym310_mqtt_set_recv_callback(rt_uint8_t client_idx,
- void (*callback)(rt_uint8_t, const char *,
- const char *, rt_uint16_t))
- {
- if (client_idx >= MQTT_INSTANCE_MAX)
- {
- return -RT_ERROR;
- }
- // rt_base_t //level = rt_hw_interrupt_disable();
- g_mqtt_callbacks[client_idx] = (mqtt_recv_handler_t)callback;
- // rt_hw_interrupt_enable(level);
- LOG_I("MQTT[%d] Callback registered", client_idx);
- return RT_EOK;
- }
- void ym310_mqtt_force_reconnect(rt_uint8_t client_idx)
- {
- if (client_idx >= MQTT_INSTANCE_MAX)
- {
- return;
- }
- mqtt_instance_t *inst = &g_mqtt_instances[client_idx];
- LOG_I("MQTT[%d] Forcing reconnection...", client_idx);
- ym310_mqtt_disconnect(client_idx);
- // rt_base_t //level = rt_hw_interrupt_disable();
- inst->state = MQTT_STATE_DISCONNECTED;
- inst->err_count = 0;
- // rt_hw_interrupt_enable(level);
- s_last_connect_attempt[client_idx] = 0;
- LOG_I("MQTT[%d] Reconnection triggered", client_idx);
- }
- void ym310_mqtt_force_reconnect_all(void)
- {
- LOG_I("Forcing reconnection for all MQTT instances...");
- for (int i = 0; i < MQTT_INSTANCE_MAX; i++)
- {
- if (g_mqtt_instances[i].config.enabled)
- {
- ym310_mqtt_force_reconnect(i);
- }
- }
- }
|